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

如何在Airflow中发布OpenTelemetry自定义指标?

Airflow中OpenTelemetry自定义指标无法上报问题

环境背景

  • 本地Docker容器部署Airflow 2.9.2、OTel Collector、Prometheus和Grafana
  • Airflow原生指标可正常上报至OTel Collector,并在Grafana中展示,确认环境配置无问题

编写的DAG代码

from datetime import datetime
import logging
import os
from random import randint

from airflow import DAG
from airflow.operators.python import PythonOperator

from opentelemetry import metrics
from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter
from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 7, 1),
    'email': [],
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 0,
    'retry_delay': 30,
}

otel_host = os.getenv('AIRFLOW__METRICS__OTEL_HOST', 'localhost')
otel_port = os.getenv('AIRFLOW__METRICS__OTEL_PORT', '24317')
endpoint = f'http://{otel_host}:{otel_port}'

# Set up the Metric Exporter
exporter = OTLPMetricExporter(endpoint=endpoint, insecure=True)

# Set up the Metric Reader
reader = PeriodicExportingMetricReader(exporter)

# Set up the Meter Provider
provider = MeterProvider(metric_readers=[reader])
metrics.set_meter_provider(provider)

# Create a Meter
meter = metrics.get_meter(__name__)

# Create a custom metric
gauge = meter.create_gauge(
    name='hanks_gauge_from_airflow',
    description=f'hanks custom gauge from airflow',
    unit='int',
)


def publish_custom_metric(**context):
    i = randint(1, 100)
    ts = datetime.now().strftime('%Y-%m-%dT%H:%M:%S.%f')
    attributes = {'ts': ts}

    # Record a value to the gauge
    gauge.set(i, attributes)

    logging.info(i)
    logging.info(ts)
    logging.info(endpoint)

    return


with DAG(
    dag_id='hanks_otel_dag',
    default_args=default_args,
    schedule='* * * * *',
    catchup=False,
) as dag:

    # 给任务起短名称是因为OTel日志会报警告:任务名超过63字符限制将被截断
    ht1 = PythonOperator(
        task_id='ht1',
        python_callable=publish_custom_metric,
    )

问题现象

  • 本地单元测试运行该DAG时,自定义指标可成功上报至OTel Collector并在Grafana中显示
  • 在Airflow集群中运行DAG时,指标未上报,但DAG执行成功无任何报错
  • 尝试调用reader.force_flush()和provider.force_flush()强制刷新指标,均无效

补充更新信息

1. Exporter类型与端口不匹配问题

原代码使用了gRPC类型的Exporter,但Airflow配置的是HTTP端口,正确配置方式如下:

HTTP方式配置

from opentelemetry.exporter.otlp.proto.http.metric_exporter import OTLPMetricExporter
exporter = OTLPMetricExporter(endpoint=f'http://{otel_host}:{http_port}/v1/metrics')

gRPC方式配置

from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter
exporter = OTLPMetricExporter(endpoint=f'http://{otel_host}:{grpc_port}')

2. 端点连通性验证

为确认Airflow任务与OTel Collector端点的连通性,在DAG中添加了以下BashOperator任务:

test_endpoint = BashOperator(
    task_id='test_endpoint',
    bash_command=f"curl -v -d {} -H 'Content-Type: application/json' {endpoint}"
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 01:12:20