分布式事务中本地消息异步确认机制解析

0 次阅读

分布式系统中的业务操作通常会同时涉及多个服务,例如订单创建后需要扣减库存、记录支付状态并发送通知。传统单体应用可以借助数据库事务保证这些操作的原子性,但服务拆分后,每个服务往往拥有独立数据库,单一数据库事务已经无法覆盖完整业务流程。此时,本地消息表结合异步确认机制成为一种常见的分布式事务最终一致性方案。

一、本地消息异步确认机制是什么

本地消息异步确认机制的核心思想,是把“业务数据变更”和“待发送消息”放进同一个本地数据库事务中。

例如订单服务创建订单时,需要通知库存服务扣减库存。订单服务不直接依赖远程调用结果,而是执行:

  1. 创建订单;

  2. 在本地消息表写入一条待发送消息;

  3. 提交数据库事务。

由于订单数据和消息记录属于同一个本地事务,只要事务成功提交,两者就一定同时存在。随后由独立的消息投递线程、任务调度器或消息中间件负责异步发送。

简化后的业务结构如下:

订单服务
   │
   ├── 写入订单表
   │
   ├── 写入本地消息表
   │
   └── 提交本地事务
          │
          ▼
      异步消息投递
          │
          ▼
      消息中间件
          │
          ▼
      库存服务

这里需要特别注意,“本地消息”并不等于消息已经成功发送给下游服务。它只是说明消息发送任务已经和本地业务数据绑定,并具备后续可靠投递的依据。

二、为什么需要异步确认

分布式事务中最麻烦的问题之一,是业务数据已经提交,但远程操作可能失败。

假设订单创建成功后直接调用库存服务:

创建订单
   ↓
调用库存服务
   ↓
扣减库存

如果订单数据库事务已经提交,但调用库存服务时发生网络超时,就会出现状态不确定:

  • 订单可能已经创建成功;

  • 库存服务可能没有收到请求;

  • 也可能已经扣减库存,但响应没有返回;

  • 调用方无法仅凭超时结果判断最终状态。

如果为了保证一致性而长时间等待远程调用,系统又会产生线程阻塞、事务持锁时间过长以及服务级联故障等问题。

本地消息异步确认机制将这个过程拆成两个阶段:

阶段一:本地事务
业务数据 + 消息记录
        ↓
    一起提交

阶段二:异步投递
消息记录
   ↓
发送消息
   ↓
等待确认
   ↓
成功 / 重试

这样既避免了跨服务长事务,又能够通过消息重试机制提高最终一致性。

三、核心数据结构设计

典型的本地消息表可以设计为:

SQL
CREATE TABLE local_message (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    message_id VARCHAR(64) NOT NULL UNIQUE,
    business_id VARCHAR(64) NOT NULL,
    message_type VARCHAR(100) NOT NULL,
    payload TEXT NOT NULL,
    status TINYINT NOT NULL DEFAULT 0,
    retry_count INT NOT NULL DEFAULT 0,
    next_retry_time DATETIME NULL,
    created_at DATETIME NOT NULL,
    updated_at DATETIME NOT NULL
);

其中几个字段非常关键。

1. message_id

消息唯一标识,用于消息幂等和链路追踪。

不要简单使用业务主键作为消息ID,因为一个业务对象可能产生多种不同类型的消息。

2. business_id

关联订单号、支付单号、用户ID等业务标识,方便排查问题。

3. status

用于记录消息当前状态,例如:

0:待发送
1:发送中
2:发送成功
3:等待重试
4:最终失败

实际项目可以根据业务需求进一步细分。

4. retry_count

记录已经重试的次数,防止异常消息无限重试。

5. next_retry_time

实现延迟重试的重要字段,可以根据指数退避策略计算下一次发送时间。

四、本地事务保证消息不丢失

本地消息机制最重要的设计原则是:

业务数据与消息记录必须处于同一个本地事务中。

例如:

Java
@Transactional
public void createOrder(OrderRequest request) {
    Order order = createOrder(request);

    LocalMessage message = new LocalMessage();
    message.setMessageId(UUID.randomUUID().toString());
    message.setBusinessId(order.getId().toString());
    message.setMessageType("ORDER_CREATED");
    message.setPayload(buildPayload(order));

    localMessageRepository.save(message);
}

数据库事务提交后,订单和消息同时存在。

如果订单插入成功、消息插入失败,事务整体回滚;如果消息插入成功、订单插入失败,同样整体回滚。

因此可以避免以下典型问题:

订单创建成功
       ↓
服务宕机
       ↓
消息没有任何记录
       ↓
库存永远不会收到通知

本地消息表实际上为后续的可靠投递提供了“事实依据”。

五、异步发送与确认流程

消息记录成功后,需要有后台任务持续扫描待发送消息。

