Grafana Tempo 中的 Kafka Receiver:基于 OpenTelemetry Collector Contrib 的 Kafka 遥测消费与管线元数据传递实践指南
【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempo
本文围绕 OpenTelemetry Collector Contrib 的kafkareceiver组件(本仓库以 vendor 形式托管其完整文档与实现)展开,讲解如何从 Kafka 消费 traces / metrics / logs / profiles 遥测数据,覆盖默认配置、topic 正则消费与排除、TLS/SASL 认证、消息元数据与 Header 传递、消息确认(marking)语义等核心机制,并结合 Grafana Tempo 仓库中的 Kafka 摄取(ingest)实现,说明 Kafka 消费模式在分布式追踪后端中的真实落地方式。读完本文,你将能独立配置一个可投入生产的 Kafka Receiver,并理解其底层消费与 ack 行为。
一、Kafka Receiver 是什么:从 Kafka 到 OTel 管线的入口组件
Kafka Receiver 是 OpenTelemetry Collector Contrib 中负责从 Kafka topic 消费遥测数据(traces、metrics、logs、profiles)并送入 Collector 下游管线的接收器(receiver)。它在 OTel 生态中的典型定位是作为数据管道的中转入口:上游采集器把遥测数据写入 Kafka 作为缓冲队列,Kafka Receiver 以消费者身份拉取消息并解码为 OTel 数据模型,再交给 processor / exporter 处理。该组件的稳定性状态为:profiles 处于 development 阶段,metrics、logs、traces 处于 beta 阶段。
该组件有一个关键特性:如果与配置了include_metadata_keys的kafkaexporter配合使用,Kafka Receiver 会把 Kafka 消息的 headers 传播到下游管线,使管线任意节点都能访问消息携带的元数据键值。此外,对于每条被消费的消息,Receiver 会将其部分记录元数据(topic / partition / offset)以及全部 Kafka headers 注入到请求上下文中,供 attributes processor 等下游组件使用。
在 Grafana Tempo 仓库中,该组件以 vendor 依赖的形式存在于vendor/github.com/open-telemetry/opentelemetry-collector-contrib/receiver/kafkareceiver/下(README 位于vendor/github.com/open-telemetry/opentelemetry-collector-contrib/receiver/kafkareceiver/README.md),同时 Tempo 自身在pkg/ingest/中维护了一套独立的 Kafka 摄取实现,二者共同构成了“Kafka 作为遥测缓冲层”的完整图景。
二、快速开始:零配置起步与核心可选项
2.1 最小配置
Kafka Receiver没有任何必填配置项。下面的配置即可让 Receiver 从localhost:9092、使用otlp_proto编码消费默认 topic:
receivers: kafka:默认情况下它会消费如下信号默认 topic(编码均为otlp_proto):
| 信号 | 默认 topic |
|---|---|
| logs | otlp_logs |
| metrics | otlp_metrics |
| traces | otlp_spans |
| profiles | otlp_profiles |
2.2 关键可选项一览
以下是 README 中给出的全部可选配置项(含默认值):
| 配置项 | 默认值 | 说明 |
|---|---|---|
brokers | localhost:9092 | Kafka broker 地址列表 |
protocol_version | 2.1.0 | Kafka 协议版本 |
resolve_canonical_bootstrap_servers_only | false | 启动时是否解析并对 broker IP 做反向查询 |
logs.topic/logs.topics | otlp_logs | 消费 logs 的 topic(topic已弃用,见下文) |
logs.encoding | otlp_proto | logs topic 的编码 |
logs.exclude_topic/logs.exclude_topics | "" | 正则 topic 模式下排除匹配的 topic |
metrics.topic/metrics.topics | otlp_metrics | 消费 metrics 的 topic |
traces.topic/traces.topics | otlp_spans | 消费 traces 的 topic |
profiles.topic/profiles.topics | otlp_profiles | 消费 profiles 的 topic |
group_id | otel-collector | 消费者组 ID |
client_id | otel-collector | 消费者客户端 ID |
rack_id | "" | 机架标识,配合 broker 的 rack-aware replica selector 从最近副本拉取 |
use_leader_epoch | true | (实验性)是否使用 KIP-320 leader epoch 检测日志截断 |
conn_idle_timeout | 9m | 空闲连接关闭时间 |
initial_offset | latest | 无已提交 offset 时的起始位置,取值latest或earliest |
session_timeout | 10s | 组管理机制下检测客户端故障的请求超时 |
heartbeat_interval | 3s | 到 consumer coordinator 的心跳间隔 |
group_rebalance_strategy | cooperative-sticky | 分区再平衡分配策略 |
group_instance_id | "" | 静态组成员 ID(非空则启用静态成员) |
min_fetch_size | 1 | 单次 fetch 请求的最小消息字节数 |
max_fetch_size | 1048576 | 单次 fetch 请求的最大消息字节数(≥ min_fetch_size) |
max_fetch_wait | 250ms | broker 等待凑满min_fetch_size的最长时间 |
max_partition_fetch_size | 1048576 | 每分区单次 fetch 的字节数(单条 record batch 更大时仍会返回以保证进度) |
metadata.full | true | 是否维护完整元数据集(关闭则启动时不向 broker 发首次请求) |
metadata.refresh_interval | 10m | 集群元数据后台刷新频率 |
metadata.retry.max | 3 | 获取元数据的重试次数 |
metadata.retry.backoff | 250ms | 元数据重试等待时间 |
autocommit.enable | true | 是否自动提交已更新 offset |
autocommit.interval | 1s | 自动提交频率 |
message_marking.after | false | 是否在管线执行完后再标记消息 |
message_marking.on_error | false | 非永久错误时是否标记(false 表示仅标记成功处理的消息) |
message_marking.on_permanent_error | 取on_error值 | 永久错误消息是否标记 |
header_extraction.extract_headers | false | 是否将 header 附加为 resource attribute |
header_extraction.headers | [] | 要提取的 header 列表(精确匹配,不支持正则) |
error_backoff.enabled | false | 消费出错时是否启用退避重试 |
telemetry.metrics.kafka_receiver_records_delay.enabled | false | 是否上报kafka_receiver_records_delay指标 |
三、topic 消费:多 topic 与正则模式
3.1 从topic到topics的演进
Kafka Receiver 底层使用franz-go客户端库,相比传统的librdkafka具备更好的性能,并原生支持现代 Kafka 特性。自 v0.142.0 起,各信号下的topic配置被弃用,统一改为topics(topic 列表);exclude_topic对应弃用为exclude_topics。兼容规则:如果设置了旧的topic字段,它会优先于topics的默认值生效。
3.2 正则 topic 消费与排除
franz-go客户端支持直接通过正则表达式消费多个 topic:在 topic 名前加上^前缀即可开启正则消费(与librdkafka行为一致)。在已弃用的topic设置中,只要任一 topic 带^前缀,就会启用正则消费。
结合exclude_topics可以过滤掉动态 topic 集合中不想要的部分,典型用法如下(注意topic与exclude_topic必须同时使用^正则前缀排除才生效,该特性仅 franz-go 客户端可用):
receivers: kafka: logs: topics: - "^logs-.*" # 消费所有 logs-* 匹配的 topic exclude_topics: - "^logs-(test|dev)$" # 排除 logs-test 与 logs-dev metrics: topics: - "^metrics-.*" exclude_topics: - "^metrics-internal-.*$"上述示例的效果:
- logs:消费
logs-prod、logs-staging、logs-app等,排除logs-test、logs-dev; - metrics:消费
metrics-app、metrics-infra等,排除任何以metrics-internal-开头的 topic。
在 Grafana Tempo 中,Kafka topic 正则消费的另一侧是分区与消费组管理。Tempo 的 metrics-generator 通过 franz-go 消费ingest.kafka.topic指定的 topic,并使用handlePartitionsAssigned/handlePartitionsRevoked/handlePartitionsLost三个回调跟踪分区的分配、吊销与丢失(详见modules/generator/generator_kafka.go)。其中对 lost 分区的处理尤为关键:会话超时或成员被 fenced 时,分区不会走 cooperative revoke 路径,若不在OnPartitionsLost中清理,分区 lag 指标将永远输出过期增长的脏数据(该逻辑在modules/generator/generator_kafka_test.go的TestHandlePartitionsLost_RemovesLostPartitions等用例中有覆盖验证)。
四、支持的编码(Supported Encodings)
除编码扩展(encoding extensions)外,Receiver 内置以下编码:
所有信号通用:
| 编码 | 说明 |
|---|---|
otlp_proto | payload 按 OTLP Protobuf 解码 |
otlp_json | payload 按 OTLP JSON 解码 |
仅 traces 可用:
| 编码 | 说明 |
|---|---|
jaeger_proto | 反序列化为单个 Jaeger protoSpan |
jaeger_json | 用jsonpb反序列化为单个 Jaeger JSON Span |
zipkin_proto | 反序列化为 Zipkin proto spans 列表 |
zipkin_json | 反序列化为 Zipkin V2 JSON spans 列表 |
zipkin_thrift | 反序列化为 Zipkin Thrift spans 列表 |
仅 logs 可用:
| 编码 | 说明 |
|---|---|
raw | payload 字节直接作为 log record 的 body |
text | payload 按文本解码后作为 log record 的 body;默认 UTF-8,可用text_<ENCODING>(如text_utf-8、text_shift_jis)定制 |
json | payload 解码为 JSON 后作为 log record 的 body |
azure_resource_logs | (v0.149.0 弃用,改用azureencodingextension)将 Azure Resource Logs 格式转换为 OTel 格式 |
从源码实现角度看,“解码”这一层在 Tempo 中对应pkg/ingest/encoding.go的GeneratorCodec接口:它定义了Decode([]byte) (iter.Seq2[*tempopb.PushSpansRequest, error], error),并有PushBytesDecoder(反序列化tempopb.PushBytesRequest)与OTLPDecoder(反序列化ptrace.Traces)两个实现,供 metrics-generator 在readCh中根据cfg.Codec选择。这印证了 README 中“编码决定 payload 如何被解释”的设计思路——无论是 OTel Collector 还是 Tempo,都通过可插拔的编解码器解耦 Kafka 字节流与上层数据模型。
五、消息元数据传播(Message metadata propagation)
每条被消费的消息,Receiver 都会把以下记录元数据作为请求元数据(context)注入到管线:
kafka.topic:消息来源 topickafka.partition:消息所在分区kafka.offset:消息在分区内的 offset
此外,消息的全部 Kafka headers 也会被包含进请求元数据。这些元数据可以在管线任意位置使用,例如通过 attributes processor 将其设置为属性。
5.1 Header 提取为资源属性
除了上述隐式传播,Receiver 还支持把指定 header显式提取并挂载为 resource attribute:
receivers: kafka: header_extraction: extract_headers: true headers: ["header1", "header2"]如果向 Kafka 生产一条携带header1: value1、header2: value2的消息,上述配置会将其作为带kafka.header.前缀的资源属性附加:
"resource": { "attributes": { "kafka.header.header1": "value1", "kafka.header.header2": "value2", } } ...注意:header 匹配目前仅支持精确匹配,暂不支持正则。
六、TLS 与认证配置
6.1 TLS + SASL/SCRAM 示例
生产环境最常见的组合是 TLS 加密传输 + SASL 认证。README 给出的示例将tls配置在顶层,auth下配置 SASL:
receivers: kafka: tls: auth: sasl: username: "user" password: "secret" mechanism: "SCRAM-SHA-512"顶层tls支持 OpenTelemetry Collector 的 configtls 全套选项(CA、证书、密钥、insecure_skip_verify等,详见 Collector 的 TLS Configuration Settings)。
6.2 SASL 机制与 Kerberos
auth.sasl.mechanism支持以下取值:
PLAIN(注意:auth.plain_text自 v0.123.0 弃用,改用 sasl 且 mechanism 设为 PLAIN)SCRAM-SHA-256SCRAM-SHA-512AWS_MSK_IAM_OAUTHBEARER(需配合auth.sasl.aws_msk.region指定 AWS 区域)
Kerberos 认证通过auth.kerberos配置:
| 配置项 | 说明 |
|---|---|
service_name | Kerberos 服务名 |
realm | Kerberos realm |
use_keytab | 是否使用 keytab 文件替代密码 |
username/password | 用于向 KDC 认证的凭据 |
config_file | Kerberos 配置路径,如/etc/krb5.conf |
keytab_file | keytab 文件路径,如/etc/security/kafka.keytab |
disable_fast_negotiation | 是否禁用 PA-FX-FAST 协商(默认false,部分 Kerberos 实现不支持 FAST 时需开启) |
另有历史遗留字段:auth.tls(v0.124.0 弃用,为顶层 tls 的别名)。
6.3 Tempo 侧的认证映射
在 Tempo 中,Kafka 的 SASL 认证由pkg/ingest/config.go的KafkaAuthConfig实现,支持的机制常量包括PLAIN、SCRAM-SHA-256、SCRAM-SHA-512、OAUTHBEARER与AWS_MSK_IAM(SASLMechanism类型定义于同文件)。其中:
- PLAIN/SCRAM 需要同时配置 username 与 password,否则校验报
ErrInconsistentSASLCredentials; - OAUTHBEARER 与 AWS_MSK_IAM 支持三种凭据来源:静态凭据、文件路径(每次重新认证时重新读取,可轮换令牌)、HTTP Unix domain socket(每次认证/重认证时通过
GET /获取令牌),且三种来源必须且只能配置一种(kafkaSASLConfig.Validate强制此约束)。
这与 README 中 SASL 机制的设计一一对应,说明 Tempo 在消费侧(metrics-generator / distributor / ingester)复用了同一套 Kafka 认证语义,只是配置入口不同(Tempo 用命令行 flag 与ingest.kafka配置块)。
七、消息确认语义:message_marking 与 error_backoff
Kafka Receiver 的消息确认(marking)行为是生产部署中最容易踩坑的配置,其语义如下:
message_marking.after(默认false):为 true 时,消息在管线执行完之后才被标记;message_marking.on_error(默认false):为 false 时,仅成功处理的消息被标记(针对非永久错误);message_marking.on_permanent_error(默认取on_error的值):为 false 时不标记产生永久错误的消息,为 true 时标记。
两个重要注意事项(README 原文强调):
- 启用
error_backoff时,重试全部耗尽后,失败记录会在下一个 poll 周期自动重试;不启用error_backoff时,分区会一直暂停,直到发生再平衡; - 永久错误不会通过
error_backoff重试,但未提交的消息会在再平衡后被重新处理——这可能阻塞整个分区。
error_backoff采用 Collector 的 configretry 退避配置,包含以下子项:
| 配置项 | 默认值 | 说明 |
|---|---|---|
enabled | false | 是否在消费出错时启用退避 |
initial_interval | - | 首次错误后的等待时间 |
max_interval | - | 连续重试间隔的上界 |
multiplier | - | 退避间隔的倍增系数 |
randomization_factor | - | 随机化因子:实际间隔 = 退避间隔 × (1 ± 随机化因子) |
max_elapsed_time | - | 放弃前的最大退避总时长;为 0 则永不停止重试 |
八、在 Grafana Tempo 中看 Kafka 消费的完整链路
虽然 kafkareceiver 本身是 OTel Collector 组件,但本仓库(Grafana Tempo)恰好给出了 Kafka 作为遥测传输层的完整闭环,可作为理解该组件价值的参照。
8.1 Tempo 的 Kafka 摄取架构
在 Tempo 中,Kafka 位于 distributor 与 metrics-generator / ingester / block-builder 之间:
- 生产侧:distributor 配置
PushSpansToKafka(modules/distributor/config.go中的KafkaConfig ingest.KafkaConfig),调用pkg/ingest/encoding.go的Encode将PushBytesRequest编码为kgo.Record(Key 为 tenant ID,Partition 由分区分配决定),超过maxSize的请求会被拆分到多条记录; - 消费侧:metrics-generator 通过
startKafka启动消费(modules/generator/generator_kafka.go),内部用IngestConcurrency个 goroutine 并行解码,以“先入 channel、多协程解码”的方式把昂贵的 proto unmarshal 从拉取循环中剥离出来;generator 还会按r.Key提取 tenant 并getOrCreateInstance; - 消费组与分区:
pkg/ingest/config.go的GetConsumerGroup(instanceID, partitionID)决定消费组命名——consumer_group为空时使用 ingester 实例 ID 保证唯一性,含<partition>占位符时替换为实际分区号;AutoCreateTopicEnabled默认开启,且会把num.partitions写入 broker 配置以控制自动创建 topic 的分区数(默认 1000)。
对应的最小配置可见example/docker-compose/distributed/tempo.yaml:
ingest: kafka: address: redpanda:9092 topic: tempo-ingest8.2 消费组协调与分区生命周期
Tempo 的生成器消费循环对“分区分配变化”非常敏感,其回调实现与 kafkareceiver 的消费组语义是同一套 franz-go 机制下的两种工程实践:
- 协作式(cooperative)再平衡下,
OnPartitionsAssigned只报告新增分区,因此handlePartitionsAssigned采用 append 而非 replace 维护已分配集合,避免丢失未移动的分区(TestHandlePartitionsAssigned_CooperativeAppend); OnPartitionsLost不提交 offset(kgo 明确警告不要在 lost 回调中提交),仅清理分区 lag 指标;- 停止时,若配置了静态成员(
InstanceID)且LeaveConsumerGroupOnShutdown为 true,会显式发送 LeaveGroup 让协调器立即再平衡,避免等待 session-timeout(TestStopKafka_LeaveGroupConditional、TestPartitionHandoff_LeaveGroupTriggersImmediateReassignment)。
8.3 启动时跳过陈旧积压(stale backlog)
Tempo 还实现了 README 未覆盖但极具实践价值的消费侧优化:skip_stale_backlog_on_startup启用时,generator 在启动阶段通过AdjustFetchOffsetsFn钩子(adjustStartupOffsets)把各分区 fetch offset 前移到“当前时间减去metrics_ingestion_time_range_slack”对应的 horizon offset,从而跳过 slack 窗口之外、注定会被丢弃的陈旧积压,避免重启后重放无用数据,同时让分区 lag 指标保持真实。horizon 查询失败时回退到已提交 offset 重放,属于 best-effort 优化(见modules/generator/generator_kafka.go中startupSeekOffset/startupSeekOffsets的实现与注释)。
九、生产配置建议与注意事项
综合 README 与 Tempo 仓库中的工程实践,以下要点值得在生产部署时重点关注:
- topic 命名与正则:优先使用
topics(列表)而非已弃用的topic;正则 topic 务必以^前缀开头,且排除规则(exclude_topics)需与包含规则同为正则才会生效; - 编码匹配:确认 producer 侧的编码与
encoding一致(如 OTLP 生态两端都用otlp_proto);logs 的text_<ENCODING>变体可满足多字节字符集场景; - 消费组与幂等性:
group_id决定 offset 提交的归属;group_instance_id静态成员适合有状态消费者,可保持重启后分区归属不变,但需注意静态成员停机期间协调器需等待 session-timeout 才能把分区让出,Tempo 通过显式 LeaveGroup 规避该延迟; - 再平衡策略:默认
cooperative-sticky采用增量再平衡,避免全停(stop-the-world)式重分配;sticky则最小化分区迁移但会触发完整再平衡;range/roundrobin是经典策略;还可通过扩展注册实现自定义kgo.GroupBalancer; - 消息确认语义:默认
message_marking.after=false、on_error=false意味着失败消息不会被标记,配合error_backoff实现重试;如需“处理后确认”或永久错误跳过,需显式调整相应开关,并留意未确认消息在再平衡后可能导致的重复消费与分区阻塞; - 认证与传输安全:TLS 统一在顶层
tls配置(auth.tls已弃用);SASL 用PLAIN、SCRAM-SHA-256/512、AWS_MSK_IAM_OAUTHBEARER(需aws_msk.region);Kerberos 场景注意disable_fast_negotiation与旧版 KDC 的兼容问题; - 观测性:
telemetry.metrics.kafka_receiver_records_delay.enabled可上报记录延迟指标,用于评估端到端消费时效;Tempo 侧则通过pkg/ingest导出分区 lag 指标并支持 KIP-714 客户端指标(DisableKafkaTelemetry默认 false)。
十、参考资源
- 组件文档(本仓库 vendored 副本):
vendor/github.com/open-telemetry/opentelemetry-collector-contrib/receiver/kafkareceiver/README.md - Tempo Kafka 摄取编码/解码:
pkg/ingest/encoding.go - Tempo Kafka 客户端配置与校验:
pkg/ingest/config.go - Tempo metrics-generator Kafka 消费循环:
modules/generator/generator_kafka.go - 消费组协调测试:
modules/generator/generator_kafka_test.go - 分布式部署示例(Kafka 摄取配置):
example/docker-compose/distributed/tempo.yaml
【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempo
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考