项目实战:订单事件流水线
最后一章把全课程知识落到一个真实场景:订单服务产生事件,库存、风控、数仓三个团队各自消费。我们要解决的问题清单:Topic 怎么设计、分区键怎么选、消息不丢不重怎么保证、坏消息怎么兜底、数仓链路怎么做 exactly-once。
1. 架构与需求
+--------------------------+
订单服务 ──(事务发送)──> | Kafka |
| orders.events (12 分区) |
| orders.events.dlq |
+--------------------------+
| | |
消费组: inventory risk dwh-etl
| | |
库存扣减 风控规则 Streams 聚合
(幂等写库) (可容忍重复) (exactly-once)
|
失败重试 3 次仍失败 -> DLQ需求映射到技术决策:
| 需求 | 决策 | 依据章节 |
|---|---|---|
| 同一订单的事件必须有序 | key = orderId | 第 10 章 |
| 事件不能丢 | acks=all + 幂等 + min.isr=2 | 第 6、11 章 |
| 库存不能重复扣减 | 消费端幂等(唯一键) | 第 8 章 |
| 坏消息不阻塞消费 | 重试 + 死信队列 | 第 19 章 |
| 数仓统计精确一次 | Streams EXACTLY_ONCE_V2 | 第 15 章 |
2. Topic 设计
# 主事件流: 12 分区 (峰值 3 万 msg/s / 单分区 3 千 ≈ 10, 取 12 留余量)
bin/kafka-topics.sh --bootstrap-server localhost:9092 --create \
--topic orders.events --partitions 12 --replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=604800000
# 死信队列: 流量极小, 保留期放长便于排查
bin/kafka-topics.sh --bootstrap-server localhost:9092 --create \
--topic orders.events.dlq --partitions 3 --replication-factor 3 \
--config retention.ms=2592000000设计要点:
- 一个 topic 承载订单全生命周期事件(created/paid/cancelled),事件类型放消息体的
eventType字段——同一 orderId 的不同事件才能保证顺序。拆成多个 topic 就没有跨事件顺序了。 - 命名规范
域.实体.事件流;DLQ 用主 topic 名 +.dlq后缀。 - 消息体带
eventId(全局唯一,幂等的依据)、occurredAt、schemaVersion。
{
"eventId": "evt_8f3ac2",
"eventType": "ORDER_PAID",
"orderId": "ord_10091",
"userId": "u_557",
"amount": 129.00,
"occurredAt": "2026-07-30T10:15:30Z",
"schemaVersion": 1
}3. 生产端:订单服务
可靠发送 + 本地 Outbox 兜底(解决"写库和发消息不原子"):
public class OrderEventProducer {
private final KafkaProducer<String, String> producer;
public OrderEventProducer() {
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092,kafka2:9092");
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
p.put(ProducerConfig.ACKS_CONFIG, "all");
p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
p.put(ProducerConfig.LINGER_MS_CONFIG, 10);
p.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
this.producer = new KafkaProducer<>(p);
}
// key = orderId: 同一订单事件永远同分区 => 局部有序
public void publish(OrderEvent e) {
var record = new ProducerRecord<>("orders.events", e.orderId(), toJson(e));
producer.send(record, (meta, ex) -> {
if (ex != null) {
// 发送失败: 落 outbox 表, 由后台任务重发, 保证最终送达
outboxDao.save(e);
}
});
}
}Outbox 模式:业务库事务里同时写业务表和 outbox 表,后台任务轮询 outbox 发 Kafka、成功后标记。这样"下单成功但事件丢了"不可能发生。(另一条路线是 Debezium 直接 CDC 业务表,见第 14 章。)
4. 库存消费组:幂等 + 死信队列
核心:至少一次消费 + 数据库唯一键幂等 + 有限重试 + DLQ。
public class InventoryConsumer {
private static final int MAX_RETRIES = 3;
private final KafkaConsumer<String, String> consumer;
private final KafkaProducer<String, String> dlqProducer;
public void run() {
consumer.subscribe(List.of("orders.events"));
while (running) {
var records = consumer.poll(Duration.ofMillis(500));
for (var r : records) {
handleWithRetry(r);
}
consumer.commitSync(); // 整批处理完(含进 DLQ)才提交
}
}
private void handleWithRetry(ConsumerRecord<String, String> r) {
for (int attempt = 1; ; attempt++) {
try {
handle(r);
return;
} catch (RetriableException e) { // 网络/超时类: 重试
if (attempt >= MAX_RETRIES) { sendToDlq(r, e); return; }
sleepBackoff(attempt); // 1s, 2s, 4s...
} catch (Exception e) { // 解析失败等毒消息: 直接 DLQ
sendToDlq(r, e);
return;
}
}
}
private void handle(ConsumerRecord<String, String> r) {
OrderEvent e = parse(r.value());
if (!"ORDER_PAID".equals(e.eventType())) return;
// 幂等核心: processed_events.event_id 唯一索引
// INSERT INTO processed_events(event_id) VALUES (?)
// ON CONFLICT DO NOTHING; 影响行数为 0 => 已处理过, 跳过
boolean firstTime = dedupDao.tryInsert(e.eventId());
if (!firstTime) return;
inventoryDao.deduct(e.orderId()); // 与上一句同一个 DB 事务!
}
private void sendToDlq(ConsumerRecord<String, String> r, Exception cause) {
var dlqRecord = new ProducerRecord<>("orders.events.dlq", r.key(), r.value());
dlqRecord.headers()
.add("dlq.error", cause.toString().getBytes(StandardCharsets.UTF_8))
.add("dlq.source.partition", String.valueOf(r.partition()).getBytes())
.add("dlq.source.offset", String.valueOf(r.offset()).getBytes());
dlqProducer.send(dlqRecord);
alert("order event moved to DLQ: " + r.key()); // DLQ 必须有告警!
}
}Python 版核心逻辑(confluent-kafka):
def handle(msg):
event = json.loads(msg.value())
if event["eventType"] != "ORDER_PAID":
return
with db.transaction():
# 唯一键幂等: 重复 eventId 直接跳过
inserted = db.execute(
"INSERT INTO processed_events(event_id) VALUES (%s) "
"ON CONFLICT DO NOTHING", (event["eventId"],))
if inserted == 0:
return
db.execute("UPDATE inventory SET stock = stock - 1 "
"WHERE order_id = %s", (event["orderId"],))
while True:
msg = consumer.poll(0.5)
if msg is None or msg.error():
continue
try:
retry(handle, msg, times=3, backoff=1.0)
except Exception as e:
send_to_dlq(msg, e)
consumer.commit(message=msg, asynchronous=False)幂等的关键:去重记录(processed_events)与业务写入(inventory)在同一个数据库事务里。分开写就会出现"记了已处理但扣减失败"的漏洞。
5. 数仓链路:Streams exactly-once 聚合
每小时 GMV 统计,消费-聚合-输出全程精确一次:
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "dwh-order-agg");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
StreamsBuilder builder = new StreamsBuilder();
builder.stream("orders.events", Consumed.with(Serdes.String(), Serdes.String()))
.filter((k, v) -> v.contains("\"ORDER_PAID\""))
.mapValues(v -> parseAmount(v))
.groupBy((k, amt) -> "gmv") // 全局聚合
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofHours(1)))
.reduce(Double::sum, Materialized.as("gmv-hourly"))
.toStream()
.map((win, gmv) -> KeyValue.pair(win.key() + "@" + win.window().start(),
String.valueOf(gmv)))
.to("dwh.gmv.hourly");
new KafkaStreams(builder.build(), props).start();EXACTLY_ONCE_V2 底层就是第 6 章的事务:消费 offset、状态 changelog、输出消息三者原子提交,任何一步失败整体回滚重放,结果 topic 里的聚合值不重不漏。
风控消费组逻辑类似库存组,但风控规则天然幂等(重复评估同一事件结果相同),可以省去去重表、用简单的 at-least-once。按下游语义选一致性等级,不要全链路无脑上事务。
6. 上线检查清单
| 类别 | 检查项 |
|---|---|
| 可靠性 | acks=all、幂等开启、min.insync.replicas=2、3 副本 |
| 顺序 | key=orderId;消费端同分区串行处理 |
| 消费 | 手动提交;CooperativeStickyAssignor;静态成员(K8s) |
| 兜底 | DLQ + 告警 + 重放工具(从 DLQ 读出修复后回灌主 topic) |
| 监控 | 三个消费组的 Lag 告警、DLQ 非空告警、UnderReplicated 告警 |
| 演练 | kill 一个 broker / 一个消费者实例,验证无丢失无重复扣减 |
- DLQ 有进无出:消息进了死信没人看,等于静默丢单。DLQ 非空必须告警,且要有回放工具与值班流程。
- 去重表无限膨胀:processed_events 按时间分区、保留期大于 Kafka retention 即可安全清理。
- 用 userId 当 key:大卖家会造成热分区。本例顺序需求是"单订单内有序",所以 orderId 是正确粒度——key 粒度永远取"顺序需求的最细粒度"。
这套骨架(有序 key + outbox 生产 + 幂等消费 + DLQ + 分级一致性)适用于绝大多数事件驱动系统:支付回调、物流轨迹、积分变动。换个领域名词就能复用。需要跨库大规模同步时,把手写 outbox 换成 Debezium(第 14 章);聚合规模超出 Streams 时换 Flink(第 15 章)。
小结
- Topic 设计:一个实体一个事件流 topic,事件类型放消息体,key 取顺序需求的最细粒度(orderId)
- 生产端:acks=all + 幂等 + Outbox 模式,保证"业务成功则事件必达"
- 消费端:at-least-once + 数据库唯一键幂等(去重与业务同事务)+ 有限重试 + DLQ 告警
- 数仓聚合用 Streams EXACTLY_ONCE_V2;风控等天然幂等的链路用简单 at-least-once
- 一致性等级按下游语义分级选择,上线前做 broker/消费者故障演练
- 完整搭建本章流水线(可用 Docker Compose 单机三 broker),实现生产者、库存消费组(含 DLQ)、Streams 聚合三个程序。
- 故障演练:处理中 kill -9 库存消费者再重启,验证扣减不重复;kill 一个 broker,验证事件不丢。
- 写一个 DLQ 回放工具:读取死信消息,修复后带上原 headers 回灌主 topic,并保证回放本身也是幂等的。