Python Polars:如何用LazyFrame实现大型Parquet文件的连接?
解决方案与
iter_slices说明 一、处理大型Parquet文件连接的可行方案
针对LazyFrame连接大文件导致崩溃的问题,推荐以下几种方案:
1. 启用Polars流式处理模式
LazyFrame本身支持延迟执行,在最终collect时开启流式处理,Polars会自动将数据分块加载到内存,避免一次性加载全量数据:
import polars as pl # 扫描两个Parquet文件 lf1 = pl.scan_parquet("large_file1.parquet") lf2 = pl.scan_parquet("large_file2.parquet") # 执行连接并启用流式处理 result = lf1.join(lf2, on="join_key").collect(streaming=True)
这种方式无需手动分块,Polars会根据内存情况自动管理分块逻辑。
2. 手动分块读取与连接
如果流式处理仍有问题,可以手动控制分块大小,逐块读取并连接:
import polars as pl # 获取每个文件的总行数 total_rows1 = lf1.select(pl.count()).collect().item() total_rows2 = lf2.select(pl.count()).collect().item() chunk_size = 10_000_000 # 根据内存情况调整分块大小 result_chunks = [] for i in range(0, total_rows1, chunk_size): # 分块读取第一个文件的切片 chunk1 = lf1.slice(i, chunk_size).collect() # 用当前切片连接第二个文件(可优化为仅连接匹配行) chunk_join = chunk1.join(lf2.collect(streaming=True), on="join_key") result_chunks.append(chunk_join) # 合并所有分块结果 final_result = pl.concat(result_chunks)
如果第二个文件也过大,可以同时对两个文件分块,或者先对第二个文件按连接键建立索引,减少每次连接的开销。
3. 按连接键分区预处理
将其中一个大文件按连接键分区存储,之后逐块读取另一个文件,仅与对应分区连接:
# 将第二个文件按连接键分区写入磁盘 lf2.write_parquet("partitioned_file2", partition_by="join_key") # 遍历第一个文件的分块,与对应分区连接 result_chunks = [] for i in range(0, total_rows1, chunk_size): chunk1 = lf1.slice(i, chunk_size).collect() # 获取当前chunk包含的唯一连接键 unique_keys = chunk1["join_key"].unique() # 仅读取匹配的分区 matching_partitions = pl.scan_parquet("partitioned_file2", filters=[("join_key", "in", unique_keys)]).collect() chunk_join = chunk1.join(matching_partitions, on="join_key") result_chunks.append(chunk_join) final_result = pl.concat(result_chunks)
二、iter_slices的作用
iter_slices是Polars DataFrame的专用方法,仅支持已加载到内存中的DataFrame,作用是将完整的DataFrame拆分为指定大小的小DataFrame迭代器,适合以下场景:
- 批量处理内存中的数据,比如分批次写入数据库、分批次执行计算逻辑,避免单次处理过大的内存负载;
- 对数据进行增量式操作,比如逐块生成统计报告、逐块进行数据清洗。
示例用法:
df = pl.DataFrame({"a": range(100)}) for slice_df in df.iter_slices(30): print(slice_df.shape) # 对每个切片执行操作
内容的提问来源于stack exchange,提问作者Hernan Diaz
相关产品推荐
相关产品推荐

