如何将输入列表与自定义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
相关产品推荐
相关产品推荐

