下一篇【第60篇】探针注册通信扩展——基于HTTP的SPI注册实现完全指南 上一篇【第62篇】通信扩展最佳实践——gRPC/HTTP/Kafka全景对比与选型决策
一、为什么是Kafka?
先问一个扎心的问题:gRPC那么好,为什么还要用Kafka?
答案取决于你的部署规模和数据量。来看两种典型场景:
+——————————————————————+
| 两种通信模式的适用场景 |
+——————————————————————+
| |
| 场景A:中小规模(<100个Agent实例) |
| ┌──────────────────────────────────────┐ |
| │ Agent → gRPC → OAP │ ✓ 简单直接 |
| │ │ ✓ 零额外组件 |
| │ Agent数量少,OAP压力可控 │ ✓ 延迟低 |
| │ 直接gRPC通信完全够用 │ |
| └──────────────────────────────────────┘ |
| |
| 场景B:大规模部署(>500个Agent实例) |
| ┌──────────────────────────────────────┐ |
| │ Agent → Kafka → OAP │ ✓ 削峰填谷 |
| │ │ ✓ OAP 可水平扩展 |
| │ Agent数量多,数据量大 │ ✓ 消息可持久化 |
| │ Kafka做缓冲,避免OAP被打爆 │ ✓ 消费可回溯 |
| └──────────────────────────────────────┘ |
| |
| Agent数量 推荐方案 |
| ───────────────────────────────────────── |
| 1-100 直接gRPC |
| 100-500 gRPC + 适当调参 |
| 500-2000 Kafka Reporter |
| 2000+ Kafka + OAP集群 + ES集群 |
| |
+——————————————————————+
Kafka的优势在大规模场景中非常明显:
- 解耦:Agent不关心OAP在不在,只管往Kafka扔数据
- 削峰:Trace数据有波峰波谷,Kafka天然抗波动
- 可回溯:Kafka可按时间范围回放数据,方便故障排查
- 可靠性:Kafka的多副本机制保证数据不丢
二、Kafka Reporter的整体架构
+——————————————————————+
| SkyWalking + Kafka 数据流全景 |
+——————————————————————+
| |
| ┌───────────────────────────┐ |
| │ App JVM 1 │ |
| │ ┌─────────────────────┐ │ |
| │ │ SkyWalking Agent │ │ |
| │ │ ┌─────────────────┐ │ │ |
| │ │ │ Trace Segment │ │ │ |
| │ │ │ Collector CPU=3 │ │ │ |
| │ │ │ Memory=521MB │ │ │ |
| │ │ └────────┬────────┘ │ │ |
| │ │ ↓ │ │ |
| │ │ ┌─────────────────┐ │ │ |
| │ │ │ Kafka Reporter │ │ │ |
| │ │ │ (KafkaProducer) │ │──┼───────→ Kafka Broker |
| │ │ └─────────────────┘ │ │ Topic: |
| │ └─────────────────────┘ │ skywalking-segments |
| └───────────────────────────┘ |
| |
| ┌───────────────────────────┐ |
| │ App JVM 2 │ |
| │ ┌─────────────────────┐ │ |
| │ │ Kafka Reporter │──┼───────→ 同上 |
| │ └─────────────────────┘ │ |
| └───────────────────────────┘ |
| |
| ┌───────────────────────────┐ |
| │ App JVM N │ |
| │ ┌─────────────────────┐ │ |
| │ │ Kafka Reporter │──┼───────→ 同上 |
| │ └─────────────────────┘ │ |
| └───────────────────────────┘ |
| |
| ┌──────────────────────────────────────────────────┐ |
| │ Kafka Cluster │ |
| │ ┌────────────────────────────────────────────┐ │ |
| │ │ Topic: skywalking-segments (Partition x N) │ │ |
| │ │ Topic: skywalking-metrics │ │ |
| │ │ Topic: skywalking-profilings │ │ |
| │ │ Topic: skywalking-managements │ │ |
| │ │ Topic: skywalking-logs │ │ |
| │ └────────────────────────────────────────────┘ │ |
| └──────────────────────┬───────────────────────────┘ |
| │ |
| ┌─────────────┼─────────────┐ |
| │ │ │ |
| ↓ ↓ ↓ |
| ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ |
| │ OAP Server 1│ │ OAP Server 2│ │ OAP Server 3│ |
| │KafkaFetcher │ │KafkaFetcher │ │KafkaFetcher │ |
| │Analyzer │ │Analyzer │ │Analyzer │ |
| └──────┬──────┘ └──────┬──────┘ └──────┬──────┘ |
| │ │ │ |
| └───────────────┼───────────────┘ |
| │ |
| ↓ |
| ┌──────────────────┐ |
| │ Elasticsearch │ |
| │ (持久化存储) │ |
| └──────────────────┘ |
| |
+——————————————————————+
三、Agent端 —— Kafka Reporter配置
SkyWalking 8.x+版本官方就已经内置了Kafka Reporter。你不需要自己写代码,只需要配置即可。
3.1 基础配置
# agent/config/agent.config
# ========== Kafka Reporter 配置 ==========
# 1. 指定使用Kafka作为Trace数据上报通道
plugin.kafka.topic_segment=skywalking–segments
# 2. Kafka集群地址
plugin.kafka.bootstrap_servers=127.0.0.1:9092,127.0.0.1:9093,127.0.0.1:9094
# 3. 其他数据类型也走Kafka
plugin.kafka.topic_metrics=skywalking–metrics
plugin.kafka.topic_profilings=skywalking–profilings
plugin.kafka.topic_managements=skywalking–managements
plugin.kafka.topic_logs=skywalking–logs
# 4. Producer配置
plugin.kafka.producer_config.max_request_size=104857600 # 100MB
plugin.kafka.producer_config.batch_size=16384
plugin.kafka.producer_config.linger_ms=10
plugin.kafka.producer_config.compression_type=lz4
plugin.kafka.producer_config.acks=–1 # all replicas
plugin.kafka.producer_config.retries=3
plugin.kafka.producer_config.max_in_flight_requests_per_connection=1
# 5. 命名空间(多环境隔离)
plugin.kafka.namespace=production
3.2 高级生产配置详解
# Kafka Producer 完整调优参数
plugin.kafka.producer_config:
# === 可靠性 ===
acks: "all" # 等待所有副本确认(最高可靠)
retries: 2147483647 # 最大重试次数
enable.idempotence: true # 幂等生产(防止重复)
# === 性能 ===
batch.size: 131072 # 批次大小(128KB)
linger.ms: 5 # 批次等待时间
buffer.memory: 67108864 # 缓冲区大小(64MB)
compression.type: "lz4" # 压缩算法(lz4/snappy/gzip)
# === 超时 ===
request.timeout.ms: 30000 # 请求超时
delivery.timeout.ms: 120000 # 投递超时
max.block.ms: 60000 # 缓冲区满时的阻塞时间
3.3 Kafka Reporter的实现原理
// Kafka Reporter的核心逻辑(简化自SkyWalking源码)
public class KafkaTraceSegmentServiceClient
implements TracingContextListener, GRPCChannelListener {
private KafkaProducer<String, Bytes> producer;
// 当有新Trace Segment产生时,TracingContext会回调此方法
@Override
public void afterFinished(TraceSegment traceSegment) {
if (traceSegment.isSkipAnalysis()) {
return; // 跳过不需要分析的Segment
}
// 将TraceSegment序列化为Protobuf
UpstreamSegment upstream = traceSegment.transform();
Bytes data = Bytes.wrap(upstream.toByteArray());
// 构建Kafka消息
String topic = String.format(
"%s-%s",
config.getTopicSegment(),
config.getNamespace()
);
// 以TraceId作为Key(保证同一Trace的所有Segment去同一个Partition)
String traceId = traceSegment.getRelatedGlobalTrace().getId();
ProducerRecord<String, Bytes> record =
new ProducerRecord<>(topic, traceId, data);
// 异步发送
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata,
Exception e) {
if (e != null) {
LOGGER.warn("Failed to send trace segment: {}",
e.getMessage());
}
}
});
}
// 初始化Kafka Producer
private void initProducer(KafkaReporterConfig config) {
Properties props = new Properties();
props.put("bootstrap.servers", config.getBootstrapServers());
// 合并自定义Producer配置
if (config.getProducerConfig() != null) {
props.putAll(config.getProducerConfig());
}
props.put("key.serializer",
"org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer",
"org.apache.kafka.common.serialization.BytesSerializer");
this.producer = new KafkaProducer<>(props);
}
}
四、OAP端 —— Kafka Fetcher配置
Agent把数据发到Kafka了,OAP怎么消费呢?
4.1 application.yml配置
# oap-server/config/application.yml
# ========== Kafka Fetcher 配置 ==========
core:
default:
# 指定从Kafka消费数据
selector: ${SW_CLUSTER:standalone}
# 不使用gRPC接收
gRPCHost: ${SW_CORE_GRPC_HOST:0.0.0.0}
gRPCPort: ${SW_CORE_GRPC_PORT:11800}
# Kafka 消费者配置
kafka:
bootstrapServers: ${SW_KAFKA_FETCHER_SERVERS:127.0.0.1:9092}
# === 消费配置 ===
# 消费者组ID
groupId: ${SW_KAFKA_FETCHER_GROUP_ID:skywalking–oap}
# 每个Topic的Partition数量(用于分配消费者)
partitions: ${SW_KAFKA_FETCHER_PARTITIONS:3}
# 每批消费的消息数
batchSize: ${SW_KAFKA_FETCHER_BATCH_SIZE:1000}
# 轮询间隔
pollInterval: ${SW_KAFKA_FETCHER_POLL_INTERVAL:1000}
# === Topic映射 ===
consumers:
– topic: ${SW_KAFKA_TOPIC_SEGMENT:skywalking–segments}
processor: "trace"
– topic: ${SW_KAFKA_TOPIC_METRICS:skywalking–metrics}
processor: "metrics"
– topic: ${SW_KAFKA_TOPIC_PROFILING:skywalking–profilings}
processor: "profiling"
– topic: ${SW_KAFKA_TOPIC_MANAGEMENT:skywalking–managements}
processor: "management"
– topic: ${SW_KAFKA_TOPIC_LOGS:skywalking–logs}
processor: "logs"
# === 消费者高级配置 ===
kafkaConsumerConfig:
enable.auto.commit: true
auto.commit.interval.ms: 5000
session.timeout.ms: 30000
max.poll.records: 500
fetch.max.bytes: 52428800 # 50MB
max.partition.fetch.bytes: 10485760 # 10MB
4.2 Kafka Fetcher的实现逻辑
// KafkaFetcherHandlerRegister.java (简化核心逻辑)
public class KafkaFetcherHandlerRegister
implements FetcherHandlerRegister {
private final Map<String, KafkaConsumer<String, Bytes>> consumers
= new ConcurrentHashMap<>();
@Override
public void register(FetcherConfig config) {
KafkaFetcherConfig kafkaConfig = (KafkaFetcherConfig) config;
for (ConsumerConfig consumerConfig : kafkaConfig.getConsumers()) {
// 为每个Topic创建一个消费者
KafkaConsumer<String, Bytes> consumer =
createConsumer(kafkaConfig, consumerConfig);
consumers.put(consumerConfig.getTopic(), consumer);
// 启动消费者线程
startConsumerThread(consumer, consumerConfig);
}
}
private void startConsumerThread(
KafkaConsumer<String, Bytes> consumer,
ConsumerConfig config) {
Thread consumerThread = new Thread(() -> {
consumer.subscribe(Collections.singletonList(config.getTopic()));
while (!stopped) {
ConsumerRecords<String, Bytes> records =
consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, Bytes> record : records) {
// 根据processor类型分发给不同的处理器
switch (config.getProcessor()) {
case "trace":
traceAnalyzer.analyze(
record.key(), record.value()
);
break;
case "metrics":
metricsAggregator.aggregate(
record.value()
);
break;
case "profiling":
profilingHandler.handle(
record.value()
);
break;
}
}
// 提交offset
consumer.commitAsync();
}
}, "kafka-fetcher-" + config.getTopic());
consumerThread.setDaemon(true);
consumerThread.start();
}
}
五、完整的端到端数据流
把Agent和OAP都配好后,数据是如何流动的?
+——————————————————————+
| Trace数据从Agent到ES的完整生命周期 |
+——————————————————————+
| |
| Step 1: Agent采集 |
| ┌────────────────────────────────────────────┐ |
| │ TracingContext生成TraceSegment │ |
| │ ↓ │ |
| │ 序列化为UpstreamSegment (Protobuf) │ |
| │ ↓ │ |
| │ KafkaReporter.send() │ |
| └────────────────────┬───────────────────────┘ |
| │ |
| Step 2: Kafka传输 |
| ┌────────────────────▼───────────────────────┐ |
| │ Topic: skywalking-segments-prod │ |
| │ Key: traceId → 路由到固定Partition │ |
| │ Value: UpstreamSegment (二进制) │ |
| │ ↓ │ |
| │ 副本同步 → Leader确认 │ |
| └────────────────────┬───────────────────────┘ |
| │ |
| Step 3: OAP消费 |
| ┌────────────────────▼───────────────────────┐ |
| │ KafkaFetcher.poll() → 批量拉取 │ |
| │ ↓ │ |
| │ 反序列化UpstreamSegment │ |
| │ ↓ │ |
| │ TraceAnalyzer.doAnalysis() │ |
| │ ├── Span聚合 │ |
| │ ├── 调用链构建 │ |
| │ ├── 指标计算(OAL) │ |
| │ └── 拓扑推断 │ |
| │ ↓ │ |
| │ 写入Elasticsearch │ |
| └────────────────────────────────────────────┘ |
| |
| Step 4: 查询展示 |
| ┌────────────────────────────────────────────┐ |
| │ SkyWalking UI → GraphQL API → ES查询 │ |
| │ ↓ │ |
| │ 展示Trace详情/拓扑图/指标图表 │ |
| └────────────────────────────────────────────┘ |
| |
+——————————————————————+
六、性能影响与注意事项
6.1 性能对比
# 不同通信方式的性能特征
gRPC (长连接, Protobuf):
延迟: ~5ms (同机房)
吞吐: ~10000 segments/秒 (单OAP)
CPU: Agent端 < 1%, OAP端 中等
内存: Agent端 ~10MB, OAP端 ~200MB
Kafka (异步, 批量):
延迟: ~10–50ms (消息队列缓冲)
吞吐: ~50000 segments/秒 (单OAP消费)
CPU: Agent端 < 1%, OAP端 低
内存: Agent端 ~20MB (Producer缓冲区), OAP端 ~300MB (Consumer缓冲区)
HTTPS (短连接, JSON):
延迟: ~20ms
吞吐: ~3000 segments/秒
CPU: Agent端 低, OAP端 高 (JSON解析)
内存: 较低
6.2 Kafka环境检查清单
# 部署前验证清单
# 1. 检查Kafka集群状态
kafka-topics.sh –bootstrap-server localhost:9092 –list
# 2. 验证Topic创建(Partition数量 >= OAP实例数)
kafka-topics.sh –bootstrap-server localhost:9092 \\
–create \\
–topic skywalking-segments \\ –partitions 6 \\
–replication-factor 3 \\
–config retention.ms=172800000 # 保留2天
# 3. 检查消费者的Lag
kafka-consumer-groups.sh –bootstrap-server localhost:9092 \\
–group skywalking-oap –describe
# 4. 监控Topic的写入速率
kafka-run-class.sh kafka.tools.GetOffsetShell \\
–broker-list localhost:9092 \\
–topic skywalking-segments –time -1
# 5. 验证网络连通性
telnet kafka-broker-1 9092
6.3 常见问题与解决
+——————————————————————+
+ Kafka集成常见问题 +
+——————————————————————+
| |
| 问题1: "Topic不存在" |
| 原因: Kafka未启用auto.create.topics.enable |
| 解决: 手动创建Topic或启用自动创建 |
| |
| 问题2: "OAP不消费数据" |
| 原因: |
| – OAP配置中未启用KafkaFetcher |
| – 消费者Group配置冲突 |
| – Offset异常(未从头消费) |
| 解决: |
| – 检查SW_CORE环境变量 |
| – 重置Consumer Group Offset |
| |
| 问题3: "数据延迟严重" |
| 原因: |
| – Kafka Partition不够(OAP实例并行度受限) |
| – OAP处理能力不足 |
| – 网络带宽不足 |
| 解决: 增加Partition数量 + 增加OAP实例 |
| |
+——————————————————————+
七、总结
用Kafka作为SkyWalking的数据上报通道,是大规模生产环境的最佳实践:
| 适用场景 | Agent实例数 > 500,或需要缓冲/削峰 |
| Agent配置 | 指定bootstrap_servers和各topic名称 |
| OAP配置 | 配置KafkaFetcher消费者,映射processor |
| 监控重点 | Consumer Lag、数据延迟、错误率 |
| 运维要点 | 确保Partition >= OAP实例数,定期清Topic |
下一篇文章将汇总所有通信方案的对比和最佳实践。
下一篇【第60篇】探针注册通信扩展——基于HTTP的SPI注册实现完全指南 上一篇【第62篇】通信扩展最佳实践——gRPC/HTTP/Kafka全景对比与选型决策
网硕互联帮助中心


评论前必须登录!
注册