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

如何配置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)

问题分析与修正方案

现有代码存在以下几个核心问题:

  1. 传感器仅监听SUCCESS状态,失败任务的指标无法上报
  2. 重复获取DagsterRun实例,context.dagster_run已包含所需信息
  3. 冗余的RunRequest返回逻辑(指向未定义的status_report op)会触发无效任务
  4. 未使用的prometheus_op与自定义指标注册表逻辑冲突
  5. 状态指标的标签设置逻辑不清晰,易导致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)

关键优化点

  1. 扩展监听范围:同时监听SUCCESS和FAILURE状态,确保所有任务结束状态都能上报
  2. 简化数据获取:直接使用context.dagster_run提供的字段,避免不必要的实例查询
  3. 明确指标维度:为每个状态标签设置明确的0/1值,方便Prometheus进行聚合查询
  4. 移除冗余逻辑:删除未使用的prometheus_op和无效的RunRequest,避免触发错误任务
  5. 增加校验逻辑:对任务的开始/结束时间做非空校验,防止计算时长时出现异常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 02:19:57