本地消息表方案将订单创建与消息记录置于同一数据库事务中,确保消息不丢失;后台定时任务扫描待发送消息,通过RocketMQ实现可靠投递,配合指数退避重试与死信队列机制,最终达成跨服务数据最终一致性。
在微服务架构实践中,跨服务数据同步始终是绕不开的难点。以Taocarts跨境电商独立站系统为例:用户一旦下单,订单数据需要迅速同步至商品服务扣减库存,同时通知物流服务生成运单、通知服务发送邮件……这一连串连锁反应中,任何环节出现问题,都可能影响用户体验。

长期稳定更新的攒劲资源: >>>点此立即查看<<<
如果采用强一致性分布式事务(如2PC两阶段提交或TCC补偿事务),性能和复杂度会急剧上升,得不偿失。更务实的做法是拥抱“最终一致性”:允许短暂的不一致窗口,但保证数据最终同步成功。这如同精密的异步协作机制,各司其职,最终达成全局一致。
如何落地呢?一个成熟可靠的方案是“本地消息表”,它本质上基于数据库的事务保障机制。核心思路是:在业务操作所在的数据库中增加一张专门记录消息的表,使业务操作和消息写入处于同一个本地事务中。这样,只要业务操作成功,消息必然被记录,从源头保证数据一致性。
消息表设计并不复杂,示例如下:
sql
CREATE TABLE local_message (
id bigint PRIMARY KEY AUTO_INCREMENT,
message_id varchar(64) NOT NULL COMMENT '消息唯一ID',
topic varchar(64) NOT NULL COMMENT 'MQ Topic',
payload text NOT NULL COMMENT '消息内容(JSON)',
status tinyint NOT NULL DEFAULT 0 COMMENT '0-待发送, 1-已发送, 2-发送失败',
retry_count int NOT NULL DEFAULT 0,
max_retries int NOT NULL DEFAULT 3,
next_retry_time datetime DEFAULT NULL,
created_at datetime DEFAULT CURRENT_TIMESTAMP,
updated_at datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_message_id (message_id),
KEY idx_status_next_retry (status, next_retry_time)
);
在订单创建代码层面,最关键的是将“保存订单”和“插入本地消息”两个动作放在一个事务里。只有订单和消息同时写入成功才算真正成功;任一失败则整个事务回滚。这是确保“不丢消息”的第一道防线。
ja va
@Service
@Transactional
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private LocalMessageMapper messageMapper;
public void createOrder(OrderDTO orderDTO) {
// 1. 保存订单
Order order = new Order();
order.setOrderNo(generateOrderNo());
order.setUserId(orderDTO.getUserId());
order.setAmount(orderDTO.getAmount());
order.setStatus(OrderStatus.PENDING_PAYMENT);
orderMapper.insert(order);
// 2. 保存本地消息
LocalMessage msg = new LocalMessage();
msg.setMessageId(UUID.randomUUID().toString());
msg.setTopic("ORDER_CREATED");
msg.setPayload(JSON.toJSONString(order));
msg.setStatus(0);
msg.setNextRetryTime(new Date());
messageMapper.insert(msg);
}
}
此处使用Spring的@Transactional注解,保证对orderMapper和messageMapper的两个插入操作在同一个数据库连接中,要么全部成功提交,要么全部失败回滚。消息唯一ID(message_id)至关重要,为后续幂等性处理埋下伏笔。
消息写入数据库只是第一步,关键要正确发送到RocketMQ。这需要定时任务完成:周期性扫描消息表,找出“待发送”记录,尝试发送到指定MQ Topic。
发送逻辑值得推敲:成功则更新状态为“已发送”;失败则采用重试策略。科学的做法是指数退避——首次重试间隔1分钟,第二次2分钟,第三次4分钟……直至达到最大重试次数。若仍失败,标记为“发送失败”状态,转入死信队列。
ja va
@Component
public class MessageSendScheduler {
@Autowired
private LocalMessageMapper messageMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Scheduled(fixedDelay = 5000)
public void sendPendingMessages() {
List messages = messageMapper.selectPendingMessages(100);
for (LocalMessage msg : messages) {
try {
SendResult result = rocketMQTemplate.syncSend(msg.getTopic(), msg.getPayload());
if (result.getSendStatus() == SendStatus.SEND_OK) {
msg.setStatus(1);
messageMapper.updateById(msg);
}
} catch (Exception e) {
msg.setRetryCount(msg.getRetryCount() + 1);
if (msg.getRetryCount() >= msg.getMaxRetries()) {
msg.setStatus(2); // 失败,进入死信
alertService.send("消息发送失败,进入死信队列: " + msg.getMessageId());
} else {
// 指数退避:2^retryCount 分钟
long delayMinutes = 1L << msg.getRetryCount();
msg.setNextRetryTime(new Date(System.currentTimeMillis() + delayMinutes * 60 * 1000));
}
messageMapper.updateById(msg);
}
}
}
}
该定时任务扮演“搬运工”角色,将消息从数据库搬运到MQ,并通过重试和错误处理机制,保证消息“不丢失、不重复、不堆积”。
消息到达下游服务后,消费端需处理“重复消息”。网络波动或MQ重试机制可能导致同一条消息被消费多次。若处理逻辑不具备幂等性,库存可能被重复扣减、重复生成运单,后果严重。
保证幂等性最常用的方法是利用消息唯一ID。消费端使用Redis记录已处理的消息ID:处理前先检查,若ID已存在则跳过;否则执行业务逻辑,处理完后写入Redis。
ja va
@Component
@RocketMQMessageListener(topic = "ORDER_CREATED", consumerGroup = "inventory-consumer")
public class InventoryConsumer implements RocketMQListener {
@Autowired
private RedisTemplate redisTemplate;
@Autowired
private InventoryService inventoryService;
@Override
public void onMessage(String message) {
JSONObject json = JSON.parseObject(message);
String messageId = json.getString("messageId");
// 幂等检查
String key = "processed:" + messageId;
Boolean success = redisTemplate.opsForValue().setIfAbsent(key, "1", Duration.ofHours(24));
if (Boolean.FALSE.equals(success)) {
log.info("消息已处理过,跳过: {}", messageId);
return;
}
try {
Order order = json.getObject("order", Order.class);
inventoryService.deductStock(order.getProductId(), order.getQuantity());
} catch (Exception e) {
// 处理失败,删除幂等标记,让消息重试
redisTemplate.delete(key);
throw e;
}
}
}
值得注意的是,setIfAbsent是原子操作,保证并发安全。若业务处理失败,需删除幂等标记,以便消息重试时再次处理。这一步虽简单,却是整个方案稳定运行的“守护神”。
即便方案完善,仍有极端情况导致消息始终无法发送成功,例如MQ服务器宕机、网络长时间中断。反复重试失败的消息最终进入死信队列。
对于死信消息,最稳妥的处理方式是告警+人工介入。例如对接钉钉群机器人,一旦有消息进入死信队列,自动发送告警通知。运维人员可通过管理后台查看死信列表,分析失败原因,手动重发或修复数据。
ja va
@RestController
@RequestMapping("/admin/messages")
public class DeadLetterController {
@Autowired
private LocalMessageMapper messageMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
@PostMapping("/retry/{id}")
public Result retry(@PathVariable Long id) {
LocalMessage msg = messageMapper.selectById(id);
if (msg.getStatus() != 2) {
return Result.error("只有死信消息可以重试");
}
rocketMQTemplate.syncSend(msg.getTopic(), msg.getPayload());
msg.setStatus(0);
msg.setRetryCount(0);
msg.setNextRetryTime(new Date());
messageMapper.updateById(msg);
return Result.success();
}
}
这种“自动化+人工兜底”的设计理念,是对系统可靠性的务实考量。它承认系统不可能100%完美,但通过机制设计将异常损失降到最低。
总的来说,Taocarts系统通过“本地消息表 + RocketMQ”组合方案,成功实现了订单创建与库存扣减、物流生成、通知发送等下游服务的解耦。在生产环境稳定运行一年后,消息送达率达到99.99%,仅有极少数死信需人工介入快速修复。这个数据充分验证了该方案的成熟与可靠。
侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述