Polars read_database搭配SQLAlchemy/OracleDB时iter_batches=True不生效
解决Polars + SQLAlchemy(Oracle)分批读取失效的问题
问题原因
Polars的pl.read_database参数iter_batches=True在Oracle的SQLAlchemy引擎下存在兼容性问题,无法正确触发流式分批读取,导致一次性加载全量数据引发内存溢出。
可行解决方案
方案1:基于SQLAlchemy流式查询手动分批
利用SQLAlchemy原生的stream_results=True开启流式查询,手动逐批获取数据并转换为Polars DataFrame处理:
from sqlalchemy import text import polars as pl batch_size = 1000000 sql = sql_text(sql) # 你的SQL语句 with e.connect() as conn: # 开启流式查询,禁止一次性加载所有结果到内存 result = conn.execute(text(sql), stream_results=True) # 逐批读取数据 while True: batch_rows = result.fetchmany(size=batch_size) if not batch_rows: break # 将批次数据转为Polars DataFrame,自动继承结果集的字段名 df_batch = pl.DataFrame(batch_rows, schema=result.keys()) # 替换为你的批次处理逻辑(比如写入Parquet、计算等) df_batch.write_parquet(f"batch_{offset}.parquet")
方案2:Polars懒加载 + Offset/Limit分批
通过Polars的scan_database创建懒加载DataFrame,配合offset和limit实现分批读取,适合需要在Polars中做链式数据处理的场景:
注意:必须为SQL语句添加稳定的ORDER BY子句(比如基于主键或时间戳),否则可能出现数据重复或遗漏
import polars as pl batch_size = 1000000 lazy_df = pl.scan_database(sql_text(sql) + " ORDER BY id", e) # 替换id为你的排序字段 offset = 0 while True: df_batch = lazy_df.offset(offset).limit(batch_size).collect() if df_batch.is_empty(): break # 处理当前批次 process_batch(df_batch) offset += batch_size
方案3:直接使用OracleDB原生连接
绕过SQLAlchemy中间层,直接用oracledb驱动实现分批读取,性能更优:
import oracledb import polars as pl # 替换为你的Oracle连接信息 conn = oracledb.connect(user="your_user", password="your_pwd", dsn="your_dsn") cursor = conn.cursor() cursor.arraysize = 1000000 # 设置每次读取的批次大小 cursor.execute(sql_text(sql)) # 获取字段名作为DataFrame schema schema = [col[0] for col in cursor.description] while True: batch_rows = cursor.fetchmany() if not batch_rows: break df_batch = pl.DataFrame(batch_rows, schema=schema) # 处理批次数据 process_batch(df_batch) cursor.close() conn.close()
内容的提问来源于stack exchange,提问作者Niels Jespersen
相关产品推荐
相关产品推荐

