多来买:Lua 库存与 RocketMQ 事务消息实战
多来买:Lua 库存与 RocketMQ 事务消息实战
本篇范围:事务消息发送、Redis Lua 原子扣库存、本地事务监听器、Broker 回查与业务结果映射
事实来源:当前PromoServiceImpl、PromoTransactionProducer、RedisStockOper及 MQ 常量
验证口径:代码已确认;未连接 Redis 或 RocketMQ,未验证消息提交、回查和重试行为
本篇聚焦秒杀链路最核心的协调段:Producer 先发 Half Message,再把 Redis Lua 扣库存作为生产者侧“本地事务”,最后依据库存结果提交或回滚消息;当结果不确定时,Broker 根据 transactionId 对 Redis 标记进行回查。
事务消息主链路
阅读目标
- 说明 RocketMQ Half Message、提交、回滚和回查四个阶段。
- 说明这里的“本地事务”是 Redis Lua,而不是 MySQL 事务。
- 说明 Lua 为什么能原子完成库存检查与扣减。
- 解释 transactionId Bucket 怎样支持 Broker 回查。
- 区分“事务消息提交成功”与“订单消费者已经落库”。
第七步:发送 RocketMQ 事务消息
代码路径:duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/service/impl/PromoServiceImpl.java
关键方法:submitOrderInTransaction
Long skuId = orderInfoParam .getOrderDetailList() .get(0) .getSkuId();
SeckillGoodsDTO goods = getSeckillGoodsDTO(skuId);
if (goods == null) { throw new BusinessException( SeckillCodeEnum.SECKILL_ILLEGAL);}
HashMap<String, Object> localParams = new HashMap<>();localParams.put("id", goods.getId());localParams.put("stock", 1);localParams.put("skuId", skuId);
MqResultEnum result = promoTransactionProducer .sendTransactionMessage( MqTopicConst .PROMO_ORDER_TOPIC, orderInfoParam, localParams);消息 body 是 OrderInfoParam JSON,本地事务参数包含 SKU、活动 ID 和扣减数量。活动 ID 在当前 Lua 分支没有使用,属于从数据库扣库存方案保留下来的参数。
发送结果分支
if (MqResultEnum.SEND_FAIL .equals(result)) { RSet<Long> set = redissonClient.getSet( RedisConst .PROMO_USER_ORDERED_FLAG + orderInfoParam.getUserId());
set.remove(skuId);
throw new BusinessException( SeckillCodeEnum .SECKILL_ORDER_TRY_AGAIN);}
if (MqResultEnum.LOCAL_TRANSACTION_FAIL .equals(result)) { throw new BusinessException( SeckillCodeEnum.SECKILL_FINISH);}
if (MqResultEnum.LOCAL_TRANSACTION_EXCEPTION .equals(result)) { throw new BusinessException( SeckillCodeEnum .SECKILL_ORDER_TRY_AGAIN);}只有 SEND_FAIL 显式删除用户 Set 标记。库存不足和本地事务未知分支没有删除:
- 库存确实不足时保留标记影响不大;
- 结果未知但实际未扣库存时,用户重试会被当成重复抢购。
Controller 在方法正常返回后响应“正在排队”,而不是“下单成功”。
RocketMQ 事务消息的真实语义
四个阶段
1. Producer 发送 Half Message2. Broker 暂不向消费者投递3. Producer 执行本地事务:Redis Lua 扣库存4. Producer 返回: COMMIT -> 消息对订单消费者可见 ROLLBACK -> Broker 删除半消息 UNKNOWN -> Broker 稍后回查当前所谓“本地事务”不是 MySQL 事务,而是 Producer 回调里的一次 Redis Lua 操作,以及把结果写入 Redis Bucket。
Producer 初始化
代码路径:duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/mq/PromoTransactionProducer.java
关键方法:init
@PostConstructpublic void init() throws MQClientException {
transactionMQProducer = new TransactionMQProducer( "tx-promo-service-group");
// 注册中心地址省略,不进入文档 transactionMQProducer .setTransactionListener( new TransactionListener() { // executeLocalTransaction // checkLocalTransaction });
transactionMQProducer.start();}生产者会随 Spring Bean 初始化启动。源码中注册中心地址是硬编码;订单消费者使用另一个硬编码值,两端不一致,这是 P0 级配置风险,但文档不回显具体地址。
第八步:Lua 原子检查并扣库存
代码路径:duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/redis/RedisStockOper.java
当前 Lua
local current = redis.call('hget', KEYS[1], ARGV[1])
if not current or tonumber(current) < tonumber(ARGV[2]) then return -1end
return redis.call( 'hincrby', KEYS[1], ARGV[1], -tonumber(ARGV[2]))参数映射:
| Lua 参数 | 当前值 |
|---|---|
KEYS[1] | promo:goods:stock |
ARGV[1] | skuId |
ARGV[2] | 本次扣减数量,当前为 1 |
为什么不能拆成 Java 三步
错误的非原子思路:
stock = HGET(...)if stock >= 1: HINCRBY(..., -1)两个线程可以同时读到 1,都通过判断,然后各扣一次。Lua 在 Redis 服务端一次执行,脚本运行期间不会与另一条命令交错,因此“检查 + 扣减”成为一个原子操作。
加载和执行 SHA
@PostConstructpublic void loadScript() { sha1 = redissonClient .getScript() .scriptLoad(STOCK_LUA);}
public Long decrRedisStock( Long skuId, Integer count) { return redissonClient .getScript(StringCodec.INSTANCE) .evalSha( RScript.Mode.READ_WRITE, sha1, RScript.ReturnType.INTEGER, List.of( RedisConst .PROMO_SECKILL_GOODS_STOCK), skuId.toString(), count.toString());}脚本启动时加载,调用时通过 SHA 执行。当前没有看到 NOSCRIPT 后自动重新加载并重试的分支;Redis 脚本缓存被清理时,影响需要运行验证。
第九步:事务监听器记录库存结果
代码路径:duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/mq/PromoTransactionProducer.java
关键方法:executeLocalTransaction
扣库存
Long skuId = (Long) localParams.get("skuId");Integer stock = (Integer) localParams.get("stock");
String transactionId = msg.getTransactionId();
RBucket<String> stateBucket = redissonClient.getBucket( RedisConst .PROMO_TRANSACTION_PREFIX + transactionId);
Long remainStock = redisStockOper .decrRedisStock( skuId, stock);决定提交或回滚
if (remainStock.intValue() < 0) { stateBucket.set("fail");
LocalCacheHelper.put( skuId.toString(), LocalStockStatus .STOCK_NOT_ENOUGH .getNum());
return LocalTransactionState .ROLLBACK_MESSAGE;}
stateBucket.set("success");
return LocalTransactionState .COMMIT_MESSAGE;Lua 返回:
-1:库存不存在或不足,回滚消息;0:刚好扣完,提交消息;- 大于 0:仍有库存,提交消息。
本地售罄状态更新慢一拍
代码中“remainStock == 0 时标记售罄”的分支被注释。当前只有下一次请求执行 Lua 并得到 -1 时,才把本地状态改成无库存。
这不会让 Redis 库存变负,因为 Lua 仍会拒绝;但会让每个实例在库存刚好归零后多放过至少一次本地检查,产生额外 Redis/MQ 事务请求。
第十步:Broker 回查本地事务
@Overridepublic LocalTransactionStatecheckLocalTransaction(MessageExt msg) {
String transactionId = msg.getTransactionId();
RBucket<String> bucket = redissonClient.getBucket( RedisConst .PROMO_TRANSACTION_PREFIX + transactionId);
String value = bucket.get();
if ("success".equals(value)) { return LocalTransactionState .COMMIT_MESSAGE; }
if ("fail".equals(value)) { return LocalTransactionState .ROLLBACK_MESSAGE; }
return LocalTransactionState.UNKNOW;}当 Producer 执行完库存逻辑却没把结果成功告诉 Broker,Broker 会回查:
- Redis 为
success:提交; - Redis 为
fail:回滚; - key 不存在或值未知:继续返回
UNKNOW。
事务状态 Bucket 没有 TTL。它便于回查,却也会长期积累;活动清理的 promo:* 删除如果与 Broker 回查重叠,也可能让已知结果变成未知。
第十一步:把事务发送结果映射为业务结果
代码路径:duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/mq/PromoTransactionProducer.java
关键方法:sendTransactionMessage
TransactionSendResult sendResult = transactionMQProducer .sendMessageInTransaction( message, localTransactionParam);
SendStatus sendStatus = sendResult.getSendStatus();
LocalTransactionState localState = sendResult .getLocalTransactionState();
if (!SendStatus.SEND_OK .equals(sendStatus)) { return MqResultEnum.SEND_FAIL;}
if (LocalTransactionState .ROLLBACK_MESSAGE .equals(localState)) { return MqResultEnum .LOCAL_TRANSACTION_FAIL;}
if (LocalTransactionState .COMMIT_MESSAGE .equals(localState)) { return MqResultEnum .Local_TRANSACTION_SUCCESS;}
return MqResultEnum .LOCAL_TRANSACTION_EXCEPTION;Java 方法返回“本次发送和本地事务观察结果”,不代表订单消费者已经完成。即使是 Local_TRANSACTION_SUCCESS,订单消息仍在异步投递途中。
本篇限制与风险收束
- Lua 只保护一个 Redis 库存 Hash,不解决用户重复、订单落库和库存回补。
- Lua 扣库存与写 transactionId 结果是两次 Redis 操作,中间存在进程退出窗口。
- transactionId Bucket 没有 TTL,清理时机依赖活动清理入口。
- Producer 与 Consumer 的硬编码地址当前不一致,消息是否到达同一 Broker 只能标记为未验证。
- 方法返回事务发送结果,只说明消息侧决策,不代表下游订单已经创建。
关键源码导航
| 阅读顺序 | 文件 | 关键方法 |
|---|---|---|
| 1 | duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/service/impl/PromoServiceImpl.java | submitOrderInTransaction |
| 2 | duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/mq/PromoTransactionProducer.java | executeLocalTransaction、checkLocalTransaction |
| 3 | duolaimall-promo/promo-service/src/main/java/com/cskaoyan/mall/promo/redis/RedisStockOper.java | decrRedisStock |
| 4 | duolaimall-common/common-mq/src/main/java/com/cskaoyan/mall/mq/constant/MqTopicConst.java | PROMO_ORDER_TOPIC |
| 5 | duolaimall-common/common-service/src/main/java/com/cskaoyan/mall/common/constant/RedisConst.java | 事务回查状态 key |
4 条短复习点
- Lua 保证 Redis 内“查库存 + 扣库存”不可被其他命令插入。
- RocketMQ 事务消息根据生产者侧库存结果决定订单消息是否可见。
- Broker 回查依赖 Redis 中 transactionId 对应的业务结果。
- 消息 COMMIT 仍不等于订单落库,消费端还需要幂等和重试。
本章总结
Lua 保证 Redis 库存检查与扣减的单实例原子性,RocketMQ 事务消息据此决定订单消息是否可见;两者没有覆盖消费建单和补偿,因此不能表述为全链路强一致。
下一篇:秒杀订单消费、轮询与补偿边界实战。
文章分享
如果这篇文章对你有帮助,欢迎分享给更多人!