OpenTelemetry Collector 处理器(Processor)完全指南:推荐顺序、数据所有权模型与自定义开发实战
【免费下载链接】opentelemetry-collectorOpenTelemetry Collector项目地址: https://gitcode.com/GitHub_Trending/op/opentelemetry-collector
OpenTelemetry Collector 的处理器(Processor)是管道(Pipeline)中在接收器(Receiver)与导出器(Exporter)之间对遥测数据执行预处理、过滤、采样、变换与聚合的关键组件。本文基于 opentelemetry-collector 仓库的 processor/README.md,系统讲解处理器的核心工作机制——包括推荐处理器及其最佳放置顺序、管道中的数据所有权模型(独占/共享)、处理器排序原则,以及如何基于 processorhelper 开发自定义处理器,并结合仓库源码与配置文件提供可落地的实战指导。读完本文,你将能够正确编排处理器管道、理解数据何时可安全修改,并独立实现一个自定义处理器。
处理器(Processor)概述
在 OpenTelemetry Collector 中,Processor 被用于管道(Pipeline)的各个阶段。通常情况下,处理器在数据被导出之前对数据进行预处理,例如修改属性(Attribute)或进行采样(Sampling)。其核心定位是:
- 预处理:在导出前对
pdata.Traces、pdata.Metrics、pdata.Logs数据进行清洗、过滤、增强; - 降载:通过采样、限流等手段降低下游导出压力;
- 聚合:通过批处理(Batching)减少网络连接数与传输开销。
需要注意的是:默认情况下,Collector 不会启用任何处理器。处理器必须针对每一种数据源(signals)单独启用,并且并非所有处理器都支持所有数据源(例如 batch processor 同时支持 traces/metrics/logs,而部分采样处理器只支持 traces)。此外,处理器的顺序至关重要——同一个管道中处理器的声明顺序就是它们被应用的实际顺序。
仓库内置(core distribution)的处理器按字母排序有两个:
- Batch Processor(批处理器)
- Memory Limiter Processor(内存限制器)
除此之外,opentelemetry-collector-contrib)将其加入 Collector 发行版中。
推荐处理器及其最佳放置顺序
官方文档给出了一套“最佳实践顺序”,按此顺序配置可以最大化管道效率并降低数据丢失风险。其核心思想是:尽早丢弃无用的数据,仅在最后阶段对将真正导出的数据做批量聚合。
下面的顺序是官方推荐的最佳实践。具体每个处理器的配置请参考其各自的文档。
- memory_limiter —— 必须放在第一位:它必须在其他处理器累积无法刷出的数据之前进行降载(shed load)。如果 Collector 内存吃紧,它能最先向上游接收器施加背压(backpressure),最大限度避免 OOM。
- 任何会丢弃数据的采样或过滤处理器(如
filter、tailsampling、probabilisticsampler):尽早丢弃不需要的数据,避免对其做进一步处理,节省 CPU 与内存。 - 任何依赖从
Context中获取来源信息的处理器(例如k8sattributes):必须放在 batch 处理器之前运行,因为批处理会清空请求上下文(request context)。 - 任何对遥测数据进行变换或增强的处理器(如
attributes、transform、resource):只对真正会被导出的数据进行增强。 - batch —— 放在最后:在过滤与变换之后再做批处理,可以确保不会对即将被丢弃或后续还会被修改的数据进行批量缓存;同时,当导出器自身具备批处理能力时,优先使用导出器的批处理能力。
一个遵循最佳实践的管道示例(伪配置结构):
service: pipelines: traces: receivers: [otlp] processors: [memory_limiter, filter, k8sattributes, attributes, batch] exporters: [otlp]处理器排序为什么重要
处理器在管道中的声明顺序就是其执行顺序。将丢弃类处理器放在前面,将批处理放在最后,可以带来三重收益:
- 避免无效批处理:先过滤/变换再批处理,不会缓存那些最终会被丢弃或还会被二次修改的数据;
- 保住请求上下文:依赖
Context的处理器(如k8sattributes)必须在批处理之前执行,因为批处理器会清空请求上下文(这一点在 batch_processor.go 的实现中体现为按批重组数据而非透传原始请求上下文); - 提前降载:
memory_limiter放在第一位,保证内存告警时背压能第一时间传递到接收器,而不是先让下游处理器堆积无法刷出的数据。
管道中的数据所有权模型
数据所有权(Data Ownership)是理解 OpenTelemetry Collector 处理器行为的关键概念。pdata.Traces、pdata.Metrics和pdata.Logs数据在管道中流动时,其所有权也随之传递:
- 数据由**接收器(Receiver)**创建;
- 当第一个处理器的
ConsumeTraces/ConsumeMetrics/ConsumeLogs函数被调用时,所有权移交给该处理器; - 处理器处理完毕后,通过调用下一个处理器的
ConsumeTraces/ConsumeMetrics/ConsumeLogs函数把所有权传递给下一环节; - 依此类推,直到数据被导出。
注意:一个接收器可能被挂接到多个管道(pipeline),此时同一份数据会通过数据扇出连接器(fan-out connector)被传递给所有关联管道。这也正是“数据所有权模式”需要区分的原因。
所有权模式如何确定
所有权模式在启动期间(startup)根据处理器报告的数据修改意图(data modification intent)来决定:
- 每个处理器通过
Capabilities函数返回的结构体中的MutatesData字段声明其修改意图; - 如果管道中任一处理器声明要修改数据(
MutatesData: true),则该管道工作于独占所有权模式(Exclusive Ownership); - 此外,任何从某个已处于独占模式的管道所挂接的接收器获取数据的其他管道,也会被强制工作于独占所有权模式(因为共享的接收器数据必须被克隆后分发给多个管道)。
源码佐证:批处理器在 batch_processor.go 中声明MutatesData: true;而内存限制器在 factory.go 中声明processorCapabilities = consumer.Capabilities{MutatesData: false},因为它只读判断内存水位并拒绝数据,不修改数据本身。处理器接口定义在 processor.go,其中Traces/Metrics/Logs三个接口均由component.Component与对应的consumer接口组合而成。
独占所有权(Exclusive Ownership)
在独占所有权模式下,数据在某一时刻被某个处理器独占拥有,该处理器可以自由修改它拥有的数据。
要点:
- 适用场景:独占模式仅适用于从同一接收器接收数据的管道。如果一个管道被标记为独占模式,那么从共享接收器收到的数据会在扇出连接器处先被克隆,再分别传递给每个管道。这保证了每个管道拥有自己独占的数据副本,可以安全地在管道内进行修改。
- 所有权持续时间:处理器对数据的所有权从自身
ConsumeTraces/ConsumeMetrics/ConsumeLogs调用开始,直到它调用下一个处理器的对应 Consume 函数把所有权移交出去为止。移交之后,该处理器不得再读写这份数据,因为新所有者可能正在并发修改它。 - 实现红利:独占模式让需要修改数据的处理器只需声明修改意图即可轻松实现,无需自己处理并发与共享问题。例如 contrib 仓库中的
attributesprocessor就依赖这一机制。
fan-out 连接器的智能克隆逻辑可以在 internal/fanoutconsumer/traces.go 中看到:它会将数据克隆后发送给所有“需要修改数据”的消费者(最后一个除外,最后一个可直接使用原始可变数据),并在发送给多个只读消费者前将数据标记为只读(td.MarkReadOnly())。
共享所有权(Shared Ownership)
在共享所有权模式下,没有任何处理器拥有数据,任何处理器都不得修改共享数据:
- 挂接到多个管道的接收器,其扇出连接器不做克隆,所有关联管道看到的是同一份共享数据副本;
- 共享模式下管道中的处理器禁止修改通过
ConsumeTraces/ConsumeMetrics/ConsumeLogs接收到的原始数据,只能读取。
如果处理器在处理过程中确实需要修改数据,但又不希望承担独占模式带来的克隆成本,可以:
- 声明自己不修改数据(
MutatesData=false); - 采用**写时复制(copy-on-write)**等技术,只对
pdata.Traces/pdata.Metrics/pdata.Logs的个别子部分进行替换,而绝不改动传入的原始数据。
只要不修改传入的原始pdata对象,任何方案都是被允许的。通过将MutatesData=false,可以避免管道被标记为独占模式,从而避免数据克隆的开销。这正是内存限制处理器(只读判断、只返回错误)与批处理器(重组数据、声明可变)所展示的两种典型能力声明的差异。
自定义处理器的开发
要为 OpenTelemetry Collector 创建自定义处理器,通常需要做三件事:实现处理器接口、定义处理器配置、向 Collector 注册。完整流程包括创建 Factory、实现处理逻辑、处理配置选项。官方推荐的开发路径是使用processorhelper包,它提供了大量工具与模式来简化处理器开发。
第一步:定义配置结构体
每个处理器需要一个配置结构体,实现component.Config接口,并提供默认配置:
type Config struct { // 自定义字段,例如采样率、属性键名等 BatchSize int `mapstructure:"batch_size"` } func createDefaultConfig() component.Config { return &Config{ BatchSize: 100, // 提供合理的默认值 } }第二步:实现处理逻辑
处理逻辑就是实现一个函数,接收数据、处理后转发给下一个消费者。以 traces 为例:
func processTraces(ctx context.Context, td ptrace.Traces) (ptrace.Traces, error) { // 在这里读取/修改 td(取决于声明的能力) return td, nil }注意:是否允许修改传入的td取决于你通过WithCapabilities声明的MutatesData。默认情况下,processorhelper 的fromOptions会将能力初始化为consumer.Capabilities{MutatesData: true}(见 processor/processorhelper/processor.go),即默认声明“会修改数据”。
第三步:创建 Factory 并注册
利用processor.NewFactory(定义于 processor/processor.go)与processorhelper的NewTraces/NewMetrics/NewLogs构造器组合出完整处理器:
func NewFactory() processor.Factory { return processor.NewFactory( component.MustNewType("myprocessor"), createDefaultConfig, processor.WithTraces(createTraces, component.StabilityLevelBeta), ) } func createTraces( ctx context.Context, set processor.Settings, cfg component.Config, next consumer.Traces, ) (processor.Traces, error) { return processorhelper.NewTraces(ctx, set, cfg, next, processTraces, processorhelper.WithCapabilities(consumer.Capabilities{MutatesData: true}), ) }可用的processorhelper选项包括:
WithCapabilities:覆盖默认能力声明(默认MutatesData: true);WithStart/WithShutdown:覆盖默认的启动/关闭函数(默认空实现,返回 nil);ErrSkipProcessingData哨兵错误(见 processor/processorhelper/processor.go):处理器可返回它来“有意丢弃”数据,而不会把错误沿管道向上传播到日志中。
完成 Factory 后,通过 cmd/otelcorecol 或builder工具将其注册进 Collector 的自定义构建中即可。
参考:memory_limiter 的工程实践
内存限制器是一个极好的自定义处理器参考范本:factory.go 展示了:
- 用
xprocessor.NewFactory声明对 traces/metrics/logs/profiles 四种信号的支持; - 通过
processorhelper.WithCapabilities(processorCapabilities)(MutatesData: false)声明只读能力; - 通过
processorhelper.WithStart/WithShutdown注入限流器的启动与关闭生命周期; - 通过工厂级缓存(
memoryLimiters map[component.Config]*memoryLimiterProcessor)复用同一配置的限流实例,避免为每个管道重复运行内存检查与 GC。
两个内置处理器的配置速查
memory_limiter:防止 Collector 内存耗尽
memory_limiter 处理器用于防止 Collector 出现 OOM(Out of Memory)。它周期性地检查内存使用情况,当超过设定阈值时开始拒绝数据并强制 GC以降低内存消耗。它使用软限制(soft limit)与硬限制(hard limit)两级水位:
- 硬限制由
limit_mib(或limit_percentage)定义,始终大于等于软限制; - 软限制 = 硬限制 −
spike_limit_mib; - 内存超过软限制时进入受限模式:向上一环节(通常是接收器)的
ConsumeLogs/ConsumeTraces/ConsumeMetrics调用返回非永久性错误,接收器应重试发送并向上游数据源施加背压; - 内存超过硬限制时额外强制执行 GC,若 GC 无效则对强制 GC 做指数退避(由
max_gc_interval_when_soft_limited/max_gc_interval_when_hard_limited控制上限,默认 30s); - 内存回落到软限制以下后恢复正常操作。
最佳实践(详见 memorylimiterprocessor/README.md):
- 在每个 Collector 上同时配置
GOMEMLIMIT环境变量与 memory_limiter,GOMEMLIMIT建议设为 Collector 硬内存限制的80%; - memory_limiter必须作为管道中的第一个处理器,确保背压第一时间传递到接收器;
spike_limit_mib建议设为硬限制的20%,以保证单个检查间隔内内存增幅不会越过硬限制;- 容器化环境(支持 cgroup,如 Docker)优先用
limit_percentage;裸机/虚拟机环境且吞吐可预期时优先用limit_mib。
常用配置示例:
processors: memory_limiter: check_interval: 1s limit_mib: 4000 spike_limit_mib: 800- 硬限制为4000 MiB;
- 软限制为 4000 − 800 =3200 MiB。
按百分比配置(容器环境):
processors: memory_limiter: check_interval: 1s limit_percentage: 80 spike_limit_percentage: 15在总内存 1000 MiB 的机器上:硬限制为 800 MiB,软限制为 650 MiB。
注意:memory_limiter 返回的拒绝错误是非永久性的,接收器必须重试,否则数据会永久丢失。另外,在 memory_limiter 拒绝数据之前,入站数据仍可能先消耗额外的内存(尤其对非 OTLP 接收器),设置限制时要为这部分留出余量。
batch:批处理压缩与减少连接数
batch 处理器将 spans、metrics 或 logs 放入批次中统一发送,以更好地压缩数据并减少传输所需的出站连接数,同时支持基于大小与基于时间两种触发方式。它应配置在memory_limiter以及所有采样处理器之后(先丢弃,再批量)。
核心配置项(详见 batchprocessor/README.md):
| 配置项 | 默认值 | 说明 |
|---|---|---|
send_batch_size | 8192 | 达到该数量的 span/指标点/日志记录后立即发送批次(无论是否超时)。它只是触发值,不限制批次大小上限 |
timeout | 200ms | 达到该时长后无论批次大小都发送;设为0s时忽略send_batch_size,数据即时发送,仅受send_batch_max_size约束 |
send_batch_max_size | 0 | 批次大小的上限(0 表示不限制),保证更大的批次被拆成更小的单元;必须大于等于send_batch_size |
metadata_keys | 空 | 非空时,为client.Metadata中键值的每个不同组合创建一个独立的 batcher 实例(多租户批处理) |
metadata_cardinality_limit | 1000 | 当metadata_keys非空时,限制进程生命周期内可处理的键值组合数量上限 |
配置示例(默认 + 自定义):
processors: batch: batch/2: send_batch_size: 10000 timeout: 10sbatch/2将缓冲最多 10000 个 span/指标点/日志记录、最长 10 秒,且不拆分数据项(不强制批次大小上限)。
无人工延迟、但强制批次上限的配置:
processors: batch: send_batch_max_size: 10000 timeout: 0s多租户元数据批处理(需在接收器上启用include_metadata: true):
processors: batch: # 按 tenant-id 分组批处理 metadata_keys: - tenant_id # 限制最多 10 个 batcher,超过后报错 metadata_cardinality_limit: 10注意:每个不同的元数据组合都会在 Collector 中分配一个运行于整个进程生命周期的后台任务,每个任务持有最多send_batch_size条记录的待发批次,因此按元数据批处理会显著增加批处理相关的内存占用;建议配合 Auth 扩展校验相关元数据键的值。当前在用的批处理器数量通过otelcol_processor_batch_metadata_cardinality指标暴露。
总结
- 处理器默认不启用,需为每种数据源显式配置,且顺序即执行顺序;
- 推荐顺序为:
memory_limiter→ 采样/过滤 → 依赖Context的处理器 → 变换/增强 →batch; - 管道的所有权模式由处理器
Capabilities().MutatesData决定:任一处理器声明可变 → 独占模式(fan-out 处克隆数据);全部只读 → 共享模式(共享同一份数据,禁止修改); - 自定义处理器遵循“配置结构体 + 处理函数 + Factory”三步走,优先基于 processorhelper 开发,并通过
WithCapabilities诚实声明数据修改意图,以避免不必要的克隆开销; - 更多处理器可查阅 contrib 仓库,并通过自定义构建加入 Collector。
相关深入阅读:
- processor/README.md(本文核心依据)
- processor/batchprocessor/README.md
- processor/memorylimiterprocessor/README.md
- processor/processor.go(处理器 Factory 接口定义)
- processor/processorhelper(自定义处理器开发工具包)
- internal/fanoutconsumer/traces.go(fan-out 智能克隆实现)
【免费下载链接】opentelemetry-collectorOpenTelemetry Collector项目地址: https://gitcode.com/GitHub_Trending/op/opentelemetry-collector
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考