如何不使用Delayed函数在多个Dask DataFrame上应用自定义有状态函数?
解决方案
首先纠正你当前代码里的一个问题:dd.read_parquet本身就返回延迟计算的Dask DataFrame,完全不需要用@dask.delayed包装加载逻辑——这也是你不符合最佳实践的原因之一,用delayed包裹Dask DataFrame会破坏它自带的分区、列裁剪等优化能力。
针对你的需求(调用有内部状态、无法并行的operation函数),无需使用Delayed函数的可行方案分两种场景:
场景1:数据量较小,能全部放入内存
直接通过dask.compute把三个Dask DataFrame的目标列计算为Pandas DataFrame,再传入operation:
import dask.dataframe as dd # 1. 直接读取数据源,不需要delayed sources_path = ["/source1", "/source2", "/source3"] ddfs = [dd.read_parquet(path, engine="pyarrow") for path in sources_path] # 2. 定义每个数据源需要的列 cols_source1 = ["col_a", "col_b"] # 替换为你的列列表 cols_source2 = ["col_c", "col_d"] cols_source3 = ["col_e", "col_f"] # 3. 提取目标列的Dask DataFrame ddf1 = ddfs[0][cols_source1] ddf2 = ddfs[1][cols_source2] ddf3 = ddfs[2][cols_source3] # 4. 计算为Pandas DataFrame并执行操作 df1, df2, df3 = dask.compute(ddf1, ddf2, ddf3) result = operation(df1, df2, df3)
场景2:数据量较大,本地内存无法容纳
借助Dask分布式集群的Client.submit,将计算任务提交到单个Worker节点,让Worker负责加载、计算数据并执行operation(避免把大数据拉到本地):
import dask.dataframe as dd from dask.distributed import Client # 初始化分布式客户端 client = Client() # 1. 读取数据源 sources_path = ["/source1", "/source2", "/source3"] ddfs = [dd.read_parquet(path, engine="pyarrow") for path in sources_path] # 2. 定义目标列 cols_source1 = ["col_a", "col_b"] cols_source2 = ["col_c", "col_d"] cols_source3 = ["col_e", "col_f"] # 3. 定义执行函数,在Worker上完成数据计算与操作 def run_operation(ddf1, ddf2, ddf3, cols1, cols2, cols3): df1 = ddf1[cols1].compute() df2 = ddf2[cols2].compute() df3 = ddf3[cols3].compute() return operation(df1, df2, df3) # 4. 提交任务到Worker并获取结果 future = client.submit(run_operation, ddfs[0], ddfs[1], ddfs[2], cols_source1, cols_source2, cols_source3) result = future.result()
核心逻辑说明
因为operation存在内部状态,必须一次性处理所有数据,不能拆分并行。上述两种方案都是确保三个数据源的全部目标数据被聚合后,再调用operation,同时避免了不必要的Delayed包装,符合Dask的最佳实践。
内容的提问来源于stack exchange,提问作者theShadow89
相关产品推荐
相关产品推荐

