本系列的前三篇文章,我们从数据库乐观锁、Redis 预扣 + 同步落库,一路演进到 Redis 预扣 + 异步落库。
每一次演进都解决了一个痛点,但也引入了新的问题,但这也就是单机防超卖或者说用于高并发场景下的比较合适的解决方案。
本篇文章就异步以后,消息一致性问题进行探讨。
单机环境(逻辑单机、部署单机)
第三篇的异步落库方案,响应时间大幅缩短,吞吐量显著提升,但代价是:系统从" 强一致性 " 退化为了 " 最终一致性 "。
原本 Redis预扣减库存成功后 进行 “ DB创建订单---DB扣减库存 ” 这两个步骤是一个事务,要么都成功,要么都失败,失败后只需要回滚Redis库存操作。
但在防止超卖,高并发场景下出现问题的根源部分 “ DB扣减库存 ” 这个步骤,我们将其改造成了异步,由MQ的消费者进行实际扣减,“DB创建订单”这个步骤还是保留在了 订单服务 里。
那问题显而易见,如果消费者消费失败,就会出现订单创建了,DB里库存并没有扣减,或者是MQ消息根本就是发送失败,也会出现订单和库存扣减动作只执行了其一。
本地消息表 + 补偿重试
在业务数据库中建一张本地消息表,将 " 业务操作 "和 " 消息记录 " 放在同一个数据库事务中。这样,要么业务操作和消息记录都成功,要么都失败——利用数据库 ACID 事务保证了强一致性。
然后,由后台定时任务扫描消息表中"未发送"或"发送失败"的消息,进行补偿重试。
本地消息表一些关键的字段:
字段名 | 解释 |
|---|---|
message_id | 业务唯一ID,用于幂等(防止消息重复消费) |
status | 消息生命周期状态 |
retry_count | 当前已重试次数 |
max_retry_count | 最大重试次数(超过则告警人工介入) |
next_retry_time | 结合指数退避策略,避免无效重试打满 CPU |
另外消费日志表
CREATE TABLE `local_message_consume_log` (
`id` BIGINT NOT NULL AUTO_INCREMENT,
`message_id` VARCHAR(64) NOT NULL COMMENT '消息唯一标识',
`biz_id` VARCHAR(64) NOT NULL COMMENT '业务ID(如订单号)',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_message_id` (`message_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='消息消费日志表(幂等表)';/**
*
* 生产者
* 这里只负责插入消息表,不直接发送 MQ。发送逻辑交给后台 Job 统一处理,保证"发送"这个动作也是可重试、可追溯的。
*/
@Service
@Slf4j
public class OrderService {
@Autowired
private RedisService redisService;
@Autowired
private OrderMapper orderMapper;
@Autowired
private LocalMessageMapper localMessageMapper;
public CreateOrderResult createOrder(Long productId, Integer buyCount, Long userId) {
// ========== 第1步:Redis 预扣减 ==========
boolean deducted = redisService.deductStock(productId, buyCount);
if (!deducted) {
throw new InsufficientStockException("库存不足");
}
// ========== 第2步:本地事务(创建订单 + 插消息表)==========
// ⚠️ 注意:这里不扣 DB 库存!DB 库存由 MQ 消费者异步扣减
String messageId;
Long orderId;
try {
orderId = saveOrderAndMessage(productId, buyCount, userId);
messageId = UUID.randomUUID().toString();
log.info("本地事务提交成功,orderId={}", orderId);
} catch (Exception e) {
// 本地事务失败 → 回滚 Redis
redisService.rollbackStock(productId, buyCount);
throw new RuntimeException("系统繁忙,请重试");
}
// ========== 第3步:返回用户成功 ==========
// ⚠️ 此时 DB 库存尚未扣减,订单状态为"待确认"
// 用户能看到订单,但支付前需要等待库存确认
// ========== 第4步:尝试发送 MQ(异步,不阻塞主流程)==========
try {
OrderMessage msg = buildOrderMessage(orderId, productId, buyCount, userId);
mqProducer.send("order-topic", msg);
localMessageMapper.updateStatus(messageId, MessageStatus.SENT.getCode());
log.info("MQ 发送成功,messageId={}", messageId);
} catch (Exception e) {
// 发送失败 → 消息表保持"待发送",后台 Job 补偿
log.warn("MQ 发送失败,等待 Job 补偿,messageId={}", messageId, e);
}
return CreateOrderResult.success(orderId);
}
@Transactional(rollbackFor = Exception.class)
public Long saveOrderAndMessage(Long productId, Integer buyCount, Long userId) {
// 1. 创建订单(状态:待确认,库存尚未扣减)
Order order = new Order();
order.setProductId(productId);
order.setBuyCount(buyCount);
order.setUserId(userId);
order.setStatus(OrderStatus.PENDING_CONFIRM.getCode()); // 待确认
order.setRedisDeducted(true); // 标记 Redis 已扣
orderMapper.insert(order);
// 2. 插入本地消息表(状态:待发送)
LocalMessage message = new LocalMessage();
message.setMessageId(UUID.randomUUID().toString());
message.setBizId(String.valueOf(order.getId()));
message.setTopic("order-topic");
message.setBody(JSON.toJSONString(buildOrderMessage(order.getId(), productId, buyCount, userId)));
message.setStatus(MessageStatus.PENDING.getCode());
message.setRetryCount(0);
message.setMaxRetryCount(5);
message.setNextRetryTime(new Date());
localMessageMapper.insert(message);
return order.getId();
}
}
/**
*
* 后台补偿job
*/
@Component
@Slf4j
public class LocalMessageJob {
@Autowired
private LocalMessageMapper localMessageMapper;
@Autowired
private MqProducer mqProducer;
/**
* 定时扫描待发送消息(每 5 秒执行一次)
*/
@Scheduled(fixedDelay = 5000)
public void processPendingMessages() {
// 查询状态为"待发送"或"发送失败",且到达重试时间的消息
// 每次取 100 条,避免单次处理太多
List<LocalMessage> messages = localMessageMapper.selectPendingMessages(
MessageStatus.PENDING.getCode(),
MessageStatus.FAILED.getCode(),
new Date(),
100
);
if (messages.isEmpty()) {
return;
}
log.info("扫描到 {} 条待发送消息", messages.size());
for (LocalMessage message : messages) {
processSingleMessage(message);
}
}
/**
* 处理单条消息
*/
private void processSingleMessage(LocalMessage message) {
try {
// 1. 发送 MQ
mqProducer.send(
message.getTopic(),
message.getTag(),
message.getBody()
);
// 2. 发送成功 → 更新状态为"已发送"
localMessageMapper.updateStatus(
message.getId(),
MessageStatus.SENT.getCode(),
message.getRetryCount() + 1
);
log.info("消息发送成功,messageId={}", message.getMessageId());
} catch (Exception e) {
log.error("消息发送失败,messageId={}, 当前重试次数={}",
message.getMessageId(), message.getRetryCount(), e);
// 3. 发送失败 → 更新重试次数 + 计算下次重试时间(指数退避)
int newRetryCount = message.getRetryCount() + 1;
int maxRetryCount = message.getMaxRetryCount();
if (newRetryCount >= maxRetryCount) {
// 超过最大重试次数 → 状态改为"最终失败",触发告警
localMessageMapper.updateStatus(
message.getId(),
MessageStatus.FINAL_FAILED.getCode(),
newRetryCount
);
// 发送告警通知(钉钉/邮件/电话)
sendAlert(message);
log.error("消息重试次数超过上限,进入最终失败状态,messageId={}", message.getMessageId());
} else {
// 计算下次重试时间:指数退避 2^n 秒
long delaySeconds = (long) Math.pow(2, newRetryCount);
Date nextRetryTime = new Date(System.currentTimeMillis() + delaySeconds * 1000);
localMessageMapper.updateStatusAndRetryTime(
message.getId(),
MessageStatus.FAILED.getCode(),
newRetryCount,
nextRetryTime
);
log.warn("消息发送失败,将在 {} 秒后重试,messageId={}", delaySeconds, message.getMessageId());
}
}
}
}
/**
*
* 消费者端
* 幂等处理
* 在消息队列场景中,重复消费是常态(网络抖动、消费者重启、ACK 超时等都可能触发重发)
* 因此,消费者必须实现幂等。
*/
@Component
@RocketMQMessageListener(topic = "order-topic", consumerGroup = "order-consumer-group")
@Slf4j
public class OrderConsumer implements RocketMQListener<OrderMessage> {
@Autowired
private ProductMapper productMapper;
@Autowired
private OrderMapper orderMapper;
@Autowired
private LocalMessageConsumeLogMapper consumeLogMapper;
@Override
public void onMessage(OrderMessage msg) {
String messageId = msg.getMessageId();
Long orderId = msg.getOrderId();
// 1. 幂等判断
if (consumeLogMapper.existsByMessageId(messageId)) {
log.info("消息已消费过,幂等跳过,messageId={}", messageId);
return;
}
try {
// 2. DB 扣减库存(乐观锁兜底)—— 这才是真正的库存扣减!
int rows = productMapper.decreaseStock(msg.getProductId(), msg.getBuyCount());
if (rows == 0) {
// DB 库存不足 → Redis 和 DB 数据不一致
// 重试无法解决,直接告警 + 人工介入
log.error("DB 库存不足!productId={}, orderId={}", msg.getProductId(), orderId);
// 不抛异常(抛了也解决不了),让消息进入死信队列
return;
}
// 3. 更新订单状态(待确认 → 已确认/待发货)
orderMapper.updateStatus(orderId, OrderStatus.CONFIRMED.getCode());
// 4. 记录消费日志(幂等)
consumeLogMapper.insert(messageId, orderId);
log.info("DB 库存扣减成功,orderId={}", orderId);
} catch (Exception e) {
// 可恢复异常(DB 连接超时等)→ 抛出异常,触发 MQ 重试
log.warn("消费者处理失败,将自动重试,orderId={}", orderId, e);
throw e;
}
}
}RocketMQ本身具有消息重试机制,如果消费者消费成功,进行ACK后,MQ才结束本消息的投递,在队列里将此消息删除,如果没有进行ACK或者自动重试超过次数,则可进入死信队列,此时就需要告警,进行人工介入。
当然,也可以另外一个job,每晚对账,也可以及时发现一些业务数据不对的问题。
至此,本系列四篇文章,从数据库乐观锁、Redis 预扣 + 同步落库、Redis 预扣 + 异步落库,再到本篇的本地消息表 + 补偿重试,我们完成了一条完整的演进路线:
阶段 | 方案 | 核心能力 | 代价 |
|---|---|---|---|
第一篇 | 数据库乐观锁 | 强一致,实现简单 | DB 行锁竞争,高并发下连接池易耗尽 |
第二篇 | Redis 预扣 + 同步落库 | 将查询压力前置到缓存 | 同步落库仍是瓶颈 |
第三篇 | Redis 预扣 + 异步落库(MQ) | 吞吐量大幅提升 | 最终一致性,消息可能丢失 |
第四篇 | 本地消息表 + 补偿重试 | 消息可靠性兜底,端到端闭环 | 侵入业务,需维护消息表 |
这套组合方案的核心思路可以概括为三句话:
单机环境下,这套方案就是解决大并发防超卖的最终方案。
当然,这套方案的价值远不止于防超卖。任何性能瓶颈在数据库层面的场景——比如订单同步、日志归档、数据聚合——都可以用同样的思路来解决:前置到缓存、异步解耦、消息表兜底。
技术方案的演进,本质是在性能、一致性、复杂度之间寻找平衡。 没有银弹,只有适合当下业务阶段的权衡。
但在实际生产环境中,这套组合拳会遇到各种各样的问题——有些是代码层面的坑,有些是中间件本身的限制,有些则是业务场景带来的挑战。
比如:
这些问题,在实际生产环境中几乎一定会遇到。如果处理不当,之前搭建的一切都可能功亏一篑。
所以,下一个系列,我将逐一列举这些实际生产环境中的问题、坑和应对策略,结合我在《八公酒业商城》的真实案例,一步一步拆解。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。