本文目录导读:

- 核心痛点与解决思路
- 分层级的可靠发送
- 领域驱动设计:基于聚合根的事务性发件箱(重点案例)
- 代码示例(基于Java + Spring Boot + RabbitMQ)
- 高级方案:基于消息中间件本身的机制
- 总结与最佳实践
我将为你详细讲解消息可靠发送的完整案例,涵盖从理论基础到代码实现的各个方面。
消息可靠发送是分布式系统和微服务架构中的核心难题,这里的“可靠”通常指不丢失、不重复(或至少能处理重复)且顺序正确。
核心痛点与解决思路
在分布式系统中,消息发送面临三大挑战:
- 网络不可靠:发送方发出消息,可能因网络超时而不知道是否送达。
- 进程故障:发送方或接收方在发送过程中可能宕机。
- 消息重复:由于重试机制,接收方可能收到重复消息。
核心解决思路是 “三管齐下”:
- 确认机制:接收方收到消息后必须发送确认(ACK)。
- 重试机制:发送方未收到ACK时,自动重发。
- 幂等性设计:接收方必须有能力处理重复消息,保证最终结果一致。
分层级的可靠发送
根据业务需求,消息发送的可靠性可以分为三个层级:
最多一次 (At-most-once)
- 场景:允许丢失少量数据(如日志、监控)。
- 实现:发送后不等待确认,不重试。
- 风险:消息可能丢失。
最少一次 (At-least-once) —— 最常用
- 场景:核心业务数据(如订单、支付)。
- 实现:发送后等待ACK,超时未收到则重试。
- 优势:保证不丢。
- 风险:可能产生重复消息,需要接收方做幂等处理。
恰好一次 (Exactly-once)
- 场景:金融交易、高价值业务。
- 实现:依赖底层系统(如Kafka的幂等生产者、RocketMQ的事务消息),或通过“最少一次”+“下游幂等”来逻辑上实现。
领域驱动设计:基于聚合根的事务性发件箱(重点案例)
这是目前业界解决“本地事务与发送消息不一致”问题的最佳实践,也是实现可靠发送的典型案例。
问题背景: 用户下单场景,你需要在数据库中创建订单,并向MQ发送一条“订单创建”消息,如果先入库,再发消息,入库成功但发消息失败,就会导致订单数据不一致。
解决方案清单:
- 创建消息表:在业务数据库中建一张
outbox(发件箱)表。 - 同事务写库:在创建订单的同一个事务中,将订单数据和待发送的消息标记一起写入数据库。
- 独立定时任务扫描:一个后台定时任务(或CDC工具,如Debezium)扫描
outbox表,将状态为PENDING的消息发送给MQ。 - 确认并标记:收到MQ确认后,更新
outbox表的状态为SENT。
流程图如下:
sequenceDiagram
participant C as 业务服务
participant DB as 业务数据库
participant S as 定时任务/CDC
participant MQ as 消息队列
C->>DB: 1.开启事务
C->>DB: 2.插入订单表
C->>DB: 3.插入Outbox表(状态:待发送)
C->>DB: 4.提交事务
S->>DB: 5.查询待发送消息
alt 消息存在
S->>MQ: 6.发送消息
MQ-->>S: 7.返回ACK
S->>DB: 8.更新Outbox状态为已发送
else 发送失败
S->>S: 等待下次轮询重试
Note over S: 由于消息在DB中,不会丢失
end
代码示例(基于Java + Spring Boot + RabbitMQ)
这是一个简化的案例,演示了发件箱与本地事务的结合。
第一步:依赖配置
// pom.xml 关键依赖 (略)
// application.yml
server:
port: 8081
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
datasource:
url: jdbc:mysql://localhost:3306/test?useSSL=false
username: root
password: root
第二步:实现发件箱机制
实体类与Mapper(MyBatis-Plus)
import com.baomidou.mybatisplus.annotation.*;
import java.time.LocalDateTime;
// 数据库表对应实体
@Data
@TableName("outbox")
public class OutboxMessage {
@TableId(value = "id", type = IdType.ASSIGN_ID)
private Long id;
private String aggregateType; // 聚合类型,"Order"
private String aggregateId; // 聚合ID,例如订单号
private String payload; // 要发送的消息体(JSON格式)
@TableField(fill = FieldFill.INSERT)
private LocalDateTime createTime;
@TableField(fill = FieldFill.INSERT)
private String status; // PENDING -> SENDING -> SENT
}
核心服务:订单创建(关键点)
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import cn.hutool.json.JSONUtil;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import java.util.List;
import java.util.stream.Collectors;
@Service
public class OrderService {
private final OrderMapper orderMapper;
private final OutboxMapper outboxMapper;
private final RabbitTemplate rabbitTemplate;
// 构造函数注入...(省略)
/**
* 创建订单,并在同一个事务中写入发件箱
* 最终一致性:保证订单数据与消息绝对一致
*/
@Transactional(rollbackFor = Exception.class)
public void createOrder(Order order) {
// 1. 核心业务逻辑:保存订单
orderMapper.insert(order);
// 2. 专门发送消息:构建消息体
OrderCreatedEvent event = new OrderCreatedEvent(
order.getId(),
order.getGoodsId(),
order.getPrice()
);
// 3. 写入发件箱表(PENDING状态)
OutboxMessage message = new OutboxMessage();
message.setAggregateType("Order");
message.setAggregateId(order.getId().toString());
message.setPayload(JSONUtil.toJsonStr(event));
message.setStatus("PENDING");
// 注意:这两步在同一个数据库事务中
// 如果插入Order失败,Outbox也不会写入
outboxMapper.insert(message);
// 事务提交后,任务将由定时任务执行发送
}
}
发件箱轮询发送器(定时任务)
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import java.util.List;
@Component
public class OutboxPoller {
private final OutboxMapper outboxMapper;
private final RabbitTemplate rabbitTemplate;
// 注入...
@Scheduled(fixedDelay = 1000) // 每秒扫描一次
public void processMessages() {
// 1. 查询待发送数据(为了防止多实例并发,通常加 limit 和 乐观锁/悲观锁)
List<OutboxMessage> pendingMessages = outboxMapper.findByStatus("PENDING");
for (OutboxMessage message : pendingMessages) {
try {
// 2. 发送前更新状态(防止重复消费)
message.setStatus("SENDING");
outboxMapper.updateById(message);
// 3. 真正发送消息
rabbitTemplate.convertAndSend(
"order.exchange",
"order.created",
message.getPayload()
);
// 4. 发送成功,更新状态
message.setStatus("SENT");
outboxMapper.updateById(message);
} catch (Exception e) {
// 记录日志并跳过(等待下次轮询重试)
log.error("消息发送失败: {}", message.getId(), e);
// 保留状态为SENDING,为了防止无限重试,可设定重试次数上限
}
}
}
}
第三步:消息接收方(保证幂等性)
在消费者端,为了防止消息重复,必须在业务逻辑中进行幂等校验。
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
public class OrderConsumer {
@Autowired
private OrderStatusMapper orderStatusMapper;
@Autowired
private MessageLogMapper messageLogMapper;
@RabbitListener(queues = "order.created.queue")
public void handleOrderCreated(OrderCreatedEvent event) {
Long eventId = event.getId();
// 幂等性检查:确保当前消息没有被处理过
// 方法1:使用Redis
// if (stringRedisTemplate.hasKey("message:processed:" + eventId)) { return; }
// 方法2:使用数据库消息表
if (messageLogMapper.findByMsgId(eventId) != null) {
log.info("消息已处理,跳过: {}", eventId);
return;
}
// 核心业务逻辑处理
// 这里执行更新订单状态等操作
// 记录已处理的消息ID
messageLogMapper.insert(new MessageLog(eventId));
}
}
高级方案:基于消息中间件本身的机制
如果不想自己写发件箱逻辑,可以直接利用成熟中间件的高级特性:
RocketMQ 事务消息
- 发送半消息。
- 执行本地事务确认。
- 回查与最终发送。
- 实现:与发件箱方案类似,但已由中间件封装,不需要自己写轮询。
Kafka 幂等生产者
- 设置
enable.idempotence=true,Kafka会在内部处理重试和去重,保证分区内不重复且顺序不丢失。
总结与最佳实践
| 场景 | 推荐方案 |
|---|---|
| 小规模单体应用 | 发件箱(Outbox)定时任务方案 |
| 中大型微服务 | 事务性发件箱 + CDC(如Debezium) |
| 高并发海量数据 | Kafka 事务 + 幂等生产者 |
| 金融分布式事务 | 阿里的 Seata(AT模式) 或 TCC 模式 |
关于可靠发送,有几个核心结论:
- 绝对不丢是核心目标,通过 数据落盘 + 重试 保证。
- 绝对不重是伪命题,网络重传无法避免,必须通过业务层幂等解决。
- 推荐使用发件箱模式将数据一致性牢牢掌握在自己手中,这也是当前后端架构师岗位的核心考点。
希望这个案例能帮你全面理解消息可靠发送机制,如果有具体的中间件(如Kafka或RabbitMQ)细节需要深入,可以继续提问。