Polars scan_ndjson为何无法在流式处理场景下工作?
问题描述
尝试读取一个300GB的换行分隔JSON(ndjson)文件,提取特定字段并写入Parquet文件。每个JSON对象相互独立,本应可流式分块处理(文件无法全部载入内存)。使用如下代码:
# define the schema... pl.scan_ndjson( 'data/input/myjson.jsonl', schema=prschema)\ .collect(streaming=True)\ .write_parquet('data/output/myparquet.parquet', compression='snappy', use_pyarrow=True )
但随着处理文件子集增大,内存消耗随文件大小线性增长。用explain(streaming=True)查看执行计划,发现并未启用流式处理:
Anonymous SCAN PROJECT */6 COLUMNS
改用sink_parquet替代write_parquet也失败,简化代码仅提取两个标量字段仍报错:
pl.scan_ndjson('data/input/myjson.jsonl')\ .select('id', 'standing')\ .sink_parquet( 'data/output/myparquet.parquet', compression='snappy' )
报错信息:InvalidOperationError: sink_Parquet(ParquetWriteOptions { compression: Snappy, statistics: false, row_group_size: None, data_pagesize_limit: None, maintain_order: true }) not yet supported in standard engine. Use 'collect().write_parquet()'
请问为何这种简单的读写场景无法使用流式处理?
核心原因与解决方案
1. collect(streaming=True)未生效的原因
collect(streaming=True)的streaming参数仅控制读取阶段的分块行为,但collect方法本身的作用是将所有分块数据聚合到内存中生成DataFrame,最终还是会把全量数据载入内存,因此内存消耗随文件大小线性增长。同时当前Polars的查询优化器不会为这种简单的扫描+投影操作生成真正的流式执行计划,导致explain结果中没有流式标识。
2. sink_parquet报错的原因
Polars的标准引擎暂未支持sink_parquet操作,该功能仅在开发中的cloud引擎或特定扩展中可用,官方报错信息已明确提示这一点。
可行的流式处理方案
要实现真正的分块流式处理,避免内存溢出,可采用以下两种方式:
方式一:手动分块读取写入
通过逐行读取文件,按固定大小分块处理,每块处理完成后追加写入Parquet:
import polars as pl # 按内存承受能力调整块大小 chunk_size = 10_000_000 # 替换为你的实际Schema schema = pl.Schema([("id", pl.UInt64), ("standing", pl.Boolean)]) first_write = True with open('data/input/myjson.jsonl', 'r') as f: while True: lines = [] # 读取当前块的行数据 for _ in range(chunk_size): line = f.readline() if not line: break lines.append(line) if not lines: break # 处理当前块并提取字段 df_chunk = pl.read_ndjson(lines, schema=schema).select('id', 'standing') # 写入Parquet,首次写入创建文件,后续追加 df_chunk.write_parquet( 'data/output/myparquet.parquet', compression='snappy', use_pyarrow=True, append=not first_write ) first_write = False
方式二:使用iter_batches迭代处理
利用Polars的iter_batches方法生成数据批次迭代器,逐批写入Parquet:
import polars as pl # 替换为你的实际Schema schema = pl.Schema([("id", pl.UInt64), ("standing", pl.Boolean)]) # 生成批次迭代器,调整batch_size适配内存 batch_iterator = pl.scan_ndjson('data/input/myjson.jsonl', schema=schema)\ .select('id', 'standing')\ .iter_batches(batch_size=10_000_000) first_write = True for batch in batch_iterator: batch.write_parquet( 'data/output/myparquet.parquet', compression='snappy', use_pyarrow=True, append=not first_write ) first_write = False
内容的提问来源于stack exchange,提问作者teejay

