如何完整合并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
相关产品推荐
相关产品推荐

