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

如何不使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 18:57:06