如何在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}")
这种方式下:
job_1执行成功后,Sensor会触发job_2运行;job_2的输出通过IO Manager存储到本地文件;job_3运行时,通过IO Manager读取job_2的输出作为输入;job_2执行成功后,Sensor触发job_3运行。
关键注意事项
- 避免直接用
job_3(job_2())的写法,这不符合Dagster的执行模型; - 如果任务逻辑紧密,优先用方案一的Graph/Job整合方式,更简洁高效;
- 若需独立调度、隔离任务执行,选择方案二的Sensor+IO Manager组合。
内容的提问来源于stack exchange,提问作者SS_Rebelious
相关产品推荐
相关产品推荐

