Метод stream_scalars
Метод stream_scalars класса AsyncSession
выполняет SQL-запрос и возвращает асинхронный
поток скалярных значений. В отличие от метода
scalars, который загружает весь результат
в память, метод stream_scalars выдает
строки по мере их поступления из базы данных.
Первым параметром метод принимает SQL-запрос
или ORM-выражение, а дополнительные параметры
передаются в метод execute.
Синтаксис
await session.stream_scalars(statement, [params])
Пример
Давайте создадим асинхронную сессию и получим
поток скалярных значений из таблицы articles:
import asyncio
from sqlalchemy import String, select
from sqlalchemy.ext.asyncio import (
AsyncSession,
create_async_engine,
)
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] = mapped_column(String(50))
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'),
Article(title='article 2'),
Article(title='article 3'),
])
await session.commit()
stmt = select(Article.title)
result = await session.stream_scalars(stmt)
async for title in result:
print(title)
asyncio.run(main())
Результат выполнения кода:
"article 1"
"article 2"
"article 3"
Пример
Давайте применим фильтрацию и получим только
те статьи, у которых статус равен published:
import asyncio
from sqlalchemy import String, select
from sqlalchemy.ext.asyncio import (
AsyncSession,
create_async_engine,
)
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] = mapped_column(String(50))
status: Mapped[str] = mapped_column(String(20))
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', status='published'),
Article(title='article 2', status='draft'),
Article(title='article 3', status='published'),
])
await session.commit()
stmt = select(Article.title).where(
Article.status == 'published'
)
result = await session.stream_scalars(stmt)
async for title in result:
print(title)
asyncio.run(main())
Результат выполнения кода:
"article 1"
"article 3"
Пример
Давайте получим поток целых чисел, используя
агрегатную функцию func.sum:
import asyncio
from sqlalchemy import String, func, select
from sqlalchemy.ext.asyncio import (
AsyncSession,
create_async_engine,
)
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] = mapped_column(String(50))
num: Mapped[int]
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', num=10),
Article(title='article 2', num=20),
Article(title='article 3', num=30),
])
await session.commit()
stmt = select(func.sum(Article.num))
result = await session.stream_scalars(stmt)
async for total in result:
print(total)
asyncio.run(main())
Результат выполнения кода:
"60"
Пример
Давайте прервем итерацию по потоку досрочно и закроем результат:
import asyncio
from sqlalchemy import String, select
from sqlalchemy.ext.asyncio import (
AsyncSession,
create_async_engine,
)
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] = mapped_column(String(50))
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'),
Article(title='article 2'),
Article(title='article 3'),
])
await session.commit()
stmt = select(Article.title)
result = await session.stream_scalars(stmt)
async for title in result:
print(title)
break
await result.close()
asyncio.run(main())
Результат выполнения кода:
"article 1"
Смотрите также
-
класс
AsyncSession,
который представляет асинхронную сессию -
метод
stream,
который возвращает поток ORM-объектов -
метод
scalars,
который возвращает все скалярные значения сразу -
метод
execute,
который выполняет SQL-запрос