使用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
相关产品推荐
相关产品推荐

