KafkaNotes

第 24 章:综合实战:日志平台与订单事件系统

zjc 于 2026-01-24 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章把前面所有知识拼成两个完整项目:一个是高吞吐日志采集平台(偏数据管道),一个是电商订单事件系统(偏业务事件驱动)。给出架构、分区设计、关键代码与上线清单。

24.1 项目一:日志采集平台

需求

架构

应用/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,进入本地重试;连续失败写入本地磁盘队列并告警(降级路径)。

监控

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

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 章。死信消息必须带:

幂等的三种实现

方式 适用
唯一事件表(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 积压

本章小结

思考题

  1. 日志 topic 为什么可以 acks=1,订单 topic 为什么必须 acks=all
  2. Outbox Relay 挂掉 10 分钟,系统数据会不一致吗?恢复后如何追平?
  3. 给订单系统加“取消订单”事件,如何保证 CREATED->CANCELLED 与并发支付的顺序正确?