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

如何将输入列表与自定义DAG字典结合?并结合get、delayed或futures使用

我来帮你解决这个问题!你已经有了处理单个文件的自定义DAG,现在要批量处理文件列表又不想重复写DAG项,同时还要用get来评估,结合delayed或者futures的话,可以这么做:

方法1:用delayed批量生成任务,转成DAG后用get执行

这是最简洁高效的方式,核心是把单个文件的处理逻辑封装成函数,用delayed包装后批量生成任务,再自动转换为DAG字典供get使用,完全不用手动重复编写每个文件的DAG节点。

示例代码

假设你已经有了单个文件的基础处理函数:

def read_file(path):
    # 你的文件读取逻辑
    return raw_data

def process_data(raw_data):
    # 你的数据处理逻辑
    return processed_data

def save_result(data, output_path):
    # 你的结果保存逻辑
    return None

第一步:把单个文件的完整流程封装成函数,用delayed装饰

from dask import delayed

@delayed
def process_single_file(input_path, output_path):
    raw = read_file(input_path)
    processed = process_data(raw)
    save_result(processed, output_path)
    return processed  # 可选:返回处理后的数据,方便后续收集结果

第二步:针对文件列表批量生成delayed任务

input_files = ["data/file1.txt", "data/file2.txt", "data/file3.txt"]
output_files = ["output/file1.txt", "output/file2.txt", "output/file3.txt"]

# 为每个文件生成对应的delayed任务
batch_tasks = [process_single_file(inp, out) for inp, out in zip(input_files, output_files)]

第三步:转换为DAG字典并执行

from dask.base import get

# 将delayed任务集合自动转为标准DAG字典,还能自动优化图结构
dag = dask.base.collections_to_dsk(batch_tasks, optimize_graph=True)

# 定义要获取的结果节点(用任务索引作为键,对应每个delayed任务的输出)
result_keys = list(range(len(batch_tasks)))

# 用本地get执行;如果是分布式集群,替换为client.get即可
results = get(dag, result_keys)

这种方式完全复用了单个文件的处理逻辑,自动帮你生成所有文件的DAG节点,完美解决重复编写DAG项的问题。

方法2:手动构建批量DAG字典(更灵活)

如果你需要对DAG节点名称、跨文件依赖有更精细的控制,可以手动封装单个文件的DAG生成逻辑,再批量合并成大DAG:

def build_single_file_dag(input_path, output_path, node_prefix):
    """生成单个文件的DAG片段,用前缀区分不同文件的节点"""
    return {
        f"{node_prefix}_read": (read_file, input_path),
        f"{node_prefix}_process": (process_data, f"{node_prefix}_read"),
        f"{node_prefix}_save": (save_result, f"{node_prefix}_process", output_path)
    }

# 构建完整的批量DAG
full_dag = {}
for idx, (inp, out) in enumerate(zip(input_files, output_files)):
    # 用file0、file1这样的前缀区分不同文件的节点
    file_dag = build_single_file_dag(inp, out, f"file{idx}")
    full_dag.update(file_dag)

# 定义要执行的结果节点(比如所有保存操作的节点)
result_nodes = [f"file{idx}_save" for idx in range(len(input_files))]

# 用get执行DAG
results = get(full_dag, result_nodes)

这种方式适合需要自定义节点标识、或者添加跨文件依赖的场景,但相比delayed会繁琐一些。

关于Futures的结合

如果用Dask Distributed集群,你可以直接把转换后的DAG提交给集群,结合Futures异步执行:

from dask.distributed import Client

client = Client()  # 连接到你的分布式集群

# 将delayed任务转为DAG后提交,获取Futures对象
dag = dask.base.collections_to_dsk(batch_tasks)
futures = client.get(dag, result_keys, asynchronous=True)

# 等待并收集结果
results = client.gather(futures)

这样既保留了用get评估DAG的逻辑,又能利用分布式集群的Futures异步能力。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 07:46:42