Метод stream класса AsyncConnection
Метод stream класса AsyncConnection
выполняет SQL-запрос и возвращает объект
асинхронного потока результатов. В отличие от
метода execute, который буферизует все
строки в памяти, метод stream позволяет
обрабатывать результаты построчно по мере их
поступления из базы данных. Это особенно важно
при работе с большими наборами данных, когда
загрузка всех строк сразу может привести
к перерасходу памяти. Первым параметром метод
принимает SQL-выражение или текстовый запрос,
вторым - словарь параметров. Возвращаемый
объект является асинхронным контекстным
менеджером и итератором одновременно.
Синтаксис
AsyncConnection.stream(statement, parameters)
Пример
Давайте создадим соединение и выполним
построчное чтение строк из таблицы
articles:
from sqlalchemy import text
from sqlalchemy.ext.asyncio import create_async_engine
engine = create_async_engine('sqlite+aiosqlite:///:memory:')
async def main():
async with engine.connect() as conn:
await conn.execute(text('CREATE TABLE articles (id INTEGER, title TEXT)'))
await conn.execute(text("INSERT INTO articles VALUES (1, 'article 1'), (2, 'article 2')"))
await conn.commit()
async with conn.stream(text('SELECT title FROM articles ORDER BY id')) as stream:
async for row in stream:
print(row.title)
import asyncio
asyncio.run(main())
Результат выполнения кода:
"article 1"
"article 2"
Пример
Давайте передадим параметры в запрос и
отфильтруем строки по значению status:
from sqlalchemy import text
from sqlalchemy.ext.asyncio import create_async_engine
engine = create_async_engine('sqlite+aiosqlite:///:memory:')
async def main():
async with engine.connect() as conn:
await conn.execute(text('CREATE TABLE articles (id INTEGER, title TEXT, status TEXT)'))
await conn.execute(text("INSERT INTO articles VALUES (1, 'article 1', 'published'), (2, 'article 2', 'draft')"))
await conn.commit()
async with conn.stream(
text('SELECT title FROM articles WHERE status = :status'),
{'status': 'published'}
) as stream:
async for row in stream:
print(row.title)
import asyncio
asyncio.run(main())
Результат выполнения кода:
"article 1"
Смотрите также
-
класс
AsyncConnection,
который представляет асинхронное соединение с базой -
метод
execute,
который выполняет запрос с буферизацией результата -
метод
stream_scalars,
который возвращает поток скалярных значений -
метод
run_sync,
который выполняет синхронную функцию в асинхронном контексте