后端开发

事务发件箱实战:MySQL Outbox + Kafka 构建可靠事件驱动系统

解决数据库与 Kafka 双写一致性:Outbox 表、并发 Relay、分区顺序、Inbox 幂等、重放和故障注入。

TY
Tycho
技术博主
• 2026-09-27 • 24 分钟阅读 • 1 次浏览
事务发件箱实战:MySQL Outbox + Kafka 构建可靠事件驱动系统

一、双写为何会丢失或制造幽灵事件

业务服务先写 MySQL 再发 Kafka 时,数据库提交后进程崩溃会丢事件;先发 Kafka 再提交数据库,则回滚时产生幽灵事件。Outbox 把业务行和事件行放进一个本地事务,由独立 relay 投递。

relay 在 Kafka 确认后、标记 published 前崩溃会重复发送,因此系统按至少一次交付设计,消费者必须幂等。Kafka 事务不能自动覆盖任意外部数据库副作用。

API → MySQL transaction
      ├─ orders
      └─ outbox_events
              ↓ relay / CDC
           Kafka topic
              ↓
       consumer + inbox/唯一键 → 下游状态

二、设计可演进的 Outbox 表

事件 ID 在事务内生成并永久稳定,schema_version 支持兼容解析,payload 保存事件发生时的事实。错误信息限长并脱敏。ready 索引支持按时间认领未发布事件。

CREATE TABLE outbox_events (
 id CHAR(36) PRIMARY KEY, aggregate_type VARCHAR(64) NOT NULL,
 aggregate_id VARCHAR(128) NOT NULL, event_type VARCHAR(128) NOT NULL,
 schema_version SMALLINT UNSIGNED NOT NULL, payload JSON NOT NULL, headers JSON NULL,
 occurred_at DATETIME(6) NOT NULL, published_at DATETIME(6) NULL, attempts INT UNSIGNED DEFAULT 0,
 next_attempt_at DATETIME(6) NOT NULL, locked_at DATETIME(6) NULL, locked_by VARCHAR(128) NULL,
 last_error VARCHAR(1000) NULL,
 KEY idx_ready (published_at,next_attempt_at,occurred_at),
 KEY idx_aggregate (aggregate_type,aggregate_id,occurred_at)
) ENGINE=InnoDB;

三、同一事务写业务和事件

远程调用不能放在数据库事务内。事务只写本地业务状态和完整事件,任何异常同时回滚。payload 不要求消费者回查可能已经改变的业务表。

DB::transaction(function () use ($command) {
  $order = Order::create(['customer_id'=>$command->customerId,'status'=>'created','total_amount'=>$command->total]);
  OutboxEvent::create([
    'id'=>(string) Str::uuid(),'aggregate_type'=>'order','aggregate_id'=>(string) $order->id,
    'event_type'=>'order.created','schema_version'=>1,
    'payload'=>['order_id'=>$order->id,'customer_id'=>$order->customer_id,'total_amount'=>(string) $order->total_amount],
    'occurred_at'=>now(),'next_attempt_at'=>now(),
  ]);
});

四、并发 Relay 的认领租约

数据库事务只负责用 SKIP LOCKED 认领一小批事件并立即提交,真实 Kafka 发送在事务外执行。多实例不会抢同一行;进程崩溃后,守护任务把超过租约的锁重新开放。

START TRANSACTION;
SELECT id FROM outbox_events
WHERE published_at IS NULL AND next_attempt_at<=NOW(6)
ORDER BY occurred_at LIMIT 100 FOR UPDATE SKIP LOCKED;
UPDATE outbox_events SET locked_at=NOW(6),locked_by='relay-a',attempts=attempts+1
WHERE id IN (...);
COMMIT;

五、发布确认、退避与毒消息

Kafka ack 成功后更新 published_at。失败按指数退避和上限安排 next_attempt_at;超过尝试阈值进入人工处理,不让同一毒消息无限占用吞吐。错误日志包含事件 ID 与异常类别,不含凭据。

