一、双写为何会丢失或制造幽灵事件
业务服务先写 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 用单库事务消除双写丢失窗口,再以至少一次投递、消费者幂等和可观测重放完成端到端可靠性。工程重点在事务边界、认领租约、分区顺序、契约版本和故障注入,而不只是创建一张表。