伪代码可以表示为:

Java
public void sendPendingMessages() {
    List<LocalMessage> messages = messageRepository.findPendingMessages();

    for (LocalMessage message : messages) {
        try {
            sendMessage(message);
            message.markSuccess();
            messageRepository.update(message);
        } catch (Exception e) {
            message.increaseRetryCount();
            message.calculateNextRetryTime();
            messageRepository.update(message);
        }
    }
}

完整流程通常是:

本地事务提交
      ↓
消息进入待发送状态
      ↓
消息扫描任务发现消息
      ↓
发送到MQ或调用可靠消息服务
      ↓
等待发送确认
   ↙       ↘
成功       失败
 ↓          ↓
已确认    重试
            ↓
       超过最大次数?
        ↙       ↘
       否        是
       ↓          ↓
    继续重试    人工介入

这里的“确认”通常不是简单理解为业务消费者已经执行成功,而是首先确认消息是否已经可靠地交给消息传输层。消费者业务执行是否成功,还需要通过消费端的幂等、重试以及状态补偿机制解决。

六、发送成功但状态更新失败怎么办

这是本地消息异步确认机制中一个非常容易被忽略的问题。

假设程序执行:

1. 发送消息成功
2. 更新本地消息状态为 SUCCESS
3. 数据库更新失败

此时消息实际上已经发送成功,但本地记录仍然显示“待发送”。

后台任务再次扫描后,就可能重复发送。

因此,本地消息机制不能保证消息绝对只发送一次

更合理的设计是:

允许消息重复投递,但必须保证消费端幂等。

例如库存服务收到:

message_id = 202608300001

第一次处理成功后记录消费状态:

SQL
CREATE TABLE consumed_message (
    message_id VARCHAR(64) PRIMARY KEY,
    consumed_at DATETIME NOT NULL
);

再次收到相同消息时:

查询message_id
    ↓
已经处理?
 ↙       ↘
是        否
↓          ↓
直接返回   执行业务
            ↓
        记录消费结果

这样即使消息重复投递,也不会造成重复扣库存、重复扣款等业务问题。

七、消息确认机制中的关键问题

1. 消息发送成功,不代表业务处理成功

生产者收到MQ确认,只能说明消息已经被消息系统接受,不能代表下游业务一定完成。

例如:

订单服务
   ↓
MQ确认成功
   ↓
库存服务
   ↓
数据库异常

因此完整的一致性体系至少包含两个层面:

  • 生产端可靠投递;

  • 消费端可靠处理。

2. 网络超时不代表发送失败

这是分布式系统中非常典型的“未知状态”。

例如生产者发送消息后:

生产者 → MQ
          ↓
       已成功接收
          ↓
       返回确认
          X
       网络断开

生产者看到的是超时,但消息可能已经进入MQ。

如果此时直接删除消息,就可能导致消息丢失;如果重新发送,则可能产生重复消息。

因此遇到超时,通常应该按照“可能成功”的原则处理,通过重试和消费幂等解决重复问题。

3. 消息状态更新必须具备并发控制

如果部署多个实例,可能出现多个节点同时扫描同一条消息:

实例A ─┐
       ├── 同一消息
实例B ─┘

可以通过数据库锁、乐观锁、状态抢占等方式控制。

例如增加版本字段:

SQL
UPDATE local_message
SET status = 1,
    version = version + 1
WHERE id = ?
  AND status = 0
  AND version = ?;

只有更新成功的实例才能获得消息处理权。

对于高并发场景,也可以通过分片扫描、任务队列等方式减少数据库竞争。

八、失败重试策略如何设计

简单的固定间隔重试:

失败 → 5秒后重试
失败 → 5秒后重试
失败 → 5秒后重试

在大量消息同时失败时容易形成重试风暴。

更推荐指数退避:

第1次:5秒
第2次:10秒
第3次:20秒
第4次:40秒
第5次:80秒

同时设置最大重试次数,例如:

retry_count >= 10

就不再自动发送,而是进入异常消息队列。

对于重要业务,还可以增加随机抖动,让不同消息的重试时间错开,避免大量任务同时触发。

九、最终失败消息不能简单删除

消息经过多次重试仍然失败时,不建议直接删除记录。

可以将其设置为:

FAILED

同时保留:

  • 消息ID;

  • 业务ID;

  • 最后一次错误原因;

  • 重试次数;

  • 最后更新时间;

  • 原始消息内容。

然后通过管理后台提供人工重试能力。

例如:

异常消息
   ↓
失败状态
   ↓
监控告警
   ↓
工程师查看原因
   ↓
修复下游问题
   ↓
重新投递

这对于支付、订单、库存等关键业务尤其重要。

十、本地消息异步确认与事务消息的区别

