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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 11:45:55