为什么 Kafka ACK 不能替代 Outbox 模式

前言

在使用 Kafka 构建事件驱动系统时,一个非常常见的误解是:

Kafka 已经支持 ACK 确认机制,只要生产者收到 Kafka 的发送确认,就能够保证数据库操作和消息发送的一致性,因此没有必要再使用 Outbox。

这个结论看起来很合理。

Kafka Producer 可以通过 acks=all 确认消息已经写入 Kafka;Kafka Consumer 也可以在业务处理成功后再提交 Offset。既然生产端和消费端都有确认机制,为什么还需要额外设计一张 Outbox 表?

问题在于:

Kafka ACK 只能确认 Kafka 内部的消息写入结果,无法保证数据库事务与 Kafka 消息发送之间的原子性

  • Outbox Pattern 解决的是数据库事务和消息发送的一致性问题
  • Kafka ACK 解决的是消息是否已经被 Kafka 接收并持久化的问题

两者解决的根本不是同一个问题。

一、先看一个典型业务场景

假设我们有一个订单服务。

用户创建订单时,需要完成两个操作:

  1. 将订单写入 MySQL;
  2. 向 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。

常见实现方式有两种:

  1. 轮询 Outbox 表;
  2. 使用 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,原因包括:

  1. Kafka ACK 无法控制数据库事务;
  2. 数据库提交和 Kafka ACK 是两个独立提交点;
  3. 两个提交点之间存在不可消除的故障窗口;
  4. 数据库提交后 Kafka 可能发送失败;
  5. Kafka 发送成功后数据库可能回滚;
  6. 应用进程可能在任意中间状态崩溃;
  7. Producer 收到超时时,消息结果甚至可能是不确定的。

Outbox 的核心价值是:

先用本地事务可靠记录业务事实,
再通过可重试机制把业务事实传播到 Kafka。

最终,一个可靠的事件驱动系统通常不是只依赖某一种机制,而是组合使用:

本地事务
+ Outbox
+ Kafka ACK
+ 失败重试
+ 消费者幂等
+ 可观测性

因此,面对“为什么 Kafka ACK 不能替代 Outbox”这个问题,最准确的回答是:

Kafka ACK 只能证明 Kafka 接收了消息,不能证明数据库事务和 Kafka 消息同时成功。Outbox 通过将业务数据和事件记录写入同一个本地事务,解决了数据库与消息系统之间的双写一致性问题。

使用 Hugo 构建
主题 StackJimmy 设计