拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Grafana Tempo 中的 Kafka Receiver:基于 OpenTelemetry Collector Contrib 的 Kafka 遥测消费与管线元数据传递实践指南

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编码消费默认 topicreceivers: kafka:默认情况下它会消费如下信号默认 topic编码均为otlp_proto信号默认 topiclogsotlp_logsmetricsotlp_metricstracesotlp_spansprofilesotlp_profiles2.2 关键可选项一览以下是 README 中给出的全部可选配置项含默认值配置项默认值说明brokerslocalhost:9092Kafka broker 地址列表protocol_version2.1.0Kafka 协议版本resolve_canonical_bootstrap_servers_onlyfalse启动时是否解析并对 broker IP 做反向查询logs.topic/logs.topicsotlp_logs消费 logs 的 topictopic已弃用见下文logs.encodingotlp_protologs topic 的编码logs.exclude_topic/logs.exclude_topics正则 topic 模式下排除匹配的 topicmetrics.topic/metrics.topicsotlp_metrics消费 metrics 的 topictraces.topic/traces.topicsotlp_spans消费 traces 的 topicprofiles.topic/profiles.topicsotlp_profiles消费 profiles 的 topicgroup_idotel-collector消费者组 IDclient_idotel-collector消费者客户端 IDrack_id机架标识配合 broker 的 rack-aware replica selector 从最近副本拉取use_leader_epochtrue实验性是否使用 KIP-320 leader epoch 检测日志截断conn_idle_timeout9m空闲连接关闭时间initial_offsetlatest无已提交 offset 时的起始位置取值latest或earliestsession_timeout10s组管理机制下检测客户端故障的请求超时heartbeat_interval3s到 consumer coordinator 的心跳间隔group_rebalance_strategycooperative-sticky分区再平衡分配策略group_instance_id静态组成员 ID非空则启用静态成员min_fetch_size1单次 fetch 请求的最小消息字节数max_fetch_size1048576单次 fetch 请求的最大消息字节数≥ min_fetch_sizemax_fetch_wait250msbroker 等待凑满min_fetch_size的最长时间max_partition_fetch_size1048576每分区单次 fetch 的字节数单条 record batch 更大时仍会返回以保证进度metadata.fulltrue是否维护完整元数据集关闭则启动时不向 broker 发首次请求metadata.refresh_interval10m集群元数据后台刷新频率metadata.retry.max3获取元数据的重试次数metadata.retry.backoff250ms元数据重试等待时间autocommit.enabletrue是否自动提交已更新 offsetautocommit.interval1s自动提交频率message_marking.afterfalse是否在管线执行完后再标记消息message_marking.on_errorfalse非永久错误时是否标记false 表示仅标记成功处理的消息message_marking.on_permanent_error取on_error值永久错误消息是否标记header_extraction.extract_headersfalse是否将 header 附加为 resource attributeheader_extraction.headers[]要提取的 header 列表精确匹配不支持正则error_backoff.enabledfalse消费出错时是否启用退避重试telemetry.metrics.kafka_receiver_records_delay.enabledfalse是否上报kafka_receiver_records_delay指标三、topic 消费多 topic 与正则模式3.1 从topic到topics的演进Kafka Receiver 底层使用franz-go客户端库相比传统的librdkafka具备更好的性能并原生支持现代 Kafka 特性。自 v0.142.0 起各信号下的topic配置被弃用统一改为topicstopic 列表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-devmetrics消费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_protopayload 按 OTLP Protobuf 解码otlp_jsonpayload 按 OTLP JSON 解码仅 traces 可用编码说明jaeger_proto反序列化为单个 Jaeger protoSpanjaeger_json用jsonpb反序列化为单个 Jaeger JSON Spanzipkin_proto反序列化为 Zipkin proto spans 列表zipkin_json反序列化为 Zipkin V2 JSON spans 列表zipkin_thrift反序列化为 Zipkin Thrift spans 列表仅 logs 可用编码说明rawpayload 字节直接作为 log record 的 bodytextpayload 按文本解码后作为 log record 的 body默认 UTF-8可用text_ENCODING如text_utf-8、text_shift_jis定制jsonpayload 解码为 JSON 后作为 log record 的 bodyazure_resource_logsv0.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 attributereceivers: 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下配置 SASLreceivers: 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 机制与 Kerberosauth.sasl.mechanism支持以下取值PLAIN注意auth.plain_text自 v0.123.0 弃用改用 sasl 且 mechanism 设为 PLAINSCRAM-SHA-256SCRAM-SHA-512AWS_MSK_IAM_OAUTHBEARER需配合auth.sasl.aws_msk.region指定 AWS 区域Kerberos 认证通过auth.kerberos配置配置项说明service_nameKerberos 服务名realmKerberos realmuse_keytab是否使用 keytab 文件替代密码username/password用于向 KDC 认证的凭据config_fileKerberos 配置路径如/etc/krb5.confkeytab_filekeytab 文件路径如/etc/security/kafka.keytabdisable_fast_negotiation是否禁用 PA-FX-FAST 协商默认false部分 Kerberos 实现不支持 FAST 时需开启另有历史遗留字段auth.tlsv0.124.0 弃用为顶层 tls 的别名。6.3 Tempo 侧的认证映射在 Tempo 中Kafka 的 SASL 认证由pkg/ingest/config.go的KafkaAuthConfig实现支持的机制常量包括PLAIN、SCRAM-SHA-256、SCRAM-SHA-512、OAUTHBEARER与AWS_MSK_IAMSASLMechanism类型定义于同文件。其中PLAIN/SCRAM 需要同时配置 username 与 password否则校验报ErrInconsistentSASLCredentialsOAUTHBEARER 与 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_backoffKafka 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 退避配置包含以下子项配置项默认值说明enabledfalse是否在消费出错时启用退避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 配置PushSpansToKafkamodules/distributor/config.go中的KafkaConfig ingest.KafkaConfig调用pkg/ingest/encoding.go的Encode将PushBytesRequest编码为kgo.RecordKey 为 tenant IDPartition 由分区分配决定超过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.yamlingest: kafka: address: redpanda:9092 topic: tempo-ingest8.2 消费组协调与分区生命周期Tempo 的生成器消费循环对“分区分配变化”非常敏感其回调实现与 kafkareceiver 的消费组语义是同一套 franz-go 机制下的两种工程实践协作式cooperative再平衡下OnPartitionsAssigned只报告新增分区因此handlePartitionsAssigned采用 append 而非 replace 维护已分配集合避免丢失未移动的分区TestHandlePartitionsAssigned_CooperativeAppendOnPartitionsLost不提交 offsetkgo 明确警告不要在 lost 回调中提交仅清理分区 lag 指标停止时若配置了静态成员InstanceID且LeaveConsumerGroupOnShutdown为 true会显式发送 LeaveGroup 让协调器立即再平衡避免等待 session-timeoutTestStopKafka_LeaveGroupConditional、TestPartitionHandoff_LeaveGroupTriggersImmediateReassignment。8.3 启动时跳过陈旧积压stale backlogTempo 还实现了 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_protologs 的text_ENCODING变体可满足多字节字符集场景消费组与幂等性group_id决定 offset 提交的归属group_instance_id静态成员适合有状态消费者可保持重启后分区归属不变但需注意静态成员停机期间协调器需等待 session-timeout 才能把分区让出Tempo 通过显式 LeaveGroup 规避该延迟再平衡策略默认cooperative-sticky采用增量再平衡避免全停stop-the-world式重分配sticky则最小化分区迁移但会触发完整再平衡range/roundrobin是经典策略还可通过扩展注册实现自定义kgo.GroupBalancer消息确认语义默认message_marking.afterfalse、on_errorfalse意味着失败消息不会被标记配合error_backoff实现重试如需“处理后确认”或永久错误跳过需显式调整相应开关并留意未确认消息在再平衡后可能导致的重复消费与分区阻塞认证与传输安全TLS 统一在顶层tls配置auth.tls已弃用SASL 用PLAIN、SCRAM-SHA-256/512、AWS_MSK_IAM_OAUTHBEARER需aws_msk.regionKerberos 场景注意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.mdTempo Kafka 摄取编码/解码pkg/ingest/encoding.goTempo Kafka 客户端配置与校验pkg/ingest/config.goTempo 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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门