try {
  $producer->send('orders.v1', $event->aggregate_id, encodeEvent($event));
  $event->update(['published_at'=>now(),'locked_at'=>null,'locked_by'=>null]);
} catch (Throwable $e) {
  $event->update([
    'next_attempt_at'=>now()->addSeconds(min(300, 2 ** $event->attempts)),
    'locked_at'=>null,'locked_by'=>null,'last_error'=>Str::limit($e->getMessage(),1000),
  ]);
}

六、主题、分区键与顺序边界

aggregate_id 作为 key,使同一订单进入同一分区并保持分区内顺序。不同聚合之间没有全局顺序保证;扩分区会改变 key 映射,必须评估消费者的顺序假设。

kafka-topics.sh --bootstrap-server kafka:9092 \
 --create --topic orders.v1 --partitions 12 --replication-factor 3 \
 --config min.insync.replicas=2
kafka-topics.sh --bootstrap-server kafka:9092 --describe --topic orders.v1

七、消费者 Inbox 幂等

Inbox 记录与业务副作用放进同一数据库事务。若事件 ID 已存在就直接返回,提交数据库后再提交 Kafka offset。外部 API 还要传递同一事件 ID 作为对方幂等键。

CREATE TABLE consumed_events (consumer_name VARCHAR(128) NOT NULL,event_id CHAR(36) NOT NULL,consumed_at DATETIME(6) NOT NULL,PRIMARY KEY (consumer_name,event_id)) ENGINE=InnoDB;
DB::transaction(function () use ($message) {
  $inserted = DB::table('consumed_events')->insertOrIgnore([
    'consumer_name'=>'inventory-reservation','event_id'=>$message->event_id,'consumed_at'=>now(),
  ]);
  if ($inserted === 0) return;
  Inventory::reserve($message->data->order_id, $message->data->items);
});
// 数据库提交后再提交 offset

八、契约版本和受控重放

事件字段默认只新增不删除,消费者忽略未知字段并按 schema_version 解析。金额使用字符串或最小货币单位。重放进入独立主题或消费组并限速,不和实时流量争夺下游容量。

{
  "event_id":"f9a6b944-6c45-4f16-8d88-024f714ccf1a",
  "type":"order.created",
  "version":1,
  "occurred_at":"2026-09-26T10:12:30.123456+08:00",
  "data":{"order_id":9182,"customer_id":42,"total_amount":"199.00"}
}

九、监控与故障注入

监控最老未发布事件年龄、失败率、消费者 lag、重复命中率和分区热点。每次只注入一种故障,证明任意崩溃点最多重复而不丢失,并保存恢复时间。

# 依次演练:提交后 relay 前崩溃、Kafka ack 后标记前崩溃、业务提交后 offset 前崩溃
mysql -e "SELECT COUNT(*),MIN(occurred_at) FROM outbox_events WHERE published_at IS NULL"
kafka-consumer-groups.sh --bootstrap-server kafka:9092 --describe --group inventory-reservation
  • 毒消息有隔离与人工恢复流程。
  • 已发布数据按保留窗口小批清理。
  • 重复率异常时先核对 offset 和 relay 确认逻辑。

十、上线验收

事务回滚时业务行和 outbox 行都不存在,提交时两者都存在。relay 任意崩溃不会丢事件;消费者重复执行不重复扣减库存或调用外部服务。积压、重放和清理均有运行手册。

  • 分区键和顺序边界写入契约。
  • 消费者外部副作用也有幂等键。
  • 故障演练覆盖三个关键崩溃窗口。
  • 清理条件只针对已发布且超过重放窗口的事件。

总结

Transactional Outbox 用单库事务消除双写丢失窗口,再以至少一次投递、消费者幂等和可观测重放完成端到端可靠性。工程重点在事务边界、认领租约、分区顺序、契约版本和故障注入,而不只是创建一张表。

官方资料与继续学习

TY

Tycho

热爱分享技术知识,帮助开发者成长。

评论 (0)

评论功能当前已关闭
暂无评论,快来抢沙发吧!