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

使用Dask DataFrame to_csv后大量数据丢失的问题求助

Dask分布式合并丢失大量行的问题修复

问题原因

  1. 错误使用scatter+submit:Dask DataFrame本身具备分布式执行能力,手动通过client.scatter广播数据集再用client.submit执行函数,会破坏Dask的任务调度逻辑,导致部分分区未被完整计算,出现行数丢失且结果不稳定的情况。
  2. 合并逻辑未明确连接键:依赖默认合并行为(按共同列名合并),容易因索引/列的定义差异引发意外匹配。

修复方案

修正后的代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 19:39:51