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

如何用PyArrow实现超内存数据集的惰性连接与新列创建?

问题描述

我有超出内存容量的数据集,能通过PyArrow Dataset API读取,但执行连接操作时会触发内存内连接,连接结果也超出内存容量。想让PyArrow强制实现惰性执行,同时还要通过行级两列求和创建新列。

以下是DuckDB中的实现示例:

import duckdb
import pyarrow as pa
import pandas as pd

rt = pa.Table.from_pandas(pd.DataFrame.from_dict({'passenger_count':[0,1,2,3,4]*5,'rand':[5,3,2,1,2]*5}),
                          schema=pa.schema([('passenger_count',pa.int8()),('rand',pa.int8())]))
duckdb.sql("""
copy 
(
    select *,a.trip_distance + b.rand as test
    from read_parquet('Z:/nyc_hive/yr=2009/*/*parquet',hive_partitioning=True) a 
    inner join rt b
    on a.passenger_count = b.passenger_count
)
TO 'nyc_duck' (FORMAT PARQUET, PARTITION_BY (yr, month));
""")

我尝试了以下PyArrow代码,但出现内存溢出错误:

import pyarrow as pa
import pyarrow.dataset as ds
# Create a PyArrow dataset for each table
table_a_dataset = ds.dataset('Z:/a77980db-7f0a-41b0-94ff-0913c98e183b/nyc_hive/yr=2009/', format='parquet')
j = table_a_dataset.join(rt,'passenger_count') #triggers the join
#what i want:
# j = table1.join(table2).project(col('a') + col('b')) # dont trigger
# j.write(XYZ.parquet) # now it triggers the data pull, join and write. never OOMing
用PyArrow复现的解决方案

PyArrow Dataset API的join方法默认会立即执行,要实现惰性执行,需要结合PyArrow的计算表达式和dataset.scan()构建查询计划,最后通过write_dataset分批处理数据,避免内存溢出。

步骤1:构建惰性查询计划

使用dataset.scan()创建扫描器,再用pyarrow.dataset.join(而非Dataset对象的join方法)构建连接计划,最后通过project添加新列的计算逻辑,整个过程不会触发实际数据加载:

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

# 构建小表rt
rt = pa.Table.from_pandas(
    pd.DataFrame.from_dict({
        'passenger_count': [0,1,2,3,4]*5,
        'rand': [5,3,2,1,2]*5
    }),
    schema=pa.schema([('passenger_count', pa.int8()), ('rand', pa.int8())])
)

# 创建大表的数据集扫描器
table_a_scan = ds.dataset(
    'Z:/a77980db-7f0a-41b0-94ff-0913c98e183b/nyc_hive/yr=2009/',
    format='parquet'
).scan()

# 构建惰性连接计划:指定连接键和连接类型
joined_scan = ds.join(
    left=table_a_scan,
    right=rt,
    keys='passenger_count',
    join_type='inner'
)

# 添加新列:计算trip_distance + rand,命名为test
projected_scan = joined_scan.project([
    pc.field('*'),  # 保留所有原有列
    pc.add(pc.field('trip_distance'), pc.field('rand')).alias('test')  # 新增计算列
])

步骤2:惰性执行并写入分区Parquet

使用ds.write_dataset执行查询计划,该方法会自动分批处理数据,避免一次性加载所有数据到内存,同时支持按指定字段分区:

# 写入分区Parquet文件,按yr和month分区
ds.write_dataset(
    projected_scan,
    destination='nyc_pyarrow',
    format='parquet',
    partitioning=['yr', 'month'],
    # 可选:根据内存情况调整批次大小
    write_options=ds.WriteOptions(batch_size=10_000_000)
)

关键说明

  • ds.join接收扫描器(Scanner)作为输入,返回的仍是扫描器,不会立即执行连接操作,实现了惰性计算。
  • project方法基于扫描器添加计算逻辑,同样不会触发数据加载,直到调用write_dataset时才会分批处理。
  • write_dataset自动处理分区,且可通过batch_size控制单批次数据量,适配内存容量,避免OOM。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 20:25:29