OpenTelemetry Collector processorhelper 内部遥测指标详解:incoming/duration/outgoing 三件套的采集机制与观测实践
2026/9/17 2:30:55 网站建设 项目流程

OpenTelemetry Collector processorhelper 内部遥测指标详解:incoming/duration/outgoing 三件套的采集机制与观测实践

【免费下载链接】opentelemetry-collectorOpenTelemetry Collector项目地址: https://gitcode.com/GitHub_Trending/op/opentelemetry-collector

processorhelper 是 OpenTelemetry Collector 提供给所有 processor(处理器)组件的通用封装层,它在转发数据、管理生命周期之外,还自动为每个处理器暴露一组标准化内部遥测指标。本文以 processorhelper/documentation.md 为核心,结合 metadata.yaml、obsreport.go 等源码,完整解析otelcol_processor_incoming_itemsotelcol_processor_internal_durationotelcol_processor_outgoing_items三项指标的含义、单位、语义与采集链路,并说明如何通过服务级遥测级别控制其导出,帮助你在实际部署中利用这些指标观测处理器吞吐与性能。

一、指标总览:processorhelper 暴露的 3 项内部遥测

processorhelper 是 OpenTelemetry Collector 中编写 processor 的脚手架包(稳定级别为 beta,覆盖 traces/metrics/logs 三种信号,见 metadata.yaml)。凡是通过NewTracesNewMetricsNewLogs三个构造器创建的自定义处理器,都会自动接入内部遥测,无需在业务代码中手动打点。

根据官方生成的 documentation.md,该组件共暴露以下 3 项指标,全部注册在 Meter 作用域go.opentelemetry.io/collector/processor/processorhelper下(见 internal/metadata/generated_telemetry.go):

1.otelcol_processor_incoming_items

  • 描述:传入处理器的数据项数量(Number of items passed to the processor)。
  • 指标属性
UnitMetric TypeValue TypeMonotonicStability
{item}SumInttrueAlpha

该指标是单调递增的 Int 类型 Counter,以“个数据项”为计量单位,用于统计进入处理器的数据规模。

2.otelcol_processor_internal_duration

  • 描述:处理器处理一批遥测数据所消耗的时间(Duration of time taken to process a batch of telemetry data through the processor)。
  • 指标属性
UnitMetric TypeValue TypeStability
sHistogramDoubleAlpha

该指标是秒为单位的直方图,反映处理器内部单批处理耗时分布,可用于观察 P50/P95/P99 延迟。

3.otelcol_processor_outgoing_items

  • 描述:处理器发出的数据项数量(Number of items emitted from the processor)。
  • 指标属性
UnitMetric TypeValue TypeMonotonicStability
{item}SumInttrueAlpha

同样是单调递增的 Int Counter,统计处理器实际输出、向后继续传递的数据项数量。

值得注意的是:三项指标本身的稳定级别为Alpha,意味着其名称、语义和属性在未来版本中仍可能调整;而processorhelper包整体(beta)与具体处理器组件(如 memorylimiterprocessor、batchprocessor)的稳定级别相互独立。

二、指标元数据:从 mdatagen 声明到生成代码

这三项指标并非手写埋点,而是通过 mdatagen 工具从声明式配置自动生成。核心元数据位于 metadata.yaml:

telemetry: metrics: processor_incoming_items: enabled: true stability: alpha description: Number of items passed to the processor. unit: "{item}" sum: value_type: int monotonic: true processor_internal_duration: enabled: true stability: alpha description: Duration of time taken to process a batch of telemetry data through the processor. unit: s histogram: async: false value_type: double processor_outgoing_items: enabled: true stability: alpha description: Number of items emitted from the processor. unit: "{item}" sum: value_type: int monotonic: true

