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

如何在Dagster中实现任务串联?含任务间数据传递需求

Dagster任务串联与数据传递实现方案

核心概念澄清

Dagster中,@job装饰的是独立执行单元,无法直接像普通函数那样返回值并传递给其他Job。要实现任务依赖与数据传递,有两种主流方式:


方案一:整合为单个Graph/Job(适合简单依赖场景)

将所有操作(Op)整合到一个Graph或Job中,直接通过Op的输出传递数据,同时定义执行顺序。

from dagster import job, op

# 定义基础Op
@op
def make_api_call_op():
    # 模拟API调用返回数据
    return {"data": "api_data_1"}

@op
def save_to_db_op(data):
    # 模拟存入数据库
    print(f"Saving {data} to DB")

@op
def make_another_api_call_op():
    return {"data": "api_data_2"}

@op(out={"out_1": str, "out_2": int})
def process_data_op(input_data):
    # 模拟处理数据,返回两个输出
    processed_1 = f"processed_{input_data['data']}"
    processed_2 = 123
    return processed_1, processed_2

@op
def process_op(input_data):
    print(f"Processing job_2 output: {input_data}")

@op
def do_some_other_staff_op():
    print("Doing other work")

# 整合为单个Job,定义执行顺序与数据传递
@job
def full_pipeline():
    # 执行job_1的逻辑
    save_to_db_op(make_api_call_op())
    
    # 执行job_2的逻辑,获取输出
    out_1, out_2 = process_data_op(make_another_api_call_op())
    save_to_db_op(out_1)
    
    # 用job_2的输出执行job_3的逻辑
    process_op(out_2)
    do_some_other_staff_op()

这种方式下,Dagster会自动处理依赖:make_api_call_op完成后才会执行save_to_db_op;make_another_api_call_op完成后执行process_data_op,其输出分别传给save_to_db_op和process_op;所有前置Op完成后才会执行后续Op。


方案二:独立Job + Sensor + IO Manager(适合需独立调度的场景)

如果需要将任务拆分为独立Job(比如单独调度、重启),可以通过Sensor监听Job的执行状态,触发后续Job;同时用IO Manager存储Job的输出,供下游Job读取。

步骤1:定义带输出的Job,配置IO Manager

from dagster import job, op, Out, In, IOManager
from dagster.core.storage.io_manager import io_manager
import pickle
import os

# 自定义IO Manager,用于存储和读取Job输出
class LocalPickleIOManager(IOManager):
    def __init__(self, base_dir):
        self.base_dir = base_dir
        os.makedirs(base_dir, exist_ok=True)
    
    def handle_output(self, context, obj):
        # 存储输出到文件
        file_path = os.path.join(self.base_dir, f"{context.job_name}_output.pkl")
        with open(file_path, "wb") as f:
            pickle.dump(obj, f)
    
    def load_input(self, context):
        # 读取上游Job的输出
        upstream_job_name = context.upstream_output.job_name
        file_path = os.path.join(self.base_dir, f"{upstream_job_name}_output.pkl")
        with open(file_path, "rb") as f:
            return pickle.load(f)

# 注册IO Manager
@io_manager(config_schema={"base_dir": str})
def local_pickle_io_manager(context):
    return LocalPickleIOManager(context.resource_config["base_dir"])

# 定义job_1
@job(resource_defs={"io_manager": local_pickle_io_manager})
def job_1():
    save_to_db_op(make_api_call_op())

# 定义job_2,返回out_2作为输出
@job(resource_defs={"io_manager": local_pickle_io_manager})
def job_2():
    out_1, out_2 = process_data_op(make_another_api_call_op())
    save_to_db_op(out_1)
    # 将out_2标记为Job输出
    return out_2

# 定义job_3,接收输入
@job(resource_defs={"io_manager": local_pickle_io_manager})
def job_3():
    # 读取job_2的输出作为输入
    job_2_output = job_2()
    process_op(job_2_output)
    do_some_other_staff_op()

步骤2:定义Sensor监听Job完成事件,触发后续Job

from dagster import SensorDefinition, RunRequest, sensor

# 监听job_1完成,触发job_2
@sensor(job=job_2)
def job1_finished_trigger_job2(context):
    # 获取最近成功的job_1运行记录
    last_run = context.instance.get_latest_job_run(job_1.name, status="SUCCESS")
    if last_run and not context.instance.has_run(job_2.name, last_run.run_id):
        yield RunRequest(run_key=f"job1_{last_run.run_id}")

# 监听job_2完成,触发job_3
@sensor(job=job_3)
def job2_finished_trigger_job3(context):
    last_run = context.instance.get_latest_job_run(job_2.name, status="SUCCESS")
    if last_run and not context.instance.has_run(job_3.name, last_run.run_id):
        yield RunRequest(run_key=f"job2_{last_run.run_id}")

这种方式下:

  1. job_1执行成功后,Sensor会触发job_2运行;
  2. job_2的输出通过IO Manager存储到本地文件;
  3. job_3运行时,通过IO Manager读取job_2的输出作为输入;
  4. job_2执行成功后,Sensor触发job_3运行。

关键注意事项

  • 避免直接用job_3(job_2())的写法,这不符合Dagster的执行模型;
  • 如果任务逻辑紧密,优先用方案一的Graph/Job整合方式,更简洁高效;
  • 若需独立调度、隔离任务执行,选择方案二的Sensor+IO Manager组合。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 18:51:14