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

Apache SkyWalking Kafka Fetcher 完整指南:用 Kafka 解耦探针数据上报与 OAP 消费链路

Apache SkyWalking Kafka Fetcher 完整指南用 Kafka 解耦探针数据上报与 OAP 消费链路【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sky/skywalking本文围绕 SkyWalking OAP 的kafka-fetcher模块展开系统讲解如何让 Java Agent 通过 Kafka Broker 上报 Trace 段、JVM 指标、实例属性、Meter 数据、Profiling 任务与日志并由 OAP 侧 Kafka Fetcher 拉取消费。读完本文你将掌握 Kafka Fetcher 的启用方式、全部配置项与环境变量、Topic 自动创建机制、多 OAP 集群的 Namespace 隔离方案、MirrorMaker 2.0 跨集群复制适配以及消息从 Kafka 到分析器的源码级处理链路可直接在真实集群中落地这套消息队列承载遥测数据的架构。一、Kafka Fetcher 是什么从 gRPC 直连到消息队列解耦默认情况下SkyWalking Agent 通过 gRPC或 HTTP协议将采集到的数据直接推送给 OAP 的 Receiver。kafka-fetcher提供了另一条完全不同的接入路径Agent 把数据写入 Kafka TopicOAP 以 Kafka Consumer 的身份从 Broker 拉取消息并送入分析管线。这种模式的价值在于削峰填谷Agent 与 OAP 之间通过 Kafka 解耦瞬时流量峰值由 Broker 缓冲OAP 不必直接面对突发上报多种传输协议共存官方文档明确指出Kafka Fetcher 可以与 gRPC/HTTP Receiver同时启用满足不同 Agent或同一集群中不同语言/版本的探针采用不同传输协议的场景下游生态复用Kafka 集群中保留的原始消息可被其他消费方如自建流处理、审计系统复用。需要注意Kafka Fetcher默认处于关闭状态在 application.yml 中selector: ${SW_KAFKA_FETCHER:-}环境变量为空时即不加载该模块必须显式配置后才能生效。二、Kafka Fetcher 支持的数据类型从 KafkaFetcherProvider.start() 的处理器注册逻辑可以确认Kafka Fetcher 覆盖了以下数据类型具体支持范围还取决于 Agent 端的实现数据类型对应 Handler说明Trace 段Tracing SegmentsTraceSegmentHandler解析SegmentObjectProtobuf 消息送入 Segment 解析服务JVM 指标JVMMetricsHandlerJVM 运行时指标CPU、内存、GC 等服务/实例属性ServiceManagementHandlerAgent 注册、实例属性、心跳等Meter 系统数据MeterServiceHandler自定义 Meter 指标Profiling 任务ProfileTaskHandler线程剖析任务下发与快照日志Protobuf 格式LogHandler由enableNativeProtoLog开关控制默认开启日志JSON 格式JsonLogHandler由enableNativeJsonLog开关控制默认开启其中日志类 Handler 是可选注册的只有enableNativeProtoLog为true时才注册LogHandler只有enableNativeJsonLog为true时才注册JsonLogHandler两个开关的默认值均为true见 KafkaFetcherConfig.java。若你的日志上报走其他通道可通过环境变量SW_KAFKA_FETCHER_ENABLE_NATIVE_PROTO_LOG/SW_KAFKA_FETCHER_ENABLE_NATIVE_JSON_LOG关闭对应消费避免 Topic 空转。三、快速启用最小配置在 OAP 的配置文件默认为 oap-server/server-starter/src/main/resources/application.yml中加入以下片段即可启用kafka-fetcher: selector: ${SW_KAFKA_FETCHER:default} default: bootstrapServers: ${SW_KAFKA_FETCHER_SERVERS:localhost:9092} namespace: ${SW_NAMESPACE:}关键参数说明selector模块选择器置为default表示加载内置实现。默认值直接取自环境变量SW_KAFKA_FETCHER不设置该变量时模块不启动bootstrapServersKafka Broker 地址列表host:port对多个用逗号分隔用于建立与集群的初始连接对应 Kafka 客户端的bootstrap.servers配置namespace命名空间用于隔离共享同一 Kafka 集群的多个 OAP 集群详见下文第四节。启用 Kafka Fetcher 后还需要在Agent 端开启 Kafka 上报即 Kafka Reporter让 Agent 将数据写入对应的 Topic而不是走 gRPC。具体开启方式参见 Agent 文档本仓库 docs 原文也提示Check the agent documentation for details on how to enable the Kafka reporter。四、Namespace 机制多 OAP 集群共享 Kafka 的隔离方案当多个 OAP 集群共用一个 Kafka 集群时Topic 会发生冲突。namespace正是为此设计的隔离手段其工作方式为在 OAP 的 Kafka Fetcher 配置中设置namespace后OAP 消费的 Topic 名会被加上前缀与此同时Agent 端也必须在agent.config中设置同名属性plugin.kafka.namespace保证 Agent 写入的 Topic 与 OAP 消费的 Topic 一致。从源码可以精确看到前缀的拼接规则。在 AbstractKafkaHandler.getTopic() 中最终 Topic 名的构造顺序是mm2SourceAlias mm2SourceSeparator见第六节→namespace -→ 原始 Topic 名。例如设置namespace: prod后skywalking-segments实际会被消费为prod-skywalking-segments。这是一个两端必须对齐的约定Agent 与 OAP 的 namespace 不一致时OAP 将收不到任何数据排查 Kafka 场景数据缺失问题时应首先核对两侧配置。五、Topic 体系与自动创建机制5.1 必需的七个 Topickafka-fetcher依赖以下七个 Topic默认名称均可在 KafkaFetcherConfig.java 中看到对应字段Topic 名称承载数据对应配置字段skywalking-segmentsTrace 段topicNameOfTracingSegmentsskywalking-metricsJVM/服务指标topicNameOfMetricsskywalking-profilingsProfiling 数据topicNameOfProfilingskywalking-managements服务/实例管理信息topicNameOfManagementsskywalking-metersMeter 指标topicNameOfMetersskywalking-logsProtobuf 格式日志topicNameOfLogsskywalking-logs-jsonJSON 格式日志topicNameOfJsonLogs5.2 自动创建与分区/副本控制如果这些 Topic 不存在Kafka Fetcher 启动时会自动创建当然你也可以在 OAP 启动前用 Kafka 管理工具如kafka-topics.sh预先创建。自动创建逻辑位于 KafkaFetcherHandlerRegister.createTopicIfNeeded()OAP 使用AdminClient复用bootstrapServers配置并移除group.id以避免告警先describeTopics探测缺失的 Topic再对缺失项按partitions与replicationFactor调用createTopics批量创建。创建参数可配置如下kafka-fetcher: selector: ${SW_KAFKA_FETCHER:default} default: bootstrapServers: ${SW_KAFKA_FETCHER_SERVERS:localhost:9092} namespace: ${SW_NAMESPACE:} partitions: ${SW_KAFKA_FETCHER_PARTITIONS:3} replicationFactor: ${SW_KAFKA_FETCHER_PARTITIONS_FACTOR:2} consumers: ${SW_KAFKA_FETCHER_CONSUMERS:1}partitions自动创建 Topic 时的分区数默认3。分区数决定单个 Topic 的并行度上限同时会直接影响消费吞吐replicationFactor每个分区的副本因子默认2。生产环境建议按 Kafka 集群规模设置不小于 3consumersOAP 侧启动的 KafkaConsumer 实例数默认1。从 KafkaFetcherHandlerRegister 构造函数可以看到每增加一个 consumer 就会 new 一个KafkaConsumerString, BytesString 反序列化 Key、Bytes 反序列化 Value所有 consumer 共享同一个group.id默认skywalking-consumer见 KafkaFetcherConfig.java因此它们是同一消费组内负载均衡的关系每个 consumer 各跑一个独立的消息拉取循环。六、MirrorMaker 2.0 跨集群复制支持当使用 Kafka MirrorMaker 2.0 在 Kafka 集群之间复制 Topic 时目标集群上的 Topic 名称会被加上源集群别名与分隔符例如source-cluster.topic-name。为了让 OAP 正确消费这些带前缀的复制 TopicKafka Fetcher 提供了两个参数kafka-fetcher: selector: ${SW_KAFKA_FETCHER:default} default: bootstrapServers: ${SW_KAFKA_FETCHER_SERVERS:localhost:9092} namespace: ${SW_NAMESPACE:} partitions: ${SW_KAFKA_FETCHER_PARTITIONS:3} replicationFactor: ${SW_KAFKA_FETCHER_PARTITIONS_FACTOR:2} consumers: ${SW_KAFKA_FETCHER_CONSUMERS:1} mm2SourceAlias: ${SW_KAFKA_MM2_SOURCE_ALIAS:} mm2SourceSeparator: ${SW_KAFKA_MM2_SOURCE_SEPARATOR:} kafkaConsumerConfig: enable.auto.commit: true ...mm2SourceAlias源 Kafka 集群在 MirrorMaker 2.0 中配置的别名对应 MM2 的source.cluster.aliasmm2SourceSeparatorMM2 拼接远程 Topic 名使用的分隔符对应 MM2 的replication.policy.separator默认通常是.。结合上文第四节可以看到getTopic()的拼接顺序是MM2 前缀在前、namespace 前缀在后最终形态如source-cluster.prod-skywalking-segments两者可同时使用。上述两个参数与namespace默认值均为空字符串见 KafkaFetcherConfig.java不配置时不影响 Topic 名。七、kafkaConsumerConfig原生 Kafka 消费参数透传kafkaConsumerConfig是一个Properties类型配置项对应 KafkaFetcherConfig.kafkaConsumerConfig允许你把任意原生 Kafka Consumer 配置透传给底层KafkaConsumer。官方示例中的enable.auto.commit: true即为典型用法kafkaConsumerConfig: enable.auto.commit: true # 其他原生 ConsumerConfig 属性例如 # auto.offset.reset: earliest # max.poll.records: 500 # fetch.max.wait.ms: 500该配置的语义与 Kafka 客户端完全一致举几个常用场景enable.auto.commit是否自动提交消费位点。OAP 默认按true处理见 KafkaFetcherHandlerRegister.java 的getOrDefault逻辑若显式设为false则 runTask() 会在每批消息派发完成后调用commitAsync()手动异步提交auto.offset.reset无已提交位点时从何处开始消费max.poll.records/fetch.max.wait.ms等调节单次 poll 的数据量与拉取节奏。注意该配置块是 YAML 键值对与原生 Kafka 属性名的直接映射属性名保持 Kafka 客户端的点分命名如enable.auto.commit不要写成驼峰式。八、源码剖析消息从 Kafka 到 OAP 分析链路的完整通路理解内部处理流程有助于定位问题与做容量规划。整个模块的消费骨架在 KafkaFetcherHandlerRegister.java 中核心流程如下初始化消费客户端构造函数L68-L98以group.idskywalking-consumer、StringDeserializerBytesDeserializer构造指定数量的KafkaConsumerString, Bytes同时创建一个固定大小线程池默认核心/最大线程数为CPU 核数 × 2队列容量默认 10000均可用kafkaHandlerThreadPoolSize/kafkaHandlerThreadPoolQueueSize覆盖见 application.yml线程工厂名为KafkaConsumer饱和策略为CallerRunsPolicy注册处理器映射L100-L102每个 Handler 绑定一个 Topic构建topic → KafkaHandler的不可变映射。KafkaHandler接口只定义两个方法——getTopic()与handle(ConsumerRecordString, Bytes)见 KafkaHandler.java启动消费start()L104-L114先调用createTopicIfNeeded()补齐缺失 Topic然后每个 consumersubscribe全部 Topic、seekToEnd跳到最新位点即只消费启动之后的新消息不重放历史数据最后向线程池提交拉取任务拉取与派发runTask()L116-L132无限循环中每 500mspoll一次对每条记录按 Topic 找到对应 Handler提交到线程池异步处理enable.auto.commitfalse时在批量派发后commitAsync()提交位点按类型反序列化并送入下游分析器以 TraceSegmentHandler.handle() 为例——将 Value 反序列化为SegmentObjectProtobuf随后调用segmentParserService.send(segment)送入 Segment 解析服务与 gRPC 通道共用同一套分析管线日志类则分别走 LogHandlerProtobuf 的LogData.parseFrom和 JsonLogHandler用ProtoBufJsonUtils.fromJSON将 JSON 文本解析为LogData再交给ILogAnalyzerService做日志分析。值得留意的是各 Handler 都注册了 Telemetry 指标如trace_in_latency、trace_analysis_error_count、log_in_latency、log_analysis_error_count标签为protocolkafka见 TraceSegmentHandler.java 与 LogHandler.java。启用 OAP 自监控so11y后可通过这些指标直接观察 Kafka 通道的处理时延与解析失败数量无需额外埋点即可完成链路健康度巡检。九、配置项与环境变量速查表以下汇总kafka-fetcher的全部配置项以 application.yml 与 KafkaFetcherConfig.java 为准配置项环境变量默认值说明selectorSW_KAFKA_FETCHER空模块关闭模块选择器设为default启用bootstrapServersSW_KAFKA_FETCHER_SERVERSlocalhost:9092Kafka Broker 地址列表namespaceSW_NAMESPACETopic 前缀多 OAP 集群隔离需与 Agent 的plugin.kafka.namespace一致partitionsSW_KAFKA_FETCHER_PARTITIONS3自动创建 Topic 的分区数replicationFactorSW_KAFKA_FETCHER_PARTITIONS_FACTOR2自动创建 Topic 的副本因子consumersSW_KAFKA_FETCHER_CONSUMERS1消费组内 KafkaConsumer 实例数enableNativeProtoLogSW_KAFKA_FETCHER_ENABLE_NATIVE_PROTO_LOGtrue是否注册 Protobuf 日志 HandlerenableNativeJsonLogSW_KAFKA_FETCHER_ENABLE_NATIVE_JSON_LOGtrue是否注册 JSON 日志 HandlerkafkaHandlerThreadPoolSizeSW_KAFKA_HANDLER_THREAD_POOL_SIZE-1取 CPU 核数 × 2消息处理线程池大小kafkaHandlerThreadPoolQueueSizeSW_KAFKA_HANDLER_THREAD_POOL_QUEUE_SIZE-1取 10000消息处理线程池队列容量mm2SourceAliasSW_KAFKA_MM2_SOURCE_ALIASMirrorMaker 2.0 源集群别名mm2SourceSeparatorSW_KAFKA_MM2_SOURCE_SEPARATORMirrorMaker 2.0 Topic 分隔符kafkaConsumerConfig—空 Properties透传原生 Kafka Consumer 配置如enable.auto.commit另有模块内部使用的groupId默认skywalking-consumer与configPath默认meter-analyzer-config供 Meter 分析配置使用等字段未在 application.yml 中显式列出按 Java 字段默认值生效。十、生产环境实践建议两端配置一致性优先开启 Kafka 通道前先核对 Agent 端plugin.kafka.namespace与 OAP 端namespace完全一致涉及 MirrorMaker 时再叠加核对mm2SourceAlias/mm2SourceSeparator否则消费侧会静默收不到数据容量参数按需调优partitions、consumers与kafkaHandlerThreadPoolSize三者共同决定消费吞吐。建议让 Topic 分区数 ≥ consumer 实例数并观察 so11y 中*_in_latencyprotocolkafka指标确认处理时延处于合理区间消费位点策略OAP 启动时会seekToEnd只消费新消息若因故障需要重放历史数据应在消费组skywalking-consumer的位点层面处理而不是重启 OAP 期待自动重放与 gRPC/HTTP Receiver 共存Kafka Fetcher 与 gRPC/HTTP Receiver 可以同时开启用于支持不同传输协议的 Agent 混合部署互不影响日志开关按需裁剪不需要日志消费时将SW_KAFKA_FETCHER_ENABLE_NATIVE_PROTO_LOG/SW_KAFKA_FETCHER_ENABLE_NATIVE_JSON_LOG置为false可减少无效 Topic 与消费负担故障排查锚点相关 FAQ 与词汇表可参考本仓库的 kafka-plugin.mdKafka 消费端手工埋点问题与 configuration-vocabulary.md配置项统一说明。综上Kafka Fetcher 为 SkyWalking 提供了一条与 gRPC/HTTP 并行的消息队列接入通道。理解其 Topic 命名规则、自动创建机制与消费线程模型是保证Agent 写、OAP 读两侧无缝衔接的关键配合 so11y 指标与上述实践建议即可在共享 Kafka 集群、跨集群复制等复杂拓扑下稳定运行。【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sky/skywalking创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

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

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