前言
在使用 Kafka 构建事件驱动系统时,一个非常常见的误解是:
Kafka 已经支持 ACK 确认机制,只要生产者收到 Kafka 的发送确认,就能够保证数据库操作和消息发送的一致性,因此没有必要再使用 Outbox。
这个结论看起来很合理。
Kafka Producer 可以通过 acks=all 确认消息已经写入 Kafka;Kafka Consumer 也可以在业务处理成功后再提交 Offset。既然生产端和消费端都有确认机制,为什么还需要额外设计一张 Outbox 表?
问题在于:
Kafka ACK 只能确认 Kafka 内部的消息写入结果,无法保证数据库事务与 Kafka 消息发送之间的原子性。
- Outbox Pattern 解决的是数据库事务和消息发送的一致性问题
- Kafka ACK 解决的是消息是否已经被 Kafka 接收并持久化的问题。
两者解决的根本不是同一个问题。
一、先看一个典型业务场景
假设我们有一个订单服务。
用户创建订单时,需要完成两个操作:
- 将订单写入 MySQL;
- 向 Kafka 发送一条
OrderCreated事件。
伪代码如下:
@Transactional
public void createOrder(CreateOrderCommand command) {
Order order = new Order();
order.setUserId(command.getUserId());
order.setAmount(command.getAmount());
orderRepository.save(order);
kafkaTemplate.send(
"order-created-topic",
order.getId(),
new OrderCreatedEvent(order.getId())
);
}
业务希望实现下面的语义:
- 订单创建成功等价于数据库中存在订单,并且 Kafka 中一定存在 OrderCreated 事件
也就是说:
- 数据库提交成功,消息必须发送成功
- 数据库回滚,消息不能被发送
这实际上是在要求两个不同系统之间保持原子性:
- MySQL
- Kafka
理想状态是:
- 要么全部成功
- 要么全部失败
但普通的本地事务只能控制 MySQL,不能同时控制 Kafka。
二、Kafka ACK 到底确认了什么?
Kafka Producer 发送消息时,可以通过 acks 参数控制 Broker 的确认条件。
常见配置如下:
spring:
kafka:
producer:
acks: all
Kafka 的 acks 主要有三个取值:
| 配置 | 含义 |
|---|---|
acks=0 |
Producer 不等待 Broker 确认 |
acks=1 |
Leader 写入本地日志后返回成功 |
acks=all |
Leader 等待 ISR 中满足条件的副本完成同步后返回成功 |
使用:
acks=all
能够显著提高消息的可靠性。
它能够回答的问题是:
这条消息是否已经按照 Kafka 的可靠性要求写入 Broker?
但是它无法回答:
这条 Kafka 消息对应的数据库事务是否已经提交?
Kafka Broker 并不知道 MySQL 的事务状态。
对于 Kafka 来说,它只负责处理这条消息:
OrderCreated(orderId=1001)
它不知道:
order_id = 1001
对应的订单记录是否已经成功写入数据库。
因此,Kafka ACK 能确认的是:
Kafka 消息写入成功
它不能确认:
数据库事务与 Kafka 消息同时成功
三、最核心的问题:双写一致性
同时写入数据库和 Kafka,本质上属于“双写”。
一次业务操作
├── 写 MySQL
└── 写 Kafka
只要两个操作不处于同一个原子事务中,就一定存在中间状态。
无论先写数据库还是先发送 Kafka,都存在失败窗口。
四、先写数据库,再发送 Kafka
第一种常见实现是:
@Transactional
public void createOrder(CreateOrderCommand command) {
Order order = orderRepository.save(...);
kafkaTemplate.send(
"order-created-topic",
new OrderCreatedEvent(order.getId())
);
}
逻辑顺序是:
1. 写入数据库
2. 发送 Kafka
3. 提交数据库事务
或者在某些实现中:
1. 写入数据库
2. 提交数据库事务
3. 发送 Kafka
这两种顺序都存在问题。
场景一:数据库成功,Kafka 发送失败
执行过程如下:
1. MySQL 事务提交成功
2. 发送 Kafka
3. 网络超时或 Kafka 不可用
4. 消息发送失败
最终状态:
MySQL:订单存在
Kafka:没有 OrderCreated 消息
这意味着下游系统无法感知订单创建。
例如:
- 库存服务没有扣减库存;
- 积分服务没有发放积分;
- 通知服务没有发送通知;
- 搜索系统没有建立订单索引。
即使使用:
acks=all
也无法解决这个问题。
因为此时 Kafka 根本没有返回成功 ACK。
数据库已经提交,而消息发送失败。Kafka ACK 只能告诉你“发送失败了”,却不能自动撤销已经提交的数据库事务。
场景二:Kafka 发送成功,数据库事务回滚
执行过程如下:
1. 写入订单数据
2. 发送 Kafka
3. Kafka 返回 ACK
4. 后续业务逻辑抛出异常
5. MySQL 事务回滚
最终状态:
MySQL:订单不存在
Kafka:存在 OrderCreated 消息
下游消费者收到消息后,可能开始处理一个根本不存在的订单。
例如:
graph TD
A[库存服务收到 OrderCreated] --> B[尝试锁定库存]
B --> C[订单服务数据库中却不存在该订单]
这里 Kafka ACK 完全正常。
Kafka 确实成功保存了消息,因此返回 ACK 没有任何问题。
真正的问题是:
Kafka 消息已经提交,但数据库事务最终回滚了。
Kafka ACK 无法预知未来的数据库事务是否会回滚。
五、先发送 Kafka,再提交数据库,同样不安全
有些开发者会尝试调整顺序:
@Transactional
public void createOrder(CreateOrderCommand command) {
kafkaTemplate.send(
"order-created-topic",
new OrderCreatedEvent(...)
).get();
orderRepository.save(...);
}
执行顺序变成:
1. 发送 Kafka
2. 等待 Kafka ACK
3. 写入数据库
4. 提交事务
这样虽然可以确保 Kafka 发送成功后再继续写数据库,但仍然存在问题。
执行过程可能是:
1. Kafka 消息发送成功
2. Kafka 返回 ACK
3. 数据库写入失败
4. 数据库事务回滚
最终状态仍然是:
MySQL:没有订单
Kafka:存在 OrderCreated 消息
Kafka ACK 只能证明:
Kafka 已经接收了消息
不能证明:
接下来的数据库操作一定成功
因此,仅仅调整执行顺序无法消除双写问题。
六、为什么捕获异常后补偿也不可靠?
一种常见做法是:
try {
orderRepository.save(order);
kafkaTemplate.send(...).get();
} catch (Exception exception) {
// 执行补偿逻辑
}
或者:
try {
kafkaTemplate.send(...).get();
orderRepository.save(order);
} catch (Exception exception) {
// 删除消息或者回滚订单
}
这种方式的问题是,补偿本身也可能失败。
例如:
graph TD
A[数据库提交成功] --> B[Kafka 发送失败]
B --> C[准备执行数据库补偿删除]
C --> D[应用进程崩溃]
最终依然是:
数据库中存在订单
Kafka 中没有消息
再比如:
graph TD
A[Kafka 发送成功] --> B[数据库提交失败]
B --> C[准备发送一条取消事件]
C --> D[服务器断电]
最终依然会留下不一致状态。
普通的 try-catch 只能处理当前线程中能够捕获的异常,无法覆盖:
- JVM 崩溃;
- 服务器宕机;
- 容器被强制终止;
- 网络中断;
- 进程被
kill -9; - 数据库连接突然断开;
- Kafka 返回 ACK 后应用进程立即退出。
双写问题最难处理的地方,不是普通异常,而是:
两个操作之间存在不可消除的故障窗口。
七、Kafka ACK 不等于业务事务提交
可以把 Kafka ACK 理解为 Kafka 自己的局部确认。
Kafka 返回 ACK 的含义类似于:
我 Kafka 已经按照你的配置接收并保存了这条消息。
但它不会承诺:
你的数据库事务也已经成功提交。
因为 Kafka 无法控制 MySQL。
同样,MySQL 提交成功的含义是:
数据库中的修改已经提交。
但 MySQL 也不会承诺:
Kafka 中对应的消息一定已经存在。
因为 MySQL 无法控制 Kafka。
因此,系统中实际存在两个独立的提交点:
MySQL COMMIT
Kafka ACK
两个提交点之间一定存在时间差。
graph TD
A[MySQL COMMIT] --> B[故障窗口]
B --> C[Kafka ACK]
或者:
graph TD
A[Kafka ACK] --> B[故障窗口]
B --> C[MySQL COMMIT]
只要故障发生在这个窗口内,就会出现数据不一致。
八、Kafka ACK 还可能存在“不确定结果”
除了明确成功和明确失败之外,分布式系统中还存在第三种状态:
结果未知
例如:
1. Producer 将消息发送给 Kafka
2. Kafka 成功写入消息
3. Kafka 返回 ACK
4. ACK 在网络中丢失
5. Producer 等待超时
此时 Kafka 中的真实状态是:
消息已经写入
但 Producer 观察到的状态是:
发送超时
Producer 无法确定消息到底有没有成功。
于是 Producer 可能重试:
第一次消息:已经成功写入
第二次消息:重试后再次写入
如果没有幂等机制,就可能产生重复消息。
即使启用了 Kafka 幂等生产者:
spring:
kafka:
producer:
properties:
enable.idempotence: true
它主要解决的也是 Producer 重试导致的 Kafka 内部重复写入问题。
它依然无法保证:
MySQL 事务和 Kafka 消息原子提交
这再次说明:
Kafka 的可靠性机制主要作用于 Kafka 内部,而 Outbox 解决的是跨系统的一致性问题。
九、Outbox Pattern 是什么?
Outbox Pattern,也叫事务发件箱模式。
核心思想是:
不在业务事务中直接发送 Kafka,而是将“业务数据”和“待发送消息”写入同一个数据库事务。
例如,订单服务有两张表:
orders
outbox_event
创建订单时:
@Transactional
public void createOrder(CreateOrderCommand command) {
Order order = new Order();
order.setUserId(command.getUserId());
order.setAmount(command.getAmount());
orderRepository.save(order);
OutboxEvent event = new OutboxEvent();
event.setAggregateType("ORDER");
event.setAggregateId(order.getId().toString());
event.setEventType("OrderCreated");
event.setPayload(toJson(new OrderCreatedEvent(order.getId())));
event.setStatus("PENDING");
outboxEventRepository.save(event);
}
此时本地事务中只操作 MySQL:
BEGIN
INSERT INTO orders ...
INSERT INTO outbox_event ...
COMMIT
由于订单记录和 Outbox 事件位于同一个数据库事务中,因此可以实现原子性:
要么订单和事件记录同时提交
要么订单和事件记录同时回滚
这解决了最关键的问题:
数据库中不会出现“订单存在但没有事件记录”的状态
十、Outbox 如何把消息发送到 Kafka?
Outbox 表中的消息提交成功后,由独立的消息发布器负责发送到 Kafka。
常见实现方式有两种:
- 轮询 Outbox 表;
- 使用 CDC 监听数据库变更。
十一、方案一:定时轮询 Outbox 表
发布器定期查询待发送事件:
SELECT id,
aggregate_type,
aggregate_id,
event_type,
payload
FROM outbox_event
WHERE status = 'PENDING'
ORDER BY id
LIMIT 100;
发送成功后更新状态:
UPDATE outbox_event
SET status = 'SENT',
sent_at = NOW()
WHERE id = ?;
示例代码:
@Scheduled(fixedDelay = 1000)
public void publishEvents() {
List<OutboxEvent> events =
outboxEventRepository.findPendingEvents(100);
for (OutboxEvent event : events) {
try {
kafkaTemplate.send(
event.getTopic(),
event.getAggregateId(),
event.getPayload()
).get();
outboxEventRepository.markAsSent(event.getId());
} catch (Exception exception) {
outboxEventRepository.recordFailure(
event.getId(),
exception.getMessage()
);
}
}
}
如果 Kafka 暂时不可用:
Outbox 事件仍然保留在数据库中
下一次轮询可以继续重试。
因此,系统将:
一次不可靠的即时双写
转换为:
一次可靠的本地事务
+
一个可持久化、可重试的异步发送过程
十二、方案二:CDC
CDC 是 Change Data Capture,即变更数据捕获。
可以通过 Debezium 等工具监听数据库 Binlog。
执行过程如下:
graph TD
A[业务事务] --> B[写入 orders]
B --> C[写入 outbox_event]
C --> D[提交 MySQL 事务]
D --> E[Binlog 记录变更]
E --> F[Debezium 捕获 Outbox 记录]
F --> G[发送到 Kafka]
这种方案通常不需要应用程序主动轮询数据库。
架构如下:
┌──────────────────┐
│ Order Service │
└────────┬─────────┘
│ 本地事务
▼
┌──────────────────┐
│ MySQL │
│ │
│ orders │
│ outbox_event │
└────────┬─────────┘
│ Binlog
▼
┌──────────────────┐
│ Debezium │
└────────┬─────────┘
│
▼
┌──────────────────┐
│ Kafka │
└──────────────────┘
CDC Outbox 通常具有以下优势:
- 实时性更高;
- 减少数据库轮询;
- 应用逻辑更简单;
- 更适合大规模事件驱动架构。
但它也会增加基础设施复杂度。
十三、Outbox 为什么比 ACK 更可靠?
Outbox 的关键不是“发送 Kafka 时更加可靠”,而是:
将必须保持一致的两个数据写入,收敛到同一个本地数据库事务中。
原来的双写模型:
MySQL + Kafka
变成:
MySQL 业务表 + MySQL Outbox 表
因为它们在同一个数据库中,可以使用同一个本地事务:
BEGIN
INSERT INTO orders ...
INSERT INTO outbox_event ...
COMMIT
此时不会存在:
订单提交成功,但事件记录丢失
之后从 Outbox 发送 Kafka,即使失败,也可以通过持久化状态不断重试。
这是一种非常重要的架构思想:
不要试图让两个无法原子提交的系统强行同时成功,而是先在一个可靠事务中记录事实,再异步传播这个事实。
十四、Kafka ACK 在 Outbox 中仍然有价值
需要注意的是,使用 Outbox 并不意味着 Kafka ACK 没有价值。
Outbox 和 Kafka ACK 不是互斥关系,而是互补关系。
在 Outbox 发布器中,仍然需要等待 Kafka ACK:
kafkaTemplate.send(
topic,
key,
payload
).get();
只有收到 Kafka 的成功确认后,才能将 Outbox 事件标记为已发送。
正确流程是:
graph TD
A[读取 PENDING 事件] --> B[发送 Kafka]
B --> C[等待 Kafka ACK]
C --> D[更新为 SENT]
因此:
Outbox 负责保证业务数据和事件记录的一致性
Kafka ACK 负责确认事件是否成功写入 Kafka
两者职责如下:
| 机制 | 解决的问题 |
|---|---|
| 数据库本地事务 | 业务表与 Outbox 表原子提交 |
| Outbox | 数据库状态变更与消息发送最终一致 |
| Kafka ACK | 确认 Kafka 是否成功接收消息 |
| Kafka 幂等 Producer | 降低 Producer 重试产生重复消息的风险 |
| 消费者幂等 | 避免重复消费造成重复业务处理 |
所以正确的理解不是:
ACK 和 Outbox 二选一
而是:
Outbox + ACK + 重试 + 幂等
十五、Outbox 也无法天然保证消息只发送一次
Outbox 通常保证的是:
At Least Once
也就是至少发送一次,而不是严格的 Exactly Once。
考虑下面的场景:
1. 发布器读取 Outbox 事件
2. 成功发送到 Kafka
3. Kafka 返回 ACK
4. 发布器准备更新状态为 SENT
5. 应用进程崩溃
此时:
Kafka:消息已经存在
Outbox:状态仍然是 PENDING
应用重启后,会再次发送这条消息。
于是 Kafka 中可能出现重复事件。
这说明 Outbox 解决的是:
消息不能丢
但通常不会直接保证:
消息绝不重复
因此下游消费者必须具备幂等能力。
十六、消费者如何实现幂等?
每条事件都应该具有唯一标识:
{
"eventId": "01JXABCDEF123456",
"eventType": "OrderCreated",
"aggregateId": "1001",
"occurredAt": "2026-07-15T10:30:00Z",
"payload": {
"orderId": 1001
}
}
消费者处理前,先检查 eventId 是否已经处理。
例如建立事件消费记录表:
CREATE TABLE consumed_event (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
event_id VARCHAR(64) NOT NULL,
consumer_name VARCHAR(128) NOT NULL,
consumed_at DATETIME NOT NULL,
UNIQUE KEY uk_event_consumer (
event_id,
consumer_name
)
);
消费时:
@Transactional
public void handle(OrderCreatedEvent event) {
boolean inserted = consumedEventRepository.tryInsert(
event.getEventId(),
"inventory-service"
);
if (!inserted) {
return;
}
inventoryService.reserve(event.getOrderId());
}
依靠唯一索引:
event_id + consumer_name
避免同一消费者重复处理同一事件。
十七、Kafka Consumer ACK 也不能替代 Outbox
有时开发者提到的 ACK,并不是 Producer 的 acks,而是消费者手动确认。
例如:
@KafkaListener(topics = "order-created-topic")
public void consume(
OrderCreatedEvent event,
Acknowledgment acknowledgment) {
orderHandler.handle(event);
acknowledgment.acknowledge();
}
消费者手动提交 Offset 能解决的问题是:
业务处理失败时,不要过早提交消费位点
它控制的是:
消费者是否认为这条消息已经处理完成
但它依然与生产端的数据库和 Kafka 双写无关。
消费者 ACK 发生在消息已经进入 Kafka 之后:
graph TD
A[生产者写数据库] --> B[生产者发送 Kafka]
B --> C[Kafka 保存消息]
C --> D[消费者读取消息]
D --> E[消费者提交 ACK]
如果消息根本没有成功发送到 Kafka,消费者就没有机会处理,更谈不上提交 ACK。
因此,消费者 ACK 解决的是消费可靠性问题,而 Outbox 解决的是生产端事件发布一致性问题。
十八、几种机制分别解决什么问题?
可以通过下面这张表进行区分。
| 机制 | 作用范围 | 主要解决的问题 |
|---|---|---|
| Kafka Producer ACK | Producer 到 Kafka | 消息是否成功写入 Kafka |
| Kafka Producer Retry | Producer 到 Kafka | 临时故障下自动重试 |
| Kafka Idempotent Producer | Kafka 内部 | 降低重试导致的重复写入 |
| Kafka Transaction | Kafka 内部或 Kafka 到 Kafka | 多个 Kafka 写入及 Offset 提交的一致性 |
| Consumer Offset Commit | Kafka Consumer | 控制消息消费进度 |
| Consumer Manual ACK | Kafka Consumer | 业务处理成功后再提交位点 |
| Outbox Pattern | 数据库到消息系统 | 业务数据与事件发布最终一致 |
| 消费者幂等 | 消费者业务层 | 防止重复消息造成重复业务操作 |
这些机制各自解决不同层面的问题,不能简单互相替代。
十九、Kafka Transaction 能不能替代 Outbox?
Kafka 支持事务生产者,例如:
kafkaTemplate.executeInTransaction(operations -> {
operations.send("topic-a", messageA);
operations.send("topic-b", messageB);
return true;
});
Kafka Transaction 可以保证:
多条 Kafka 消息要么都可见,要么都不可见
也可以在消费—处理—再生产模型中,实现:
消费 Offset 提交
+
发送新的 Kafka 消息
之间的原子性。
例如:
graph TD
A[Kafka Topic A] --> B[Consumer]
B --> C[Kafka Topic B]
Kafka Transaction 很适合 Kafka 内部的流式处理。
但是普通 Kafka Transaction 不能直接把 MySQL 本地事务纳入同一个原子提交中。
系统仍然存在:
MySQL COMMIT
Kafka Transaction COMMIT
两个独立提交点。
因此:
Kafka Transaction
不能天然保证:
MySQL + Kafka
的跨系统原子性。
只要业务核心数据存储在 MySQL、PostgreSQL 等外部数据库中,Outbox 仍然是非常常见且可靠的解决方案。
二十、Spring 的事务同步回调能替代 Outbox 吗?
有些项目会使用:
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronization() {
@Override
public void afterCommit() {
kafkaTemplate.send(...);
}
}
);
这种做法可以确保:
只有数据库事务提交成功后,才发送 Kafka
它确实避免了:
数据库回滚,但 Kafka 消息已经发送
但仍然无法解决:
graph TD
A[数据库已经提交] --> B[应用崩溃]
B --> C[afterCommit 逻辑未执行]
或者:
graph TD
A[数据库已经提交] --> B[发送 Kafka 失败]
最终状态仍然可能是:
数据库有数据
Kafka 无消息
所以 afterCommit 适合:
- 缓存清理;
- 本地通知;
- 非关键异步操作;
- 允许偶尔丢失的逻辑。
但对于订单、支付、库存等不能丢失的领域事件,afterCommit 不能代替 Outbox。
二十一、典型 Outbox 表设计
一个相对完整的 Outbox 表可以设计如下:
CREATE TABLE outbox_event (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
event_id VARCHAR(64) NOT NULL,
aggregate_type VARCHAR(64) NOT NULL,
aggregate_id VARCHAR(64) NOT NULL,
event_type VARCHAR(128) NOT NULL,
topic VARCHAR(128) NOT NULL,
message_key VARCHAR(128) DEFAULT NULL,
payload JSON NOT NULL,
status VARCHAR(32) NOT NULL DEFAULT 'PENDING',
retry_count INT NOT NULL DEFAULT 0,
next_retry_at DATETIME DEFAULT NULL,
sent_at DATETIME DEFAULT NULL,
created_at DATETIME NOT NULL,
updated_at DATETIME NOT NULL,
UNIQUE KEY uk_event_id (event_id),
KEY idx_publish_scan (
status,
next_retry_at,
id
),
KEY idx_aggregate (
aggregate_type,
aggregate_id
)
);
字段说明:
| 字段 | 作用 |
|---|---|
event_id |
事件唯一标识,用于幂等 |
aggregate_type |
聚合类型,例如 ORDER |
aggregate_id |
聚合根 ID,例如订单 ID |
event_type |
事件类型,例如 OrderCreated |
topic |
Kafka Topic |
message_key |
Kafka 消息 Key |
payload |
事件内容 |
status |
发送状态 |
retry_count |
重试次数 |
next_retry_at |
下次重试时间 |
sent_at |
成功发送时间 |
状态可以定义为:
PENDING
PROCESSING
SENT
FAILED
二十二、轮询 Outbox 时需要注意并发问题
如果部署了多个消息发布器实例,必须避免多个实例同时发送同一条事件。
可以使用数据库行锁:
SELECT *
FROM outbox_event
WHERE status = 'PENDING'
AND (
next_retry_at IS NULL
OR next_retry_at <= NOW()
)
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED;
处理流程:
1. 开启事务
2. 使用 FOR UPDATE SKIP LOCKED 获取任务
3. 将状态改为 PROCESSING
4. 提交事务
5. 发送 Kafka
6. 成功后改为 SENT
7. 失败后恢复为 PENDING 或改为 FAILED
也可以通过条件更新抢占任务:
UPDATE outbox_event
SET status = 'PROCESSING',
updated_at = NOW()
WHERE id = ?
AND status = 'PENDING';
只有受影响行数为 1 的实例才能继续处理。
但即便有任务抢占机制,仍然不能完全消除“发送成功但状态更新失败”带来的重复消息,所以消费者幂等依然不可省略。
二十三、失败重试不能无限立即执行
如果 Kafka 长时间不可用,发布器不能无间隔地重试。
可以使用指数退避:
第 1 次失败:10 秒后重试
第 2 次失败:30 秒后重试
第 3 次失败:1 分钟后重试
第 4 次失败:5 分钟后重试
第 5 次失败:30 分钟后重试
示例:
private Duration calculateBackoff(int retryCount) {
long seconds = Math.min(
1800,
(long) Math.pow(2, retryCount) * 10
);
return Duration.ofSeconds(seconds);
}
超过最大重试次数后,可以将事件标记为:
FAILED
并触发告警或进入人工补偿流程。
二十四、Outbox 需要做好监控
Outbox 不只是增加一张表,还需要配套可观测性。
至少应监控以下指标:
outbox_pending_count
outbox_failed_count
outbox_publish_success_total
outbox_publish_failure_total
outbox_publish_latency
outbox_oldest_pending_age
其中非常关键的是:
最老一条 PENDING 事件已经积压了多久
因为即使待发送数量不多,也可能有某一条关键事件长期无法发送。
建议设置告警:
PENDING 数量超过阈值
FAILED 数量大于 0
最老待发送事件超过 5 分钟
连续发送失败率过高
二十五、什么时候不一定需要 Outbox?
并不是所有 Kafka 消息都必须使用 Outbox。
如果消息属于以下类型,可以根据业务容忍度简化:
- 日志消息;
- 非关键埋点;
- 用户行为分析;
- 可丢失的监控事件;
- 非关键缓存刷新通知;
- 可以通过定时任务重新计算的数据。
例如:
用户打开了某个页面
即使偶尔丢失一条埋点,通常不会破坏核心业务状态。
此时直接发送 Kafka,并结合 ACK 和重试,可能已经足够。
但如果消息承载的是关键业务事实,例如:
- 订单已创建;
- 支付已成功;
- 库存已锁定;
- 退款已完成;
- 账户余额已变更;
- 优惠券已核销;
那么消息丢失通常会造成严重的数据不一致。
这类场景更适合使用 Outbox。
二十六、什么时候应该优先使用 Outbox?
当系统同时满足以下条件时,应该认真考虑 Outbox:
1. 核心业务数据写入数据库
2. 数据变更后必须发布事件
3. 事件不能丢失
4. 数据库与 Kafka 不属于同一个事务资源
5. 可以接受最终一致性
典型场景包括:
订单创建 → 库存锁定
支付成功 → 订单状态更新
退款成功 → 账户余额恢复
用户注册 → 初始化用户权益
商品更新 → 搜索索引同步
这些场景的共同特点是:
消息不是附加通知,而是业务状态传播链路的一部分。
二十七、正确的系统设计应该是什么?
一个完整的可靠事件发布链路通常包括:
graph TD
A[业务本地事务] --> B[业务表 + Outbox 表原子写入]
B --> C[Outbox Publisher 或 CDC]
C --> D[Kafka ACK 确认]
D --> E[失败重试]
E --> F[Kafka 消费者]
F --> G[消费者幂等处理]
G --> H[成功后提交 Offset]
对应的可靠性机制如下:
- 数据库本地事务保证业务记录和事件记录不会分离
- Outbox 保证消息发送失败后仍然能够重试
- Kafka ACK 确认消息是否成功写入 Kafka
- Producer Retry 处理临时网络和 Broker 故障
- 消费者幂等处理至少一次投递带来的重复消息
- Offset Commit 控制消费者处理进度
缺少任何一个环节,都可能留下可靠性缺口。
二十八、总结
Kafka ACK 和 Outbox 的根本区别在于它们解决的问题不同。
Kafka ACK 解决的是:
消息有没有成功写入 Kafka
Outbox 解决的是:
数据库业务状态和待发布事件能否保持一致
Kafka ACK 无法替代 Outbox,原因包括:
- Kafka ACK 无法控制数据库事务;
- 数据库提交和 Kafka ACK 是两个独立提交点;
- 两个提交点之间存在不可消除的故障窗口;
- 数据库提交后 Kafka 可能发送失败;
- Kafka 发送成功后数据库可能回滚;
- 应用进程可能在任意中间状态崩溃;
- Producer 收到超时时,消息结果甚至可能是不确定的。
Outbox 的核心价值是:
先用本地事务可靠记录业务事实,
再通过可重试机制把业务事实传播到 Kafka。
最终,一个可靠的事件驱动系统通常不是只依赖某一种机制,而是组合使用:
本地事务
+ Outbox
+ Kafka ACK
+ 失败重试
+ 消费者幂等
+ 可观测性
因此,面对“为什么 Kafka ACK 不能替代 Outbox”这个问题,最准确的回答是:
Kafka ACK 只能证明 Kafka 接收了消息,不能证明数据库事务和 Kafka 消息同时成功。Outbox 通过将业务数据和事件记录写入同一个本地事务,解决了数据库与消息系统之间的双写一致性问题。