Dagster工作流中RunRequest对象target参数不被识别问题排查
问题:RunStatusSensor返回RunRequest时提示缺少目标的解决办法
场景与代码
要实现一个run_status_sensor,监听etl_job执行成功事件并触发dbt_job运行,传感器代码如下(已移除部分配置细节):
from dagster import run_status_sensor, RunRequest, SkipReason, DagsterRunStatus from analytics.jobs.etl_job import dbt_job @run_status_sensor(run_status=DagsterRunStatus.SUCCESS) def trigger_dbt_job_after_etl(context): context.log.info(f"Sensor triggered for job: {context.dagster_run.job_name}") if context.dagster_run.job_name == "etl_job": context.log.info("etl_job completed successfully. Preparing to trigger dbt_job.") run_config = { "ops": { "run_dbt_task": { "config": { "cluster_name": "", "launch_type": "", "security_groups": [""], "subnets": [ "subnet-", "subnet-", "subnet-", "subnet-", "subnet-" ], "task_definition": "" } } } } run_request = RunRequest(run_key=None, run_config=run_config, job=dbt_job, job_name="dbt_job") context.log.info(f"Created RunRequest: {run_request}") return run_request else: context.log.info(f"Skipping sensor run because the job is not etl_job. Current job: {context.dagster_run.job_name}") return SkipReason("The job is not etl_job, skipping.")
传感器已正确注册到Definitions配置:
from dagster import Definitions, EnvVar from analytics.jobs.etl_job import etl_job, dbt_job from analytics.resources.ecs_resource import ecs_client from analytics.sensors.job_sensor import trigger_dbt_job_after_etl defs = Definitions( jobs=[etl_job, dbt_job], resources={ "ecs_client": ecs_client.configured({ "region_name":"us-east-1", "aws_access_key_id":"", "aws_secret_access_key":"" } ) }, sensors=[trigger_dbt_job_after_etl] )
报错信息
运行时触发如下错误:
Traceback (most recent call last): File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_daemon/sensor.py", line 471, in _process_tick_generator yield from _evaluate_sensor( File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_daemon/sensor.py", line 614, in _evaluate_sensor sensor_runtime_data = code_location.get_external_sensor_execution_data( File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_core/remote_representation/code_location.py", line 930, in get_external_sensor_execution_data return sync_get_external_sensor_execution_data_grpc( File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_api/snapshot_sensor.py", line 87, in sync_get_external_sensor_execution_data_grpc raise DagsterUserCodeProcessError.from_error_info(result.error) dagster._core.errors.DagsterUserCodeProcessError: dagster._core.errors.SensorExecutionError: Error occurred during the execution of evaluation_fn for sensor trigger_dbt_job_after_etl Stack Trace: File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_grpc/impl.py", line 400, in get_external_sensor_execution return sensor_def.evaluate_tick(sensor_context) File "/usr/lib64/python3.9/contextlib.py", line 137, in __exit__ self.gen.throw(typ, value, traceback) File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_core/errors.py", line 297, in user_code_error_boundary raise new_error from e The above exception was caused by the following exception: Exception: Error in sensor trigger_dbt_job_after_etl: Sensor evaluation function returned a RunRequest for a sensor lacking a specified target (job_name, job, or jobs). Targets can be specified by providing job, jobs, or job_name to the @sensor decorator. Stack Trace: File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_core/errors.py", line 287, in user_code_error_boundary yield File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_grpc/impl.py", line 400, in get_external_sensor_execution return sensor_def.evaluate_tick(sensor_context) File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_core/definitions/sensor_definition.py", line 921, in evaluate_tick *self.resolve_run_requests( File "/home/ec2-user/.local/lib/python3.9/site-packages/dagster/_core/definitions/sensor_definition.py", line 965, in resolve_run_requests raise Exception(
日志显示生成的RunRequest已正确指定目标:
Created RunRequest: RunRequest(run_key=None, run_config={'ops': {'run_dbt_task': {'config': {'cluster_name': 'dbt-', 'launch_type': 'FARGATE', 'security_groups': ['sg-'], 'subnets': ['subnet-', 'subnet-', 'subnet-', 'subnet-', 'subnet-'], 'task_definition': 'arn:aws:ecs:us-east-1:'}}}}, tags={}, job_name='dbt_job', asset_selection=None, stale_assets_only=False, partition_key=None, asset_check_keys=None, asset_graph_subset=None)
原因分析
Dagster的@run_status_sensor装饰器要求显式声明传感器允许触发的目标作业,尽管RunRequest中指定了job和job_name,但传感器本身未在装饰器层面配置可触发的目标,导致系统校验失败。
解决方法
在@run_status_sensor装饰器中添加job参数,明确指定该传感器可触发的作业:
修改后的传感器代码:
from dagster import run_status_sensor, RunRequest, SkipReason, DagsterRunStatus from analytics.jobs.etl_job import dbt_job # 新增job参数,指定传感器可触发的目标作业 @run_status_sensor(run_status=DagsterRunStatus.SUCCESS, job=dbt_job) def trigger_dbt_job_after_etl(context): context.log.info(f"Sensor triggered for job: {context.dagster_run.job_name}") if context.dagster_run.job_name == "etl_job": context.log.info("etl_job completed successfully. Preparing to trigger dbt_job.") run_config = { "ops": { "run_dbt_task": { "config": { "cluster_name": "", "launch_type": "", "security_groups": [""], "subnets": [ "subnet-", "subnet-", "subnet-", "subnet-", "subnet-" ], "task_definition": "" } } } } # 此处job参数可保留也可省略,因为装饰器已指定目标 run_request = RunRequest(run_key=None, run_config=run_config, job=dbt_job, job_name="dbt_job") context.log.info(f"Created RunRequest: {run_request}") return run_request else: context.log.info(f"Skipping sensor run because the job is not etl_job. Current job: {context.dagster_run.job_name}") return SkipReason("The job is not etl_job, skipping.")
如果需要触发多个作业,可以使用jobs=[dbt_job, ...]参数替代job=dbt_job。
内容的提问来源于stack exchange,提问作者Wesson
相关产品推荐
相关产品推荐

