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_items、otelcol_processor_internal_duration、otelcol_processor_outgoing_items三项指标的含义、单位、语义与采集链路,并说明如何通过服务级遥测级别控制其导出,帮助你在实际部署中利用这些指标观测处理器吞吐与性能。
一、指标总览:processorhelper 暴露的 3 项内部遥测
processorhelper 是 OpenTelemetry Collector 中编写 processor 的脚手架包(稳定级别为 beta,覆盖 traces/metrics/logs 三种信号,见 metadata.yaml)。凡是通过NewTraces、NewMetrics、NewLogs三个构造器创建的自定义处理器,都会自动接入内部遥测,无需在业务代码中手动打点。
根据官方生成的 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)。
- 指标属性:
| Unit | Metric Type | Value Type | Monotonic | Stability |
|---|---|---|---|---|
| {item} | Sum | Int | true | Alpha |
该指标是单调递增的 Int 类型 Counter,以“个数据项”为计量单位,用于统计进入处理器的数据规模。
2.otelcol_processor_internal_duration
- 描述:处理器处理一批遥测数据所消耗的时间(Duration of time taken to process a batch of telemetry data through the processor)。
- 指标属性:
| Unit | Metric Type | Value Type | Stability |
|---|---|---|---|
| s | Histogram | Double | Alpha |
该指标是秒为单位的直方图,反映处理器内部单批处理耗时分布,可用于观察 P50/P95/P99 延迟。
3.otelcol_processor_outgoing_items
- 描述:处理器发出的数据项数量(Number of items emitted from the processor)。
- 指标属性:
| Unit | Metric Type | Value Type | Monotonic | Stability |
|---|---|---|---|---|
| {item} | Sum | Int | true | Alpha |
同样是单调递增的 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:
ProcessorIncomingItems:metric.Int64Counter,注册名otelcol_processor_incoming_items,单位{item};ProcessorInternalDuration:metric.Float64Histogram,注册名otelcol_processor_internal_duration,单位s;ProcessorOutgoingItems:metric.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(如batch、memory_limiter) | 区分具体是哪个处理器 |
otel.signal | traces/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) }由此可以确认三条重要事实:
- 三个指标必须带上
processor与otel.signal两个属性查询才有意义,否则无法区分是哪个处理器的数据; incoming_items与outgoing_items成对记录,一次消费调用同时累加两者;internal_duration记录的是time.Since(startTime)的秒数,startTime取自处理函数调用前一刻。
四、三种信号下的“数据项”计量语义
“数据项”(item)在不同信号下代表不同的最小单位,这是理解该指标的关键。三个构造器(traces.go、metrics.go、logs.go)分别通过 pdata 的计数方法统计输入输出规模:
| 信号 | 构造器 | 输入计数方法 | 输出计数方法 | “数据项”含义 |
|---|---|---|---|---|
| traces | NewTraces | td.SpanCount() | td.SpanCount() | Span 数 |
| metrics | NewMetrics | md.DataPointCount() | md.DataPointCount() | 数据点(DataPoint)数 |
| logs | NewLogs | ld.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_items与outgoing_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_items与outgoing_items两个计数指标则不受该限制。若发现延迟直方图缺失,请优先检查遥测级别配置。
七、并发安全与测试验证
processorhelper 的指标记录是并发安全的:TelemetryBuilder内部使用sync.Mutex保护异步仪器注册(见 generated_telemetry.go),而 Counter/Histogram 本身即线程安全类型。
对应测试集中在 metrics_test.go:
TestMetricsConcurrency:10 个 goroutine 各并发消费 10000 批数据,验证高并发下无竞态;TestMetrics_RecordInOut:输入 2 个 DataPoint、处理函数输出 3 个,断言incoming_items=2、outgoing_items=3,并验证属性为processor="processorhelper"、otel.signal="metrics";TestMetrics_RecordIn_ErrorOut:处理返回错误时断言incoming_items=2、outgoing_items=0,印证第五节“出错输出记 0”的语义;TestMetrics_ProcessInternalDuration:断言直方图Count=1,验证每次消费都记录一次耗时分布。
如果你编写自定义处理器并依赖 processorhelper,这些测试同时是processor/otel.signal属性值以及计数语义的最佳参考样例。
八、小结
otelcol_processor_incoming_items、otelcol_processor_internal_duration、otelcol_processor_outgoing_items构成了 processorhelper 统一的处理器可观测性模型:
- 吞吐维度由两个单调 Counter 提供,配合
processor、otel.signal属性可精确到“某个处理器、某类信号”的进出流量; - 性能维度由秒级直方图提供,但需以
detailed遥测级别为前提; - 计量口径按信号区分(Span / DataPoint / LogRecord),解读差值时应结合处理器业务语义(过滤、采样、聚合);
- 所有指标当前稳定级别为 Alpha,接口演进需关注版本升级说明(见 CHANGELOG.md 中相关记录)。
建议在部署 Collector 后,将上述 PromQL 查询接入统一监控看板,即可获得每个处理器的实时吞吐、丢弃量与延迟画像,为容量评估和调优提供数据支撑。
【免费下载链接】opentelemetry-collectorOpenTelemetry Collector项目地址: https://gitcode.com/GitHub_Trending/op/opentelemetry-collector
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考