在源码根目录通过go:generate mdatagen metadata.yaml指令(见 processor.go 文件头)生成 internal/metadata/generated_telemetry.go。生成代码中,三个指标被封装进TelemetryBuilder

  • ProcessorIncomingItemsmetric.Int64Counter,注册名otelcol_processor_incoming_items,单位{item}
  • ProcessorInternalDurationmetric.Float64Histogram,注册名otelcol_processor_internal_duration,单位s
  • ProcessorOutgoingItemsmetric.Int64Counter,注册名otelcol_processor_outgoing_items,单位{item}

这正是文档表格中 Metric Type、Value Type 与 Monotonic 标记在代码层面的直接对应:两个 Sum(Counter)对应Int64Counter,一个 Histogram 对应Float64Histogram

三、采集原理:obsreport 报告器如何打点

指标的实际写入集中在 obsreport.go。newObsReport创建报告器时,为所有指标统一附加了一组标签属性:

otelAttrs: metric.WithAttributeSet(attribute.NewSet( attribute.String(internal.ProcessorKey, set.ID.String()), // 属性键 processor attribute.String(signalKey, signal.String()), // 属性键 otel.signal )),

即每条数据序列都带有两个固定维度:

属性键取值含义
processor处理器组件 ID(如batchmemory_limiter区分具体是哪个处理器
otel.signaltraces/metrics/logs区分遥测信号类型

记录逻辑分两个方法:

func (or *obsReport) recordInOut(ctx context.Context, incoming, outgoing int) { or.telemetryBuilder.ProcessorIncomingItems.Add(ctx, int64(incoming), or.otelAttrs) or.telemetryBuilder.ProcessorOutgoingItems.Add(ctx, int64(outgoing), or.otelAttrs) } func (or *obsReport) recordInternalDuration(ctx context.Context, startTime time.Time) { duration := time.Since(startTime) or.telemetryBuilder.ProcessorInternalDuration.Record(ctx, duration.Seconds(), or.otelAttrs) }

由此可以确认三条重要事实:

  1. 三个指标必须带上processorotel.signal两个属性查询才有意义,否则无法区分是哪个处理器的数据;
  2. incoming_itemsoutgoing_items成对记录,一次消费调用同时累加两者;
  3. internal_duration记录的是time.Since(startTime)的秒数,startTime取自处理函数调用前一刻。

四、三种信号下的“数据项”计量语义

“数据项”(item)在不同信号下代表不同的最小单位,这是理解该指标的关键。三个构造器(traces.go、metrics.go、logs.go)分别通过 pdata 的计数方法统计输入输出规模:

信号构造器输入计数方法输出计数方法“数据项”含义
tracesNewTracestd.SpanCount()td.SpanCount()Span 数
metricsNewMetricsmd.DataPointCount()md.DataPointCount()数据点(DataPoint)数
logsNewLogsld.LogRecordCount()ld.LogRecordCount()LogRecord 数

以 traces 为例,traces.go 中的核心处理流程为:

spansIn := td.SpanCount() // 处理前统计输入 Span 数 td, errFunc = tracesFunc(ctx, td) // 调用用户处理逻辑 obs.recordInternalDuration(ctx, startTime) if errFunc != nil { obs.recordInOut(ctx, spansIn, 0) // 出错时输出记 0 ... return errFunc } spansOut := td.SpanCount() // 处理后统计输出 Span 数 obs.recordInOut(ctx, spansIn, spansOut) return nextConsumer.ConsumeTraces(ctx, td)

这意味着incoming_itemsoutgoing_items的差值可以直观反映处理器“丢弃/合并”了多少数据(例如采样、去重、聚合类处理器天然会造成差值)。

五、错误处理与 sentinel:ErrSkipProcessingData 的特殊语义

处理器返回错误时,指标记录遵循明确规则:

  • 处理函数返回普通错误recordInOut(ctx, in, 0),即outgoing_items记 0,错误继续向上传播;
  • 处理函数返回哨兵错误ErrSkipProcessingData:同样输出记 0,但错误被吞掉、不向上传播,表现为数据被“静默丢弃”。

ErrSkipProcessingData定义在 processor.go:

var ErrSkipProcessingData = errors.New("sentinel error to skip processing data from the remainder of the pipeline")

其设计意图是:当处理器判定某批数据无关紧要(如时间戳过期、内容不符合过滤条件)时,可以主动丢弃且不污染上层日志——这也是 filter 等处理器“丢弃数据”的标准做法。观测时注意:这类被丢弃的数据会计入incoming_items而不再计入outgoing_items,两者差值即是被 sentinel 静默丢弃的量。

六、观测实践:如何查询这三个指标

processorhelper 的指标经 Collector 自身遥测(service::telemetry)导出,可被 Prometheus 抓取。典型查询方式:

# 各处理器处理吞吐(每秒) rate(otelcol_processor_incoming_items[1m]) rate(otelcol_processor_outgoing_items[1m]) # 各处理器处理延迟分布 histogram_quantile(0.95, sum(rate(otelcol_processor_internal_duration_bucket[5m])) by (le, processor, otel.signal)) # 丢弃率(outgoing 与 incoming 的比值) otelcol_processor_outgoing_items / otelcol_processor_incoming_items

重要前提——遥测级别otelcol_processor_internal_duration并非在所有遥测级别下都会导出。在 service/internal/metricviews/views.go 中,DefaultViews定义了默认的指标视图裁剪规则:

if level < configtelemetry.LevelDetailed { // Drop duration metric if the level is not detailed dropViewOption(&config.ViewSelector{ MeterName: new("go.opentelemetry.io/collector/processor/processorhelper"), InstrumentName: new("otelcol_processor_internal_duration"), }), }

也就是说,只有将 Collector 的遥测级别配置为detailed(通过service::telemetry::metrics::level--metrics-level=detailed)时,otelcol_processor_internal_duration才会被导出;incoming_itemsoutgoing_items两个计数指标则不受该限制。若发现延迟直方图缺失,请优先检查遥测级别配置。

七、并发安全与测试验证

processorhelper 的指标记录是并发安全的:TelemetryBuilder内部使用sync.Mutex保护异步仪器注册(见 generated_telemetry.go),而 Counter/Histogram 本身即线程安全类型。

对应测试集中在 metrics_test.go:

  • TestMetricsConcurrency:10 个 goroutine 各并发消费 10000 批数据,验证高并发下无竞态;
  • TestMetrics_RecordInOut:输入 2 个 DataPoint、处理函数输出 3 个,断言incoming_items=2outgoing_items=3,并验证属性为processor="processorhelper"otel.signal="metrics"
  • TestMetrics_RecordIn_ErrorOut:处理返回错误时断言incoming_items=2outgoing_items=0,印证第五节“出错输出记 0”的语义;
  • TestMetrics_ProcessInternalDuration:断言直方图Count=1,验证每次消费都记录一次耗时分布。

如果你编写自定义处理器并依赖 processorhelper,这些测试同时是processor/otel.signal属性值以及计数语义的最佳参考样例。

八、小结

otelcol_processor_incoming_itemsotelcol_processor_internal_durationotelcol_processor_outgoing_items构成了 processorhelper 统一的处理器可观测性模型:

  • 吞吐维度由两个单调 Counter 提供,配合processorotel.signal属性可精确到“某个处理器、某类信号”的进出流量;
  • 性能维度由秒级直方图提供,但需以detailed遥测级别为前提;
  • 计量口径按信号区分(Span / DataPoint / LogRecord),解读差值时应结合处理器业务语义(过滤、采样、聚合);
  • 所有指标当前稳定级别为 Alpha,接口演进需关注版本升级说明(见 CHANGELOG.md 中相关记录)。

建议在部署 Collector 后,将上述 PromQL 查询接入统一监控看板,即可获得每个处理器的实时吞吐、丢弃量与延迟画像,为容量评估和调优提供数据支撑。

【免费下载链接】opentelemetry-collectorOpenTelemetry Collector项目地址: https://gitcode.com/GitHub_Trending/op/opentelemetry-collector

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询