如何用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
相关产品推荐
相关产品推荐

