如何在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
相关产品推荐
相关产品推荐

