Метод stream класса AsyncSession
Метод stream класса AsyncSession
выполняет SQL-запрос и возвращает объект
AsyncResult, который позволяет
итерироваться по результатам потоково, не
загружая весь набор строк в память сразу.
Это особенно полезно при работе с большими
таблицами. Первым параметром метод принимает
SQL-запрос или ORM-запрос, а дополнительные
параметры передаются в метод execute.
Синтаксис
await session.stream(statement, [params])
Пример
Давайте создадим таблицу и выполним потоковую выборку всех записей:
import asyncio
from sqlalchemy import text
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
async def main():
engine = create_async_engine('sqlite+aiosqlite:///:memory:')
async with engine.begin() as conn:
await conn.execute(text(
'CREATE TABLE articles ('
'id INTEGER PRIMARY KEY, '
'title VARCHAR, '
'text VARCHAR, '
'status VARCHAR, '
'num INTEGER)'
))
await conn.execute(text(
"INSERT INTO articles (title, num) VALUES "
"('article 1', 1), "
"('article 2', 2), "
"('article 3', 3)"
))
async with AsyncSession(engine) as session:
res = await session.stream(text('SELECT title FROM articles'))
async for row in res:
print(row.title)
asyncio.run(main())
Результат выполнения кода:
"article 1"
"article 2"
"article 3"
Пример
Давайте выполним потоковую выборку через ORM-модель. Сначала определим модель:
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
class Base(DeclarativeBase):
pass
class Article(Base):
__tablename__ = 'articles'
id: Mapped[int] = mapped_column(primary_key=True)
title: Mapped[str]
text: Mapped[str]
status: Mapped[str]
num: Mapped[int]
Теперь используем stream для потоковой
выборки объектов:
import asyncio
from sqlalchemy import select
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from models import Base, Article
async def main():
engine = create_async_engine('sqlite+aiosqlite:///:memory:')
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
async with AsyncSession(engine) as session:
session.add_all([
Article(title='article 1', text='text 1', status='new', num=1),
Article(title='article 2', text='text 2', status='new', num=2),
Article(title='article 3', text='text 3', status='new', num=3),
])
await session.commit()
res = await session.stream(select(Article))
async for article in res:
print(article.title)
asyncio.run(main())
Результат выполнения кода:
"article 1"
"article 2"
"article 3"
Смотрите также
-
класс
AsyncSession,
который представляет асинхронную сессию -
метод
execute,
который выполняет SQL-запрос и возвращает результат -
метод
stream_scalars,
который выполняет потоковую выборку скалярных значений -
метод
scalars,
который возвращает только первые колонки результата