Dask自定义任务图:Delayed.compute后DataFrame未物化,如何一步完成计算?
解决Dask自定义任务图首次compute返回延迟DataFrame的问题
你的问题核心在于:任务图里的节点返回的是延迟状态的Dask DataFrame,而非物化后的Pandas DataFrame。Delayed的.compute()仅执行了任务图中的函数逻辑,得到了Dask DataFrame对象,但Dask DataFrame内部的计算任务并未触发。要实现首次compute就得到物化结果,有两种可行方案:
方案1:在任务图中直接包含DataFrame物化步骤
修改任务图内的函数,让每个节点直接返回物化后的Pandas DataFrame:
import dask.dataframe as dd import pandas as pd from dask.delayed import Delayed from functools import partial person = pd.DataFrame({ 'name': ['John', 'Jane'], 'age': [30, 25] }) task_graph = { 'df': ( lambda df: dd.from_pandas(df, npartitions=1).compute(), # 直接完成物化 person ), 'person_old': ( lambda x: x.assign(age=x.age*3), # 此时x是Pandas DF,操作后直接返回Pandas DF 'df' ), } request = Delayed(['df','person_old'], task_graph) result = request.compute() display(result[0]) # 直接显示物化后的Pandas DataFrame display(result[1])
方案2:用dask.compute()替代Delayed对象的.compute()
保持原任务图不变,通过dask.compute()触发所有Dask DataFrame的计算:
import dask.dataframe as dd import pandas as pd from dask.delayed import Delayed from functools import partial import dask person = pd.DataFrame({ 'name': ['John', 'Jane'], 'age': [30, 25] }) task_graph = { 'df': ( partial(dd.from_pandas, npartitions=1), person ), 'person_old': ( lambda x: x.assign(age=x.age*3), 'df' ), } request = Delayed(['df','person_old'], task_graph) result = dask.compute(*request) # 用dask.compute处理Delayed对象 display(result[0]) # 直接得到物化后的Pandas DataFrame display(result[1])
两种方案都能确保首次计算就得到完全物化的结果,同时Dask仍会自动优化公共依赖(比如df节点只会被计算一次)。
内容的提问来源于stack exchange,提问作者Konstantin
相关产品推荐
相关产品推荐

