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

PyArrow 18.0使用__fragment_index/__batch_index跳行时触发ArrowInvalid错误

问题:PyArrow扫描Parquet数据集时无法识别__fragment_index特殊字段

处理超大规模Parquet数据集时,需要实现跳过前n行且不加载到内存的功能。根据PyArrow 18.0文档,可通过特殊字段__fragment_index和__batch_index构造计算表达式,但运行代码时触发字段不存在的错误。

原实现代码

from pyarrow import dataset as ds
from pyarrow import compute as pc

dataset = ds.dataset(input_dir, format="parquet")

scan_columns['my_column'] = pc.field('my_column') # Arbitrary column with my data
scan_columns['__proj_index'] = (
    pc.add(pc.multiply(pc.field('__fragment_index'), pc.scalar(10000)),
                       pc.field('__batch_index')
))
filter = pc.field('__proj_index') >= pc.scalar(42000042)

batches = dataset.scanner(columns=scan_columns, filter=filter).to_batches()

触发的错误

File "pyarrow/_dataset.pyx", line 399, in pyarrow._dataset.Dataset.scanner
File "pyarrow/_dataset.pyx", line 3557, in pyarrow._dataset.Scanner.from_dataset
File "pyarrow/_dataset.pyx", line 3475, in pyarrow._dataset.Scanner._make_scan_options
File "pyarrow/_dataset.pyx", line 3422, in pyarrow._dataset._populate_builder
File "pyarrow/error.pxi", line 92, in pyarrow.lib.check_status
pyarrow.lib.ArrowInvalid: No match for FieldRef.Name(__fragment_index) in column_name: int64
问题原因及解决方案

核心原因

__fragment_index和__batch_index是PyArrow扫描过程中生成的内部元数据字段,不属于数据集本身的物理列:

  • 你使用pc.field()引用这两个字段,但pc.field()只能识别数据集已存在的物理列,因此会报错找不到字段。
  • 默认情况下,扫描器不会暴露这些元数据字段,必须显式开启开关才能使用。

修正方案

  1. 初始化扫描器时添加include_fragment_metadata=True参数,开启元数据字段的可见性
  2. 使用ds.field()(而非pc.field())引用这两个特殊元数据字段
  3. 调整表达式构造逻辑,避免在columns参数中直接构造包含元数据字段的计算列(应在过滤阶段使用这些字段)

修正后的代码

from pyarrow import dataset as ds
from pyarrow import compute as pc

dataset = ds.dataset(input_dir, format="parquet")

# 开启元数据字段可见性,指定需要加载的业务列
scanner = dataset.scanner(
    columns=["my_column"],
    include_fragment_metadata=True
)

# 构造过滤表达式:使用ds.field()引用元数据字段
filter_expr = (pc.multiply(ds.field("__fragment_index"), pc.scalar(10000)) + 
               ds.field("__batch_index")) >= pc.scalar(42000042)

# 应用过滤并获取结果批次
batches = scanner.filter(filter_expr).to_batches()

额外注意事项

__batch_index是单个Parquet文件(fragment)内的批次索引,你代码中假设每个fragment有10000个批次,这个值需要根据实际数据的批次大小调整,否则会导致跳过的行数计算不准确。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 23:04:52