Python中小DataFrame排序后与超大规模DataFrame的最快左连接方案
大规模Parquet数据左连接优化实现
问题背景
需要处理两个Parquet格式的数据集:
small_df:约50万行,未按列A排序huge_df:超100亿行,已按列A排序
在内存与CPU资源充足(数百核)的前提下,要最快完成small_df按列A排序后与huge_df的左连接操作。
最优实现方案(基于Polars)
Polars是高性能列式数据处理库,支持懒加载查询优化、多核并行处理,能最大化利用硬件资源。以下是完整可复现代码:
import polars as pl import numpy as np import time import os # 格式化文件大小为可读格式 def format_bytes(bytes_size): """Convert bytes to human readable format""" for unit in ['B', 'KB', 'MB', 'GB']: if bytes_size < 1024.0: return f"{bytes_size:.1f} {unit}" bytes_size /= 1024.0 return f"{bytes_size:.1f} TB" def create_test_data(small_size=500_000, huge_size=200_000_000, temp_dir="./temp_data"): """生成测试数据并保存为Parquet文件,模拟真实I/O场景""" os.makedirs(temp_dir, exist_ok=True) # 生成已按A列排序的huge_df(10列) huge_df = pl.DataFrame({ 'A': range(1, huge_size + 1), # 连续值即已排序 'huge_value': np.random.randn(huge_size), 'huge_category': np.random.choice(['X', 'Y', 'Z'], huge_size), 'huge_score': np.random.uniform(0, 100, huge_size), 'huge_flag': np.random.choice([True, False], huge_size), 'huge_price': np.random.exponential(50, huge_size), 'huge_rating': np.random.choice([1, 2, 3, 4, 5], huge_size), 'huge_region': np.random.choice(['North', 'South', 'East', 'West'], huge_size), 'huge_count': np.random.poisson(10, huge_size), 'huge_percentage': np.random.beta(2, 5, huge_size) * 100 }) # 生成未按A列排序的small_df(20列,与huge_df有80%键重叠) overlap_size = int(small_size * 0.8) small_keys = ( list(np.random.choice(huge_size, overlap_size, replace=False)) + # 存在于huge_df的键 list(range(huge_size + 1, huge_size + small_size - overlap_size + 1)) # 不存在的键 ) np.random.shuffle(small_keys) # 打乱顺序,模拟未排序 small_df = pl.DataFrame({ 'A': small_keys, 'small_value': np.random.randn(small_size), 'small_flag': np.random.choice([True, False], small_size), 'small_category': np.random.choice(['Alpha', 'Beta', 'Gamma', 'Delta'], small_size), 'small_score': np.random.uniform(0, 1000, small_size), 'small_price': np.random.lognormal(4, 1, small_size), 'small_quantity': np.random.randint(1, 100, small_size), 'small_status': np.random.choice(['Active', 'Inactive', 'Pending'], small_size), 'small_weight': np.random.gamma(2, 2, small_size), 'small_temperature': np.random.normal(20, 5, small_size), 'small_pressure': np.random.exponential(1.5, small_size), 'small_density': np.random.uniform(0.8, 1.2, small_size), 'small_ph': np.random.normal(7, 0.5, small_size), 'small_conductivity': np.random.exponential(10, small_size), 'small_viscosity': np.random.gamma(1.5, 2, small_size), 'small_opacity': np.random.beta(2, 3, small_size), 'small_hardness': np.random.choice([1, 2, 3, 4, 5, 6, 7, 8, 9, 10], small_size), 'small_color': np.random.choice(['Red', 'Blue', 'Green', 'Yellow', 'Orange'], small_size), 'small_texture': np.random.choice(['Smooth', 'Rough', 'Grainy', 'Silky'], small_size), 'small_timestamp': np.random.randint(1600000000, 1700000000, small_size) }) small_path = f"{temp_dir}/small_df.parquet" huge_path = f"{temp_dir}/huge_df.parquet" small_df.write_parquet(small_path) huge_df.write_parquet(huge_path) small_size_bytes = os.path.getsize(small_path) huge_size_bytes = os.path.getsize(huge_path) print(f"Small parquet file: {format_bytes(small_size_bytes)}") print(f"Huge parquet file: {format_bytes(huge_size_bytes)}") return small_path, huge_path def optimized_join(small_parquet_path, huge_parquet_path): start_time = time.time() # 懒加载Parquet文件,此时未加载数据到内存 small_lf = pl.scan_parquet(small_parquet_path) huge_lf = pl.scan_parquet(huge_parquet_path) # 构建查询计划,Polars会自动优化整个执行流程 result_lf = ( small_lf .sort("A") # 对small_df按A列排序 .join(huge_lf, on="A", how="left") # 与huge_df执行左连接 ) # 执行优化后的查询计划(包含I/O时间) result_df = result_lf.collect() print(f"连接完成耗时: {time.time() - start_time:.2f}s") return result_df def main(): print("生成测试数据...") small_path, huge_path = create_test_data() # 快速验证数据行数 small_row_count = pl.scan_parquet(small_path).select(pl.len()).collect().item() huge_row_count = pl.scan_parquet(huge_path).select(pl.len()).collect().item() print(f"Small数据集: {small_row_count:,} 行(未排序)") print(f"Huge数据集: {huge_row_count:,} 行(已排序)") print("\n--- Polars排序并连接 ---") result = optimized_join(small_path, huge_path) print(f"结果形状: {result.shape}") print(f"成功匹配行数: {result.filter(pl.col('huge_value').is_not_null()).height}") # 清理临时文件 try: os.remove(small_path) os.remove(huge_path) os.rmdir(os.path.dirname(small_path)) print("已清理临时文件") except OSError: print(f"注意: 临时文件残留于 {os.path.dirname(small_path)}") if __name__ == "__main__": # 依赖安装: pip install polars numpy main()
测试结果(ThinkPad笔记本非生产环境)
生成测试数据... Small parquet file: 44.8 MB Huge parquet file: 6.1 GB Small数据集: 500,000 行(未排序) Huge数据集: 200,000,000 行(已排序) --- Polars排序并连接 --- 连接完成耗时: 12.71s 结果形状: (500000, 29) 成功匹配行数: 400000 已清理临时文件
来源标注
内容的提问来源于stack exchange,提问作者Nicolò Cavalleri
相关产品推荐
相关产品推荐

