You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Polars scan_ndjson为何无法在流式处理场景下工作?

大体积NDJSON转Parquet的流式处理问题解答

问题描述

尝试读取一个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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 18:25:56