使用Dask DataFrame to_csv后大量数据丢失的问题求助
Dask分布式合并丢失大量行的问题修复
问题原因
- 错误使用
scatter+submit:Dask DataFrame本身具备分布式执行能力,手动通过client.scatter广播数据集再用client.submit执行函数,会破坏Dask的任务调度逻辑,导致部分分区未被完整计算,出现行数丢失且结果不稳定的情况。 - 合并逻辑未明确连接键:依赖默认合并行为(按共同列名合并),容易因索引/列的定义差异引发意外匹配。
修复方案
修正后的代码
import numpy as np import pandas as pd import dask.dataframe as dd from dask.distributed import Client import dask # 数据生成逻辑保持不变 data_vol = 2000 index = pd.date_range("2021-09-01", periods=data_vol, freq="1h") df = pd.DataFrame({"a": np.arange(data_vol), "b": ["abcaddbe"] * data_vol, 'time': index}) ddf = dd.from_pandas(df, npartitions=10) df2 = pd.DataFrame({"c": np.arange(data_vol), "d": ["xyzopq"] * data_vol, 'time': reversed(index)}) ddf2 = dd.from_pandas(df2, npartitions=10) ddf['timestamp'] = ddf.time.apply(lambda x: int(x.timestamp()), meta=('time', 'int64')) ddf2['timestamp'] = ddf2.time.apply(lambda x: int(x.timestamp()), meta=('time', 'int64')) # 修改合并函数,显式指定按索引合并 def merge_onindex(ddf, ddf2): ret = ddf.merge(ddf2, left_index=True, right_index=True) ret["add"] = ret.a + ret.c + 1 return ret dask.config.set({"dataframe.shuffle.method": "tasks"}) client = Client("tcp://172.17.0.2:8786") # 直接对Dask DataFrame执行操作,无需手动scatter/submit ddf_indexed = ddf.set_index('timestamp') ddf2_indexed = ddf2.set_index('timestamp') dd_merge_res = merge_onindex(ddf_indexed, ddf2_indexed) # 直接触发计算并写入文件,Dask自动处理分布式调度 dd_merge_res.to_csv("/jupyter/merge_single.csv", single_file=True)
关键修改点
- 移除
scatter和submit:Dask会自动协调分布式任务,无需手动传递数据集到worker。 - 显式指定合并键:通过
left_index=True, right_index=True明确按timestamp索引合并,避免默认行为的不确定性。
验证结果
执行完成后用wc -l merge_single.csv检查,应得到2001行(含表头)的完整结果,每行的add值均为2000。
内容的提问来源于stack exchange,提问作者Flybywind
相关产品推荐
相关产品推荐

