这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章把前面所有知识拼成两个完整项目:一个是高吞吐日志采集平台(偏数据管道),一个是电商订单事件系统(偏业务事件驱动)。给出架构、分区设计、关键代码与上线清单。
24.1 项目一:日志采集平台
需求
- 2000 台机器,Nginx/应用日志峰值 50 万条/s,均单条 500B;
- 写入 ClickHouse 供查询,写 S3/HDFS 归档;
- 允许极少丢失,但不允许堆积压垮采集端;
- 端到端延迟 P99 < 10s。
架构
应用/Nginx
│ Filebeat/Fluent-bit 采集
v
Kafka 集群
topic: logs.app (分区 24, RF=3)
logs.nginx (分区 24, RF=3)
│ │
│ Consumer A (批量) │ Consumer B
v v
ClickHouse S3 归档任务
│
└--> 查询/看板
容量估算
峰值流量 = 500,000 条/s * 500B = 250MB/s
单分区经验写入 15MB/s -> 至少 17 个分区
考虑消费并行与增长 -> 规划 24 个分区
保留 3 天(未压缩) = 250 * 86400 * 3 ≈ 64.8TB(单副本)
RF=3 + 压缩比 0.6 -> 64.8 * 3 * 0.6 ≈ 116TB
Topic 设计
bin/kafka-topics.sh --create --topic logs.app \
--partitions 24 --replication-factor 3 \
--config retention.ms=259200000 \
--config min.insync.replicas=2 \
--config compression.type=producer \
--config segment.bytes=1073741824
key 用 host + 服务名 均衡写入;日志不需要严格顺序,也可无 key 让粘性分区攒批。
采集端要点(Filebeat 示例)
filebeat.inputs:
- type: log
paths:
- /var/log/app/*.log
fields:
service: order-api
env: prod
fields_under_root: true
output.kafka:
hosts: ["broker1:9092", "broker2:9092", "broker3:9092"]
topic: "logs.app"
required_acks: 1
compression: zstd
bulk_max_size: 4096
worker: 2
采集端对可靠性要求略低于业务事件,required_acks=1 换取吞吐;核心审计日志则用 all。
ClickHouse 消费端
批量拉取、攒批插入是关键:
@KafkaListener(topics = "logs.app", groupId = "log-to-clickhouse", batch = "true")
public void onBatch(List<ConsumerRecord<String, String>> records, Acknowledgment ack) {
List<LogRow> rows = records.stream().map(this::parse).toList();
int chunk = 5000;
for (int i = 0; i < rows.size(); i += chunk) {
clickhouseDao.insertBatch(rows.subList(i, Math.min(i + chunk, rows.size())));
}
ack.acknowledge();
}
配置要点:
spring:
kafka:
consumer:
max-poll-records: 2000
fetch-min-size: 1MB
fetch-max-wait: 2s
listener:
concurrency: 12 # < 分区数 24,可后续扩
ack-mode: manual
ClickHouse 写入失败时不要 ack,进入本地重试;连续失败写入本地磁盘队列并告警(降级路径)。
监控
- 采集端:Filebeat 发送失败率、背压;
- Kafka:BytesIn/Out、lag、ISR;
- 消费端:批大小分布、写库耗时、死信量;
- 业务:端到端延迟 = 日志时间戳 - 入库时间。
24.2 项目二:电商订单事件系统
需求
- 下单成功后:扣库存、加积分、发通知、更新数仓;
- 订单事件绝不能丢,允许重复但必须幂等;
- 同一订单事件必须有序处理;
- 处理失败要有重试与死信,可人工重放。
架构
订单服务
├─ MySQL: orders 表 + outbox 表 (同一本地事务)
└─ Relay 进程读 outbox -> Kafka
│
v
topic: order.events (分区 12, RF=3)
│
┌─────────────┬───┴────────┬──────────────┐
v v v v
库存消费者 积分消费者 通知消费者 数仓消费者
(group=stock) (group=points) (group=notify) (group=dwh)
│
└--失败--> retry topic --> DLT --> 人工/自动补偿
Topic 与顺序设计
bin/kafka-topics.sh --create --topic order.events \
--partitions 12 --replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=604800000
- key =
orderId,保证同一订单的事件进同一分区、严格有序;消费端按 key 做幂等(订单状态机); - 分区 12 > 下游最慢消费者的实例数上限,预留扩容。
Outbox:解决“双写不一致”
直接在业务代码里“写库 + 发 Kafka”会出现一个成功一个失败的分裂。Outbox 把这两步变成“一个本地事务 + 一个可靠投递”:
BEGIN;
INSERT INTO orders(...) VALUES (...); -- 业务数据
INSERT INTO outbox(
aggregate_id, event_type, payload, created_at
) VALUES (
'order-1001', 'OrderCreated', '{"amount":99}', NOW()
);
COMMIT;
Relay 进程(或 Debezium)读取 outbox 表并发布到 Kafka,成功后标记:
@Scheduled(fixedDelay = 200)
public void relay() {
List<OutboxEvent> batch = outboxRepo.lockPendingBatch(500);
if (batch.isEmpty()) return;
for (OutboxEvent e : batch) {
kafkaTemplate.send("order.events", e.aggregateId(), e.payload())
.get(3, TimeUnit.SECONDS); // 同步确认
}
outboxRepo.markPublished(batch);
}
发送可能重复(标记前宕机),所以消费端仍然幂等——这是“至少一次”的正确姿势。
消费者骨架
@KafkaListener(topics = "order.events", groupId = "points-service")
public void onOrderEvent(ConsumerRecord<String, String> record, Acknowledgment ack) {
OrderEvent event = parse(record);
// 1. 去重表/状态机保证幂等
if (dedupService.seen(event.eventId())) {
ack.acknowledge();
return;
}
// 2. 业务处理 + 去重记录放同一事务
pointsService.grantWithDedup(event);
// 3. 处理成功才提交
ack.acknowledge();
}
重试与死信
order.events
-> 处理失败
-> order.events.retry (delay 1m, 重试 3 次, 指数退避)
-> 仍失败
-> order.events.DLT (人工处理/自动补偿)
Spring Kafka 配置见第 9 章。死信消息必须带:
- 原始 topic/partition/offset;
- 异常类型与堆栈摘要;
- 业务主键(orderId),便于按单查询与重放。
幂等的三种实现
| 方式 | 适用 |
|---|---|
| 唯一事件表(event_id 唯一键) | 通用,强一致 |
| 状态机(CREATED->PAID->SHIPPED) | 状态流转型业务 |
| upsert by 业务主键 | 最终状态型(用户资料、积分余额) |
端到端校验
每日对账:
orders 表 vs Kafka order.events (按天 count/sum)
Kafka vs 下游结果表 (库存流水/积分流水)
差异处理:
缺事件 -> 检查 outbox 卡住的记录
缺结果 -> 检查消费者 lag、DLT
多结果 -> 检查幂等键是否失效
24.3 上线检查清单
Topic:
[ ] 分区数覆盖峰值吞吐与消费并行
[ ] RF=3, min.insync.replicas=2
[ ] 保留期 > 故障恢复窗口
生产:
[ ] acks=all + 幂等 + 回调
[ ] Outbox 而非业务代码直发
[ ] key 设计与顺序要求匹配
消费:
[ ] 手动提交,先处理后提交
[ ] 幂等键 + 状态机
[ ] 重试/死信闭环 + 告警
[ ] 静态成员 + 协作式 rebalance
观测:
[ ] lag、端到端延迟、DLT 数量、对账差异
[ ] 生产失败率、outbox 积压
本章小结
- 日志管道优先吞吐:批量、压缩、分区预留、消费端攒批写库;
- 业务事件优先可靠:Outbox + 同 key 有序 + 幂等 + 死信闭环;
- 双写问题不要硬扛,用 Outbox 把跨系统一致性转化为本地事务;
- 对账是事件驱动系统的最后防线,必须自动化。
思考题
- 日志 topic 为什么可以
acks=1,订单 topic 为什么必须acks=all? - Outbox Relay 挂掉 10 分钟,系统数据会不一致吗?恢复后如何追平?
- 给订单系统加“取消订单”事件,如何保证 CREATED->CANCELLED 与并发支付的顺序正确?