Используй «расширяемые» параметры IN в SQLAlchemy + потоковую выборку — это безопаснее, быстрее и не бьёт по памяти. Вместо склейки списка ID в строку передай Python-список как параметр с
expanding=True: драйвер сам подставит нужное число плейсхолдеров и корректно закэширует план. Совмести это с итерацией по результатам (chunks) — и обрабатывай миллионы строк без OOM. Ставь лайк и подпишись на нас, каждый день мы публикуем полезные и не банальные советы для разработчиков.
pip install sqlalchemy pandas
from typing import Iterable, Sequence, Dict, Any
import pandas as pd
from sqlalchemy import create_engine, text, bindparam, tuple_
from sqlalchemy.engine import Engine, Result
# Пример: engine для Postgres/Oracle/MySQL и т.д.
# engine = create_engine("postgresql+psycopg://user:pass@host/db")
# engine = create_engine("oracle+oracledb://user:pass@tnsname")
def stream_query_in(engine: Engine, table: str, id_list: Sequence[int], chunk_rows: int = 50_000) -> Iterable[pd.DataFrame]:
"""
Потоково выбирает строки по большому списку ID:
- безопасно подставляет массив в IN (:ids) через expanding=True
- возвращает пачки DataFrame фиксированного размера (chunk_rows)
"""
# 1) Готовим SQL с расширяемым параметром
stm = text(f"SELECT * FROM {table} WHERE id IN :ids").bindparams(
bindparam("ids", expanding=True)
)
with engine.connect() as conn:
# 2) Выполняем запрос одним списком (драйвер сам развернёт IN)
result: Result = conn.execution_options(stream_results=True).execute(stm, {"ids": list(id_list)})
# 3) Потоково собираем фреймы по chunk_rows строк
batch = []
for row in result:
batch.append(dict(row._mapping))
if len(batch) >= chunk_rows:
yield pd.DataFrame.from_records(batch)
batch.clear()
if batch:
yield pd.DataFrame.from_records(batch)
def fetch_df_in(engine: Engine, table: str, id_list: Sequence[int]) -> pd.DataFrame:
"""Если нужен единый DataFrame — аккуратно склеиваем поток."""
frames = list(stream_query_in(engine, table, id_list))
return pd.concat(frames, ignore_index=True) if frames else pd.DataFrame()
# Пример использования:
# big_ids = range(1, 2_000_000)
# for df_chunk in stream_query_in(engine, "orders", big_ids, chunk_rows=100_000):
# process(df_chunk) # обрабатываем пачками без OOM
# Бонус: многоколонночный IN (tuple IN) без ручной генерации SQL
def stream_query_tuple_in(engine: Engine, table: str, pairs: Sequence[Dict[str, Any]]):
"""
pairs: [{ "country":"US", "city":"NYC" }, ...]
SELECT * FROM table WHERE (country, city) IN ((:country_1,:city_1), ...)
"""
stm = text(f"""
SELECT * FROM {table}
WHERE (country, city) IN :keys
""").bindparams(bindparam("keys", expanding=True))
# Подготовим данные как кортежи
key_tuples = [(p["country"], p["city"]) for p in pairs]
with engine.connect() as conn:
for row in conn.execution_options(stream_results=True).execute(stm, {"keys": key_tuples}):
yield dict(row._mapping)
# Пример:
# pairs = [{"country":"US","city":"NYC"}, {"country":"JP","city":"TOKYO"}]
# for rec in stream_query_tuple_in(engine, "warehouses", pairs):
# ...