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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:45:04