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

如何完整合并Dask DataFrame?求高效合并方案

解决Dask DataFrame左连接后行数大幅减少的问题

咱们先拆解问题根源,再一步步找到最优解决方案:

为什么合并后行数仅90万?

你的Sales表有2800万行,左连接理论上应该保留所有行(哪怕部分行没匹配到Product表记录,对应Product列会显示NaN)。出现行数骤减的核心原因是小表分区策略错误,导致Dask没有采用高效的广播合并,反而做了分区对分区的shuffle合并;另外也可能是关联键的数据类型不匹配,导致大部分行无法匹配。

分步解决方案

1. 调整小表分区数

Product表仅600行,完全不需要分成3个分区。将其设为1个分区后,Dask会自动识别这是小表,将其广播到Sales表的每个分区,确保每个Sales分区都能和完整的Product表做连接:

sales_dd = dd.from_pandas(Sales, npartitions=3)  # 大表分区保持不变
product_dd = dd.from_pandas(Product, npartitions=1)  # 小表仅需1个分区

2. 检查并统一关联键数据类型

如果Sales的ProductNo和Product的ProductNo类型不一致(比如一个是整数、一个是字符串),会导致大量行匹配失败。先检查类型:

print("Sales ProductNo 类型:", sales_dd['ProductNo'].dtype)
print("Product ProductNo 类型:", product_dd['ProductNo'].dtype)

若类型不同,统一转换为相同类型:

# 将Product的ProductNo转换为与Sales一致的类型
product_dd['ProductNo'] = product_dd['ProductNo'].astype(sales_dd['ProductNo'].dtype)

3. 显式启用广播合并(可选)

如果Dask未自动识别小表,可在merge时强制指定broadcast=True,确保小表被广播到每个大表分区:

productsales = dd.merge(
    sales_dd, 
    product_dd, 
    on='ProductNo', 
    how='left', 
    broadcast=True
)

4. 验证合并结果行数

注意Dask是延迟计算的,直接看tail()或len()无法得到真实总行数,需触发计算:

total_rows = productsales.shape[0].compute()
print(f"合并后总行数:{total_rows}")

最快的合并方式

当一张表远小于另一张表时,广播合并是最优选择:

  • 无需对大表做shuffle(shuffle会产生大量磁盘IO和数据传输,耗时极长)
  • 小表仅需加载一次,分发到每个大表分区进行本地连接
  • 大表分区可并行处理,充分利用多核CPU资源

完全不需要取消大表的分区,2800万行分成3个分区是合理的,能保证并行处理的效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:05:53