РЕПЕТИТОР математика физика информатика
Для школьников и студентов. Подтягивание пробелов. ЦЭ, ЦТ, ОГЭ, ЕГЭ.
Идет набор на ЛЕТО. Жмите для подробностей:)
585 of 711 menu

Метод 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-запрос
Мы используем cookie для работы сайта, аналитики и персонализации. Обработка данных происходит согласно Политике конфиденциальности.
принять все настроить отклонить