如何配置Dagster向Prometheus推送任务状态与时长指标?
问题
我正在尝试为Dagster中运行的任务配置监控,希望通过Prometheus采集以下指标:
- 任务运行成功或失败状态
- 任务运行时长
我已搭建Prometheus Pushgateway,对Prometheus有基础了解,但无法正确配置Dagster将上述指标推送至Pushgateway。我尝试用自定义传感器监控任务状态和时长,使用prometheus_client库定义并推送指标,现有代码如下:
from dagster_prometheus import PrometheusResource from dagster import DagsterRunStatus, RunRequest, op, run_status_sensor, SkipReason from orchestration.jobs import all_dbt_assets_job from prometheus_client import CollectorRegistry, Gauge, push_to_gateway from dagster import DagsterInstance # Prometheus configuration PUSHGATEWAY_URL = "http://localhost:9091" registry = CollectorRegistry() # Define Prometheus metrics run_status_gauge = Gauge('dagster_run_status', 'Status of Dagster runs', ['job_name', 'status'], registry=registry) run_duration_gauge = Gauge('dagster_run_duration_seconds', 'Duration of Dagster runs in seconds', ['job_name'], registry=registry) @op def prometheus_op(prometheus: PrometheusResource): push_to_gateway(PUSHGATEWAY_URL, job=str(all_dbt_assets_job), registry=prometheus.registry) @run_status_sensor( run_status=DagsterRunStatus.SUCCESS, monitored_jobs=[all_dbt_assets_job], ) def report_status_sensor(context): instance = DagsterInstance.get() # Ensure we're not monitoring the status reporting job itself if context.dagster_run.job_name == 'status_reporting_job': return SkipReason("Don't report status of status_reporting_job") # Get the run details run = instance.get_run_by_id(context.dagster_run.run_id) job_name = run.pipeline_name status = 'success' if run.status == DagsterRunStatus.SUCCESS else 'failure' duration = (run.end_time - run.start_time).total_seconds() if run.end_time and run.start_time else 0 # Update Prometheus metrics run_status_gauge.labels(job_name=job_name, status=status).set(1 if status == 'success' else 0) run_duration_gauge.labels(job_name=job_name).set(duration) # Push metrics to Prometheus Pushgateway push_to_gateway(PUSHGATEWAY_URL, job=job_name, registry=registry) run_config = { "ops": { "status_report": {"config": {"job_name": job_name}} } } return RunRequest(run_key=None, run_config=run_config)
问题分析与修正方案
现有代码存在以下几个核心问题:
- 传感器仅监听
SUCCESS状态,失败任务的指标无法上报 - 重复获取DagsterRun实例,
context.dagster_run已包含所需信息 - 冗余的
RunRequest返回逻辑(指向未定义的status_reportop)会触发无效任务 - 未使用的
prometheus_op与自定义指标注册表逻辑冲突 - 状态指标的标签设置逻辑不清晰,易导致Prometheus维度混乱
修正后的代码
from dagster import DagsterRunStatus, run_status_sensor, SkipReason from orchestration.jobs import all_dbt_assets_job from prometheus_client import CollectorRegistry, Gauge, push_to_gateway # Prometheus Pushgateway配置 PUSHGATEWAY_URL = "http://localhost:9091" registry = CollectorRegistry() # 定义Prometheus指标 # 任务状态:1表示对应状态触发,0表示未触发 run_status_gauge = Gauge( 'dagster_run_status', 'Dagster任务运行状态', ['job_name', 'status'], registry=registry ) # 任务运行时长:仅在任务结束时更新 run_duration_gauge = Gauge( 'dagster_run_duration_seconds', 'Dagster任务运行时长(秒)', ['job_name'], registry=registry ) @run_status_sensor( # 同时监听成功和失败状态 run_status=[DagsterRunStatus.SUCCESS, DagsterRunStatus.FAILURE], monitored_jobs=[all_dbt_assets_job], ) def report_dagster_run_metrics(context): # 跳过监控自身(如果有专门的上报任务) if context.dagster_run.job_name == 'status_reporting_job': return SkipReason("跳过监控状态上报任务本身") # 直接从context获取运行信息,无需重复查询 job_name = context.dagster_run.job_name run_status = context.dagster_run.status # 设置状态指标:明确区分成功/失败标签的取值 if run_status == DagsterRunStatus.SUCCESS: run_status_gauge.labels(job_name=job_name, status='success').set(1) run_status_gauge.labels(job_name=job_name, status='failure').set(0) elif run_status == DagsterRunStatus.FAILURE: run_status_gauge.labels(job_name=job_name, status='success').set(0) run_status_gauge.labels(job_name=job_name, status='failure').set(1) # 计算并设置运行时长(增加非空校验避免报错) if context.dagster_run.start_time and context.dagster_run.end_time: duration = (context.dagster_run.end_time - context.dagster_run.start_time).total_seconds() run_duration_gauge.labels(job_name=job_name).set(duration) # 推送指标到Pushgateway push_to_gateway(PUSHGATEWAY_URL, job=job_name, registry=registry)
关键优化点
- 扩展监听范围:同时监听
SUCCESS和FAILURE状态,确保所有任务结束状态都能上报 - 简化数据获取:直接使用
context.dagster_run提供的字段,避免不必要的实例查询 - 明确指标维度:为每个状态标签设置明确的0/1值,方便Prometheus进行聚合查询
- 移除冗余逻辑:删除未使用的
prometheus_op和无效的RunRequest,避免触发错误任务 - 增加校验逻辑:对任务的开始/结束时间做非空校验,防止计算时长时出现异常
内容的提问来源于stack exchange,提问作者AERL
相关产品推荐
相关产品推荐

