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

如何将现有Prometheus指标推送(而非拉取)至OTLP Collector?

实现已有Prometheus指标转OTLP推送至Collector的方案

针对你已经用prometheus.CounterVec等Prometheus客户端创建的指标,要直接转换为OTLP格式实时推送,有两种可行方案,分别适合不同场景:

方案一:代码层直接转换推送(无额外组件)

核心思路是读取Prometheus注册表中的指标样本,转换为OTLP Metric格式后,通过OTLP客户端发送到Collector。

步骤1:复用现有Prometheus注册表

你代码里应该已经将指标注册到默认注册表或自定义注册表,先获取这个注册表实例:

// 若使用默认注册表
reg := prometheus.DefaultRegisterer.(*prometheus.Registry)

// 若使用自定义注册表
// reg := prometheus.NewRegistry()
// prometheus.MustRegister(requestCounter) // 注册你的CounterVec等指标

步骤2:编写Prometheus到OTLP的转换逻辑

实现函数将prometheus.Metric转换为OTLP的metricpb.Metric,覆盖常见的Counter、Gauge类型:

import (
    "github.com/prometheus/client_golang/prometheus"
    "google.golang.org/protobuf/types/known/timestamppb"
    metricpb "go.opentelemetry.io/proto/otlp/metrics/v1"
)

func convertPromToOTLP(promMetric prometheus.Metric) (*metricpb.Metric, error) {
    desc := promMetric.Desc()
    otlpMetric := &metricpb.Metric{
        Name:        desc.Name(),
        Description: desc.Help(),
    }

    samples, err := promMetric.Write()
    if err != nil {
        return nil, err
    }

    for _, sample := range samples {
        // 转换Prometheus标签为OTLP属性
        attrs := make([]*metricpb.KeyValue, 0, len(sample.Metric))
        for k, v := range sample.Metric {
            attrs = append(attrs, &metricpb.KeyValue{
                Key:   k,
                Value: &metricpb.AnyValue{StringValue: v},
            })
        }

        // 根据指标类型映射到OTLP对应类型
        switch desc.Type() {
        case prometheus.CounterValue:
            sum := &metricpb.Sum{
                IsMonotonic: true,
                AggregationTemporality: metricpb.AggregationTemporality_AGGREGATION_TEMPORALITY_CUMULATIVE,
                DataPoints: []*metricpb.NumberDataPoint{{
                    Attributes:    attrs,
                    TimeUnixNano:  timestamppb.New(sample.Timestamp).AsTime().UnixNano(),
                    AsDouble:      sample.Value,
                }},
            }
            otlpMetric.Data = &metricpb.Metric_Sum{Sum: sum}
        case prometheus.GaugeValue:
            gauge := &metricpb.Gauge{
                DataPoints: []*metricpb.NumberDataPoint{{
                    Attributes:    attrs,
                    TimeUnixNano:  timestamppb.New(sample.Timestamp).AsTime().UnixNano(),
                    AsDouble:      sample.Value,
                }},
            }
            otlpMetric.Data = &metricpb.Metric_Gauge{Gauge: gauge}
        // 如需支持Histogram、Summary,可在此扩展转换逻辑
        }
    }

    return otlpMetric, nil
}

步骤3:定时采集并发送到OTLP Collector

使用OTLP gRPC exporter定时从注册表拉取指标,转换后发送:

import (
    "context"
    "time"

    "go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc"
    "go.opentelemetry.io/otel/sdk/resource"
    semconv "go.opentelemetry.io/otel/semconv/v1.24.0"
)

func startOTLPPusher(ctx context.Context, reg *prometheus.Registry, collectorAddr string) error {
    // 创建OTLP exporter
    exporter, err := otlpmetricgrpc.New(ctx,
        otlpmetricgrpc.WithEndpoint(collectorAddr),
        otlpmetricgrpc.WithInsecure(), // 生产环境请启用TLS
    )
    if err != nil {
        return err
    }

    // 配置服务资源信息
    res, err := resource.New(ctx, resource.WithAttributes(semconv.ServiceName("your-service")))
    if err != nil {
        return err
    }

    // 定时拉取并发送
    ticker := time.NewTicker(1 * time.Second) // 可根据实时性要求调整间隔
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return exporter.Shutdown(ctx)
        case <-ticker.C:
            // 从注册表获取所有指标
            promMetrics, err := reg.Gather()
            if err != nil {
                continue // 可添加日志记录错误
            }

            // 批量转换为OTLP格式
            otlpMetrics := make([]*metricpb.Metric, 0, len(promMetrics))
            for _, pm := range promMetrics {
                om, err := convertPromToOTLP(pm)
                if err == nil {
                    otlpMetrics = append(otlpMetrics, om)
                }
            }

            // 构造OTLP请求并发送
            req := &metricpb.ExportMetricsServiceRequest{
                ResourceMetrics: []*metricpb.ResourceMetrics{{
                    Resource: res.Resource(),
                    ScopeMetrics: []*metricpb.ScopeMetrics{{
                        Scope:   &metricpb.InstrumentationScope{Name: "prometheus-to-otlp"},
                        Metrics: otlpMetrics,
                    }},
                }},
            }

            if err := exporter.Export(ctx, req); err != nil {
                // 处理发送错误,如重试或日志
            }
        }
    }
}

最后在服务启动时启动推送协程:

func main() {
    ctx := context.Background()
    reg := prometheus.DefaultRegisterer.(*prometheus.Registry)

    go func() {
        if err := startOTLPPusher(ctx, reg, "localhost:4317"); err != nil {
            // 处理启动错误
        }
    }()

    // 启动你的服务逻辑...
}

方案二:通过OTel Collector中转(无需修改服务代码)

如果不想改动服务代码,可部署OpenTelemetry Collector,配置它采集你的服务/metrics端点,再转发为OTLP格式到目标Collector。

Collector配置示例

receivers:
  prometheus:
    config:
      scrape_configs:
        - job_name: "your-service"
          scrape_interval: 1s # 实时性要求高则缩短间隔
          static_configs:
            - targets: ["your-service:8080"] # 你的服务地址和端口

exporters:
  otlp:
    endpoint: "target-collector:4317"
    tls:
      insecure: true # 生产环境禁用

service:
  pipelines:
    metrics:
      receivers: [prometheus]
      exporters: [otlp]

方案对比

  • 代码转换方案:实时性更强(可按需触发或高频定时),无额外组件依赖,但需要维护转换逻辑,适合对延迟敏感的场景。
  • Collector中转方案:零代码侵入,配置简单,适合快速落地,但实时性受采集间隔限制,多了一层转发环节。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 20:55:44