两者都用于解决分布式环境下的消息可靠性问题,但实现方式存在差异。

本地消息表主要依赖:

业务数据库
+
本地消息表
+
后台投递任务
+
消息中间件
+
消费幂等

事务消息则通常由消息中间件提供事务状态协调能力。

本地消息表的优点是实现逻辑清晰、技术门槛相对较低,而且对现有数据库体系改造较容易。

不足之处是需要自行维护消息表、扫描任务、重试策略、状态管理和异常告警。当消息量较大时,本地消息表还可能给业务数据库带来额外压力。

十一、如何保证本地消息表本身不成为性能瓶颈

消息表属于高频写入和高频查询表,数据量增长后需要重点优化。

首先,为状态和重试时间建立合适的索引:

SQL
CREATE INDEX idx_status_retry
ON local_message(status, next_retry_time);

其次,不要每次扫描都查询整个消息表:

SQL
SELECT *
FROM local_message;

而应该限制条件和批量大小:

SQL
SELECT *
FROM local_message
WHERE status IN (0, 3)
  AND (next_retry_time IS NULL OR next_retry_time <= NOW())
ORDER BY id
LIMIT 100;

对于长期运行的系统,还应该考虑消息归档。

例如:

活跃消息表
    ↓
定期归档
    ↓
历史消息表

避免单张表无限增长。

十二、推荐的完整架构

一个相对完整的本地消息异步确认架构可以设计成:

                ┌──────────────┐
                │   业务请求    │
                └──────┬───────┘
                       ↓
              ┌─────────────────┐
              │   本地数据库事务 │
              │                 │
              │  业务数据       │
              │      +          │
              │  本地消息       │
              └────────┬────────┘
                       ↓
                 事务提交成功
                       ↓
              ┌─────────────────┐
              │ 消息投递任务     │
              └────────┬────────┘
                       ↓
              ┌─────────────────┐
              │   消息中间件     │
              └────────┬────────┘
                       ↓
              ┌─────────────────┐
              │   消费者服务     │
              └────────┬────────┘
                       ↓
                幂等检查与处理
                       ↓
                  业务事务提交

监控系统则应该覆盖整个链路:

消息堆积
发送失败
重试次数
消费失败
处理延迟
异常消息

只有把监控和人工补偿纳入整体设计,分布式事务方案才真正具备生产可用性。

十三、适用场景与不适用场景

本地消息异步确认机制比较适合最终一致性要求较高、但不要求多个服务同时提交的业务。

常见场景包括:

  • 订单创建后通知库存服务;

  • 用户注册后发送欢迎消息;

  • 商品变更后同步搜索索引;

  • 支付完成后更新订单状态;

  • 数据变更后刷新缓存;

  • 业务操作完成后发送异步通知。

但如果业务要求多个资源必须在同一时刻完成提交,例如严格金融结算场景,就不能简单依赖本地消息表,需要结合业务特点选择更严格的分布式事务方案。

十四、实施时容易踩的几个坑

只做生产端,不做消费端幂等。

这是最常见的问题。消息重试本身就意味着重复投递是正常情况,消费者必须能够安全处理重复消息。

把发送成功和业务成功混为一谈。

MQ确认只解决消息是否可靠进入消息系统的问题,无法直接证明下游业务已经完成。

无限重试。

下游服务持续故障时,无限重试不仅不能解决问题,还可能拖垮数据库和消息系统。

失败后直接删除消息。

删除意味着失去恢复依据。对于重要业务,应保留失败记录并提供补偿能力。

事务提交前发送消息。

如果消息发送成功后本地事务回滚,就可能产生“消息存在、业务数据不存在”的问题。因此核心业务场景下应优先保证业务数据与消息记录在同一事务中提交。

十五、最佳实践总结

设计分布式事务中的本地消息异步确认机制,可以遵循以下原则:

  1. 业务数据和本地消息必须使用同一个数据库事务;

  2. 消息必须拥有全局唯一的消息ID;

  3. 消息发送采用异步方式,避免阻塞主业务流程;

  4. 发送状态需要持久化;

  5. 网络超时不能简单判定为消息发送失败;

  6. 允许消息重复投递,通过幂等机制保证业务安全;

  7. 使用指数退避控制重试频率;

  8. 设置最大重试次数;

  9. 超过重试次数后进入异常状态并触发告警;

  10. 提供人工补偿和重新投递能力;

  11. 对消息表建立合理索引并定期归档;

  12. 对消息发送、消费和延迟建立完整监控。

本地消息异步确认机制的本质,并不是追求“消息只发送一次”,而是通过本地事务保证消息记录不丢失、可靠投递保证消息最终送达、消费幂等保证重复消息安全、失败重试与人工补偿保证异常可恢复,最终形成一套完整的分布式事务最终一致性闭环。