使用Dagster的asset API时,如何并行执行列表项对应的asset/op?
在Dagster中实现列表项的并行处理
可以实现,核心是利用Dagster的**动态输出(Dynamic Outputs)**机制,它支持基于运行时生成的列表动态创建并行任务。以下是两种适配Asset API的实现方案:
方案1:结合动态Op与Asset API
先定义生成文件列表的Asset,再通过带动态输出的Op处理每个文件,实现并行执行:
from pathlib import Path from dagster import asset, op, job, DynamicOut, DynamicOutput, OpExecutionContext # 生成文件列表的Asset @asset def list_of_files(): path = Path("/some/path") return list(path.iterdir()) # 分发文件的动态Op:为每个文件生成独立的并行任务 @op(out=DynamicOut()) def dispatch_files(context: OpExecutionContext, files: list[Path]): for idx, file in enumerate(files): yield DynamicOutput(file, mapping_key=f"file_{idx}") # 单个文件的处理逻辑 @op def process_single_file(context: OpExecutionContext, file: Path): context.log.info(f"Processing file: {file.name}") with open(file, 'rt') as f: for line in f: # 执行具体业务操作,比如写入数据库 pass return f"Completed processing {file.name}" # 定义Job串联Asset与Op @job def file_processing_job(): files = list_of_files.to_source_asset() # 将分发的动态输出映射到处理Op,自动并行执行 dispatch_files(files).map(process_single_file)
方案2:生成动态子Asset(追踪单个文件处理状态)
如果需要将每个文件的处理结果作为独立Asset(方便单独追踪、重新执行),可以用multi_asset生成动态子Asset:
from pathlib import Path from dagster import asset, multi_asset, AssetOut, AssetKey, AssetExecutionContext # 生成文件名列表的Asset @asset def file_names(): path = Path("/some/path") return [file.name for file in path.iterdir()] # 为每个文件生成独立的子Asset @multi_asset( outs={name: AssetOut(key=AssetKey(f"processed_{name}")) for name in file_names()} ) def process_files_as_assets(context: AssetExecutionContext, file_names: list[str]): base_path = Path("/some/path") results = {} for name in file_names: file_path = base_path / name context.log.info(f"Processing asset: {name}") with open(file_path, 'rt') as f: for line in f: # 业务操作逻辑 pass results[name] = f"Processed asset {name}" return results
关键配置说明
要真正实现并行,需确保执行器支持多进程/分布式执行:
- 默认的
in_process_executor是串行的,需切换为multiprocess_executor,在dagster.yaml中配置:
execution: multiprocess: max_concurrent: 4 # 设置最大并行任务数
内容的提问来源于stack exchange,提问作者Martin Preusse
相关产品推荐
相关产品推荐

