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

PyArrow:Dataset.to_batches与ParquetFile.iter_batches内存消耗差异问询

PyArrow中Dataset.to_batches与ParquetFile.iter_batches的内存差异问题

使用pyarrow.dataset的场景

>>> import pyarrow as pa
>>> import pyarrow.dataset as ds
>>>
>>> filelist = [
...     'local_data/file1.parquet'
... ]
>>> input_dataset = ds.dataset(source=filelist)
>>> sink = pa.BufferOutputStream()
>>> writer = pa.ipc.new_stream(sink, input_dataset.schema)
>>> ds_batches = input_dataset.to_batches(batch_size=500, batch_readahead=0, fragment_readahead=0)
>>> for _ in range(4):
...     writer.write_batch(next(ds_batches))
...
>>> buf = sink.getvalue()
>>> buf.size
5637392

内存观测:执行ds_batches = input_dataset.to_batches(...)前内存约44MB,执行后跃升至约815MB;迭代过程中内存保持稳定。

使用pyarrow.parquet.ParquetFile的场景

>>> import pyarrow as pa
>>> import pyarrow.parquet as pq
>>>
>>> filelist = [
...     'local_data/file1.parquet'
... ]
>>> pqfile = pq.ParquetFile(filelist[0])
>>> sink = pa.BufferOutputStream()
>>> writer = pa.ipc.new_stream(sink, pqfile.schema_arrow)
>>> pqiterator = pq.ParquetFile(filelist[0]).iter_batches(batch_size=500)
>>> for _ in range(4):
...     writer.write_batch(next(pqiterator))
...
>>> buf = sink.getvalue()
>>> buf.size
5636632

内存观测:执行pqiterator = pq.ParquetFile(...).iter_batches(...)后内存无明显变化,始终维持在约19MB的低水平。

注意:两种场景最终生成的缓冲区大小均约为5MB,读取的记录数量一致且符合预期。

疑问

  1. 是什么导致了Dataset.to_batches的高内存消耗?
  2. 为何Dataset.to_batches与ParquetFile.iter_batches同样用于迭代记录批次,内存使用却存在如此大的差异?

补充说明

  1. to_batches文档描述为:“Read the dataset as materialized record batches.”,曾疑惑“materialized”是否意味着返回所有批次,但文件总批次的实际占用空间远大于观测到的内存值,此推测不成立。
  2. 即使改用Dataset.scanner()替代Dataset.to_batches,内存表现仍一致,尽管scanner文档明确说明:“Data is not loaded immediately. Instead, this produces a Scanner, which exposes further operations”。

问题解答

1. Dataset.to_batches高内存消耗的原因

PyArrow的Dataset模块在初始化扫描器(to_batches会隐式创建扫描器)时,会预先加载文件的完整元数据与统计信息,并为后续查询执行做预处理:

  • 对于Parquet文件,Dataset会解析所有行组、列的统计信息、Schema细节等全量元数据,这类数据的内存占用会随着文件规模、列数、行组数的增加而显著上升。
  • 即便设置了batch_readahead=0和fragment_readahead=0,Dataset的执行引擎仍会完成查询规划工作,比如确定行组读取顺序、校验Schema一致性等,这些操作会占用额外内存。
  • Dataset是为多文件、多格式的通用数据场景设计的,内部结构比单一Parquet文件的专用封装更复杂,初始化阶段的内存开销必然更高。

2. 两者内存差异的核心原因

ParquetFile与Dataset的设计目标和实现逻辑存在本质区别:

  • ParquetFile是针对单一Parquet文件的轻量级工具,iter_batches初始化时仅读取最基础的文件头部信息,不会加载全量元数据或执行预处理,只有在调用next()迭代时才会读取对应批次的数据,因此内存始终维持在低水平。
  • Dataset面向大规模多源数据的查询处理,需要支持多文件合并、分区过滤、跨格式Schema统一等复杂场景,因此在创建扫描器(或调用to_batches)时,必须预先加载所有关联文件的元数据并完成查询规划,这部分内存开销是其通用性带来的必然代价。

补充说明中scanner()内存表现一致的原因是:to_batches本质是scanner().to_batches()的简写,两者共享相同的初始化逻辑,都会预先加载元数据与执行规划,因此内存表现完全一致。

内容的提问来源于stack exchange,提问作者teejay

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 01:23:38