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

使用dask.dataframe.from_delayed后persist重复加载数据的问题咨询

解决Dask from_delayed persist后重复执行加载任务的问题

问题核心

用dd.from_delayed加载数据后调用persist(),后续计算依然会重新执行load任务、重复从网络加载数据,而dd.read_parquet则不会出现这个问题。

解决方案1:给delayed任务设置pure=True

修改生成delayed对象的代码,给delayed装饰器加上pure=True参数,告诉Dask这个函数是纯函数(相同输入返回相同结果、无副作用),这样Dask会缓存任务结果,后续计算直接复用:

def load(file):
    # 实际加载逻辑更复杂,此处为示例
    return pd.read_parquet(file)

filenames = ['my_file']
# 添加pure=True标记任务为纯函数
df = dd.from_delayed([delayed(load, pure=True)(f) for f in filenames])
df = df.persist()
print(df.count().compute())

解决方案2:手动预缓存delayed结果

如果你的load函数有副作用(比如每次调用会修改外部状态),无法设置pure=True,可以先手动计算并缓存delayed结果,再传入from_delayed:

from dask import delayed, compute

def load(file):
    # 实际加载逻辑更复杂,此处为示例
    return pd.read_parquet(file)

filenames = ['my_file']
delayed_dfs = [delayed(load)(f) for f in filenames]
# 先执行delayed任务并缓存结果
cached_dfs = compute(*delayed_dfs)[0]
# 从缓存的结果生成Dask DataFrame
df = dd.from_delayed(cached_dfs)
df = df.persist()
print(df.count().compute())

原因解释

  • delayed默认pure=False,Dask会认为每个调用都是独立的、可能产生不同结果的任务,即使参数相同。所以persist()执行后,后续计算还是会重新调度任务。
  • dd.read_parquet内部生成的任务默认被标记为纯任务,Dask会自动缓存这些任务的结果,persist()后后续计算直接复用缓存,不会重复加载数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 11:35:14