秒杀系统最致命的问题:库存扣了,订单没了怎么办?
SryYa Three

Redis + Lua扣库存就安全了….吗?

在电商秒杀场景中,我们可能会用到Redis + Lua脚本判断秒杀资格、Redis预扣减库存和Redis Stream异步生成订单的秒杀方案。

我们先把链路拆开看:

1
Redis 扣库存  →  发送 MQ  →  消费生成订单  →  DB 扣库存

任何一环出问题,都会导致不一致:

  • MQ 消息丢失
    • Redis 已扣库存
    • 但消息没进入队列
    • 结果:订单永远不会生成

  • MQ 消费失败
    • 消息到了,但消费挂了
    • 或者消费过程中服务崩溃
    • 结果:库存减少,但订单不存在

  • DB 扣库存失败
    • Redis 扣成功

    • 但数据库 update 失败

    • 结果:数据彻底混乱

我们第一反应可能是:下单失败了,那我把 Redis 库存加回去不就行了?

听起来合理,但实际根本难以实现:

怎么判断“失败”? 回滚哪部分?消费顺序会不会不可知?

在这个由 “Redis + 消息队列 + 数据库” 构建的分布式系统,“回滚”基本不可控。

那我们不如换个思路:既然不能保证每一步都成功,我们就接受不一致,但保证最终能修复一致

下面,我们就用利用Kafka + 对账日志,对秒杀系统进行重构,把每一次库存变化,变成“可追溯事件”


重构秒杀

流程图

image


代码实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
public Result<Long> doSeckillVoucherV2(Long voucherId) {
// 查询秒杀优惠券
SeckillVoucherFullModel seckillVoucherFullModel = seckillVoucherService.queryByVoucherId(voucherId);
// 加载优惠券库存
seckillVoucherService.loadVoucherStock(voucherId);
Long userId = UserHolder.getUser().getId();
// 雪花算法生成全局唯一ID
long orderId = snowflakeIdGenerator.nextId();
long traceId = snowflakeIdGenerator.nextId();
// 执行lua脚本需要的key(单槽位Hash Tag键,不分片)
List<String> keys = ListUtil.of(
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_STOCK_TAG_KEY, voucherId).getRelKey(),
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_USER_TAG_KEY, voucherId).getRelKey(),
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_TRACE_LOG_TAG_KEY, voucherId).getRelKey()
);
// 构造Lua脚本中需要的数据
String[] args = new String[9];
args[0] = voucherId.toString();
args[1] = userId.toString();
args[2] = String.valueOf(LocalDateTimeUtil.toEpochMilli(seckillVoucherFullModel.getBeginTime()));
args[3] = String.valueOf(LocalDateTimeUtil.toEpochMilli(seckillVoucherFullModel.getEndTime()));
args[4] = String.valueOf(seckillVoucherFullModel.getStatus());
args[5] = String.valueOf(orderId);
args[6] = String.valueOf(traceId);
args[7] = String.valueOf(LogType.DEDUCT.getCode());
long secondsUntilEnd = Duration.between(LocalDateTimeUtil.now(), seckillVoucherFullModel.getEndTime()).getSeconds();
long ttlSeconds = Math.max(1L, secondsUntilEnd + Duration.ofDays(1).getSeconds());
args[8] = String.valueOf(ttlSeconds);
// 执行脚本
SeckillVoucherDomain seckillVoucherDomain = seckillVoucherOperate.execute(
keys,
args
);
// 获取执行结果,若返回失败则抛出异常
if (!seckillVoucherDomain.getCode().equals(BaseCode.SUCCESS.getCode())) {
throw new FrameException(Objects.requireNonNull(BaseCode.getRc(seckillVoucherDomain.getCode())));
}
// 构造下单消息
SeckillVoucherMessage seckillVoucherMessage = new SeckillVoucherMessage(
userId,
voucherId,
orderId,
traceId,
seckillVoucherDomain.getBeforeQty(),
seckillVoucherDomain.getDeductQty(),
seckillVoucherDomain.getAfterQty(),
Boolean.FALSE
);
// 发送kafka
seckillVoucherProducer.sendPayload(
SpringUtil.getPrefixDistinctionName() + "-" + SECKILL_VOUCHER_TOPIC,
seckillVoucherMessage);

// 返回订单id
return Result.ok(orderId);
}

核心流程

前置准备

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
// 查询秒杀优惠券
SeckillVoucherFullModel seckillVoucherFullModel = seckillVoucherService.queryByVoucherId(voucherId);
// 加载优惠券库存到Redis
seckillVoucherService.loadVoucherStock(voucherId);
Long userId = UserHolder.getUser().getId();
// 雪花算法生成全局唯一ID
long orderId = snowflakeIdGenerator.nextId();
long traceId = snowflakeIdGenerator.nextId();
// 执行Lua脚本需要的key(单槽位Hash Tag键,不分片)
List<String> keys = ListUtil.of(
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_STOCK_TAG_KEY, voucherId).getRelKey(),
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_USER_TAG_KEY, voucherId).getRelKey(),
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_TRACE_LOG_TAG_KEY, voucherId).getRelKey()
);
// 构造Lua脚本中需要的参数
String[] args = new String[9];
args[0] = voucherId.toString();
args[1] = userId.toString();
args[2] = String.valueOf(LocalDateTimeUtil.toEpochMilli(seckillVoucherFullModel.getBeginTime()));
args[3] = String.valueOf(LocalDateTimeUtil.toEpochMilli(seckillVoucherFullModel.getEndTime()));
args[4] = String.valueOf(seckillVoucherFullModel.getStatus());
args[5] = String.valueOf(orderId);
args[6] = String.valueOf(traceId);
args[7] = String.valueOf(LogType.DEDUCT.getCode());
long secondsUntilEnd = Duration.between(LocalDateTimeUtil.now(), seckillVoucherFullModel.getEndTime()).getSeconds();
long ttlSeconds = Math.max(1L, secondsUntilEnd + Duration.ofDays(1).getSeconds());
args[8] = String.valueOf(ttlSeconds);

这段代码做了抢购优惠券前的准备:查询优惠券并将其加载到Redis,雪花算法生成分布式下的全局唯一ID,构造Lua脚本所需要的参数。

同时,我们可以看到Lua脚本的Key使用到了Hash Tag,下面简单解释一下。

Redis 能保证 Lua 脚本的原子性,但前提是脚本操作的所有 Key 都必须在同一个物理节点上。在集群模式下,Key可能会分配在不同物理节点上,而Hash
Tag便能解决这个问题。

Hash Tag: 确保相关的 Key 能够被分配到同一个哈希槽(Slot)中,从而在同一台物理节点上执行。

具体使用: 使用{}包裹需要进行Hash的Key即可。例如”seckill:stock:{%s}”,这里就是根据{%s}进行Hash分配物理节点。


Lua脚本

1
2
3
4
5
6
7
8
9
// 执行脚本
SeckillVoucherDomain seckillVoucherDomain = seckillVoucherOperate.execute(
keys,
args
);
// 获取执行结果,若返回失败则抛出异常
if (!seckillVoucherDomain.getCode().equals(BaseCode.SUCCESS.getCode())) {
throw new FrameException(Objects.requireNonNull(BaseCode.getRc(seckillVoucherDomain.getCode())));
}

下面给出Lua脚本具体内容:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
-- 1.参数列表
-- 单槽位库存key(已格式化,带HashTag)
local stockKey = KEYS[1]
-- 单槽位用户集合key(已格式化,带HashTag)
local seckillUserKey = KEYS[2]
-- 单槽位扣减日志key(已格式化,带HashTag)
local traceLogKey = KEYS[3]
-- 优惠券id
local voucherId = ARGV[1]
-- 用户id
local userId = ARGV[2]
-- 活动开始/结束时间(毫秒)
local beginTime = tonumber(ARGV[3])
local endTime = tonumber(ARGV[4])
-- 优惠券的状态
local status = tonumber(ARGV[5])
-- 订单id与日志TTL(秒)
local orderId = ARGV[6]
-- traceId
local traceId = ARGV[7]
-- 操作类型
local logType = ARGV[8]
local ttlSeconds = tonumber(ARGV[9])
-- 2.脚本业务
-- 当前时间(毫秒)
local timeArr = redis.call('TIME')
local nowMillis = tonumber(timeArr[1]) * 1000 + math.floor(tonumber(timeArr[2]) / 1000)

-- 时间范围判断:未开始或已结束直接返回
if nowMillis < beginTime then
return string.format('{"%s": %d}', 'code', 10002)
end
if nowMillis > endTime then
return string.format('{"%s": %d}', 'code', 10003)
end
-- 判断优惠券状态
if (status == 2) then
return string.format('{"%s": %d}', 'code', 10011)
end
if (status == 3) then
return string.format('{"%s": %d}', 'code', 10012)
end
local stock = redis.call('get', stockKey);
-- 缓存中的秒杀券库存为空,则直接返回
if not stock then
return string.format('{"%s": %d}', 'code', 10004)
end
-- 判断库存是否充足
if (tonumber(stock) <= 0) then
-- 库存不足,则直接返回
return string.format('{"%s": %d}', 'code', 10005)
end
-- 3.2.判断用户是否下单
if (redis.call('sismember', seckillUserKey, userId) == 1) then
-- 3.3.存在,说明是重复下单,返回2
return string.format('{"%s": %d}', 'code', 10006)
end
-- 3.4.扣库存 incrby stockKey -1
-- 记录扣减前数量与扣减量/扣减后数量
local beforeQty = tonumber(stock)
local changeQty = 1
local afterQty = beforeQty - changeQty
redis.call('incrby', stockKey, -changeQty)
-- 3.5.下单(保存用户)
redis.call('sadd', seckillUserKey, userId)
-- 3.6.记录扣减日志
local timeArr2 = redis.call('TIME')
-- 时间戳
local logNowMillis = tonumber(timeArr2[1]) * 1000 + math.floor(tonumber(timeArr2[2]) / 1000)
-- 扣减日志的信息
local logEntry = cjson.encode({
logType = logType,
ts = logNowMillis, -- 时间戳
traceId = traceId, -- 操作id,和数据库中的记录相关联
orderId = orderId, -- 订单id
userId = userId, -- 用户id
voucherId = voucherId, -- 优惠券id
beforeQty = beforeQty, -- 扣减前库存数量
changeQty = changeQty, -- 库存数量差值
afterQty = afterQty, -- 扣减后库存数量
})
-- 设置扣减信息,使用hash类型,key:链路id,方便和数据库中的记录关联。value:信息
redis.call('hset', traceLogKey, traceId, logEntry)
if ttlSeconds and ttlSeconds > 0 then
redis.call('expire', traceLogKey, ttlSeconds)
end
-- 返回信息,包括code码、扣减前数量,扣减数量,扣减后数量
return string.format('{"%s": %d, "%s": %s, "%s": %s, "%s": %s}', 'code', 0, 'beforeQty', beforeQty, 'deductQty', changeQty, 'afterQty', afterQty)

执行流程

  • 获取服务器当前时间

    • 使用 redis.call('TIME') 获取服务器秒与微秒并转换为毫秒
  • 时间窗口校验

    • 未开始:nowMillis < beginTime → 返回 code=10002
    • 已结束:nowMillis > endTime → 返回 code=10003
  • 状态校验

    • 下架:status == 2code=10011
    • 过期:status == 3code=10012
  • 读取与校验库存

    • 获取缓存库存值 get stockKey
    • 缓存不存在:返回 code=10004(库存缓存不存在)
    • 库存不足:tonumber(stock) <= 0code=10005(库存不足,扣减失败)
  • 一人一单校验

    • 已购集合:将已购用户id存入set集合,并用sismember seckillUserKey userId 判断是否存在,若存在则返回 code=10006
      (重复下单)
  • 原子扣减与记录

    • 计算扣减前/扣减量/扣减后:beforeQtychangeQty=1afterQty=beforeQty-1
    • 扣库存:incrby stockKey -1
    • 记录用户已购:sadd seckillUserKey userId
  • 扣减日志

    • 再次获取服务器时间作为日志时间戳
    • 构造 JSON 日志
    • 哈希写入:hset traceLogKey traceId logEntry,并设置过期 expire traceLogKey ttlSeconds(>0 时)
  • 返回结果

  • 成功返回 JSON 字符串:{"code":0,"beforeQty":...,"deductQty":1,"afterQty":...}

流程图:

image


异步下单

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// 构造下单消息
SeckillVoucherMessage seckillVoucherMessage = new SeckillVoucherMessage(
userId,
voucherId,
orderId,
traceId,
seckillVoucherDomain.getBeforeQty(),
seckillVoucherDomain.getDeductQty(),
seckillVoucherDomain.getAfterQty(),
Boolean.FALSE
);
// 发送kafka
seckillVoucherProducer.sendPayload(
SpringUtil.getPrefixDistinctionName() + "-" + SECKILL_VOUCHER_TOPIC,
seckillVoucherMessage);

// 返回订单id
return Result.ok(orderId);
}

在用Lua脚本完成一系列操作后,正式进入到下单环节,这里我们用到了Kafka去异步下单,先返回订单ID的方案。

需要注意的是,Kafka 只是“传输层”,不是可靠性保障!

用上了Kafka,不代表就稳了,它也会有三个不可靠:

  1. Producer 发送失败

  2. Broker 持久化异常

  3. Consumer 消费失败

我们要想系统真正的可靠,那就要让 Kafka 即使出问题,也能恢复。

上面的Lua脚本里有一个记录扣减日志的操作,将来我们就可以基于这个日志实现Kafka的消费可靠性。

可靠地生成订单

下面我们先看成功的生成订单是怎么样的。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
// RepeatExecuteLimit是封装好的幂等性组件,根据消息的UUID防止重复消费
@RepeatExecuteLimit(name = SECKILL_VOUCHER_ORDER,keys = {"#message.uuid"})
@Transactional(rollbackFor = Exception.class)
public boolean createVoucherOrder(MessageExtend<SeckillVoucherMessage> message) {
// 获取消息体
SeckillVoucherMessage messageBody = message.getMessageBody();
Long userId = messageBody.getUserId();
// 根据优惠券id和用户id查询是否已经存在正常的订单
VoucherOrder normalVoucherOrder = lambdaQuery()
.eq(VoucherOrder::getVoucherId, messageBody.getVoucherId())
.eq(VoucherOrder::getUserId, userId)
.eq(VoucherOrder::getStatus,OrderStatus.NORMAL.getCode())
.one();
// 如果存在,则直接结束运行
if (Objects.nonNull(normalVoucherOrder)) {
log.warn("已存在此订单,voucherId:{},userId:{}", normalVoucherOrder.getVoucherId(), userId);
throw new HmdpFrameException(BaseCode.VOUCHER_ORDER_EXIST);
}
// 扣减库存
boolean success = seckillVoucherService.update()
// set stock = stock - 1
.setSql("stock = stock - 1")
// where id = ? and stock > 0
.eq("voucher_id", messageBody.getVoucherId())
.gt("stock", 0)
.update();
if (!success) {
// 扣减失败:触发消息侧回滚Redis数据
throw new FrameException("优惠券库存不足!优惠券id:" + messageBody.getVoucherId());
}
// 构建订单
VoucherOrder voucherOrder = new VoucherOrder();
voucherOrder.setId(messageBody.getOrderId());
voucherOrder.setUserId(messageBody.getUserId());
voucherOrder.setVoucherId(messageBody.getVoucherId());
voucherOrder.setCreateTime(LocalDateTimeUtil.now());
save(voucherOrder);
// 创建订单路由
VoucherOrderRouter voucherOrderRouter = new VoucherOrderRouter();
voucherOrderRouter.setId(snowflakeIdGenerator.nextId());
voucherOrderRouter.setOrderId(voucherOrder.getId());
voucherOrderRouter.setUserId(userId);
voucherOrderRouter.setVoucherId(voucherOrder.getVoucherId());
voucherOrderRouter.setCreateTime(LocalDateTimeUtil.now());
voucherOrderRouter.setUpdateTime(LocalDateTimeUtil.now());
voucherOrderRouterService.save(voucherOrderRouter);
// 订单存放到redis
redisCache.set(RedisKeyBuild.createRedisKey(
RedisKeyManage.DB_SECKILL_ORDER_KEY,messageBody.getOrderId()),
voucherOrder,
60,
TimeUnit.SECONDS
);
// 对账日志:一致-消费成功
voucherReconcileLogService.saveReconcileLog(
LogType.DEDUCT.getCode(),
BusinessType.SUCCESS.getCode(),
"order created",
message
);
return true;
}
  • 提取消息体

    提取关键字段:userIdvoucherIdorderId库存变更前后数等,方便后续操作。

  • 一人一单幂等校验

    以“userId+voucherId+status”查询数据库是否已存在订单,存在则记录告警并抛出业务异常,防止重复生成订单。

  • 数据库扣减库存

    扣减语句:setSql("stock = stock - 1"),执行条件:voucher_id = ? AND stock > 0

    扣减失败(库存不足或并发竞争导致条件不满足)抛异常,交由上游消息侧回滚 Redis 数据。

  • 创建订单与路由

    构造 VoucherOrder:设置 iduserIdvoucherIdcreateTime,并保存到数据库。

    构造 VoucherOrderRouter:新建路由记录,包含 orderIduserIdvoucherId 及时间戳。

  • 写入订单短期缓存

    参数:DB_SECKILL_ORDER_KEYorderId 存储订单对象,设置过期时间TTL 60s。

    作用:支持前端轮询“订单是否生成”与页面快速反馈。

  • 记录对账日志

    类型:LogType.DEDUCT、业务结果:BusinessType.SUCCESS、详情 "order created"

    作用:与 Redis 扣减日志(Lua 写入的 traceLog)对齐,便于后续审计与补偿。


对账日志

我们在Kafka异步下单中的日志表结构如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
CREATE TABLE `tb_voucher_reconcile_log`
(
`id` bigint NOT NULL COMMENT '主键',
`order_id` bigint NOT NULL COMMENT '订单id',
`user_id` bigint unsigned NOT NULL COMMENT '下单的用户id',
`voucher_id` bigint unsigned NOT NULL COMMENT '购买的代金券id',
`message_id` varchar(64) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT 'Kafka消息uuid',
`detail` varchar(1024) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT '差异说明',
`before_qty` int DEFAULT NULL COMMENT '改变之前库存数量',
`change_qty` int DEFAULT NULL COMMENT '本次改变数量',
`after_qty` int DEFAULT NULL COMMENT '改变之后库存数量',
`trace_id` bigint DEFAULT NULL COMMENT '追踪唯一标识',
`log_type` int DEFAULT '-1' COMMENT '记录类型 -1:扣减 1:恢复',
`business_type` int DEFAULT '1' COMMENT '业务类型:1创建订单成功;2创建订单超时;3创建订单失败',
`reconciliation_status` int NOT NULL DEFAULT '1' COMMENT '对账状态:1待处理;2异常;3不一致;4一致',
`create_time` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`) USING BTREE,
KEY `idx_order_id` (`order_id`) USING BTREE,
KEY `idx_message_id` (`message_id`) USING BTREE,
KEY `idx_trace_id` (`trace_id`)
) ENGINE = InnoDB
DEFAULT CHARSET = utf8mb4
COLLATE = utf8mb4_general_ci
ROW_FORMAT = COMPACT;

这与我们上文Redis中保存的日志十分相似,为了方便对比,我们可以看看贴出的Redis中的日志:

1
2
3
4
5
6
7
8
9
10
11
{
"ts": 1763443683096,
"afterQty": 199,
"orderId": "1990653204235288577",
"userId": "1987041610793484289",
"logType": "-1",
"beforeQty": 200,
"traceId": "1990653204235288578",
"changeQty": 1,
"voucherId": "1"
}

数据库中的日志表有traceId,对应着Redis中日志的traceId,这两个日志便链接了起来。

  • 每次扣减库存,Redis 中有一条扣减记录日志,数据库中有一条扣减记录日志。

  • 每次恢复库存,Redis 中有一条恢复记录日志,数据库中有一条恢复记录日志。

流程图:

image

对账结果

有了Redis和MySQL的traceId之后,我们便能由此设计出一系列对账补偿机制。

我们可以利用定时任务,定期扫描数据库中reconciliation_status1(待处理)2(异常) 的记录。

让系统通过 trace_id 这一唯一标识,将 Redis 扣减日志数据库对账日志 进行对齐:

  • 有 Redis 日志,无 DB 订单:

    • 说明 Kafka 消息丢失或消费过程中程序崩溃。

    • 处理方案: 触发自动补单逻辑,重新尝试创建订单;若重试多次失败,则调用 rollbackRedisVoucherData 进行 Redis
      库存回滚,确保库存不消失。

  • Redis 与 DB 数据不符: 例如 before_qtyafter_qty 逻辑对不上。

    • 处理方案: 标记 reconciliation_status3(不一致),并推送到告警系统,转为人工介入。
  • 一致处理: 当确认数据库订单已生成、库存已扣减,且两端日志对齐时,将状态更新为 4(一致)

  • 回滚处理: 若因库存不足或业务限制导致下单失败,并已成功执行 Redis 回滚,则记录为 log_type1(恢复)
    的对账记录,并将该链路标记为已关闭。

通过这种 “日志留痕 + 定时溯源” 的机制,我们将原本脆弱的异步链路,转化为了具有自我修复能力的分布式事务补偿方案。即使
Kafka 挂了或者 Redis 抖动了,系统也能在分钟级内通过对账完成最终的一致性修复。


失败处理

在利用Kafka异步生成订单时有可能消费失败,对应的,我们就要对失败的情况有预案。

下面给出一种消费失败的处理方式:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
@Override
protected void afterConsumeFailure(final MessageExtend<SeckillVoucherMessage> message,
final Throwable throwable) {
super.afterConsumeFailure(message, throwable); // 打印日志
// 删除优惠卷秒杀已购集合中的用户的标记,默认为YES删除
SeckillVoucherOrderOperate seckillVoucherOrderOperate = SeckillVoucherOrderOperate.YES;
// 判断是否需要删除,若订单正常则不删
if (throwable instanceof FrameException frameException) {
if (Objects.nonNull(frameException.getCode()) &&
frameException.getCode().equals(BaseCode.VOUCHER_ORDER_EXIST.getCode())){
seckillVoucherOrderOperate = SeckillVoucherOrderOperate.NO;
}
}
// 生成日志的追踪ID
long traceId = snowflakeIdGenerator.nextId();
// 回滚Redis数据
redisVoucherData.rollbackRedisVoucherData(
seckillVoucherOrderOperate,
traceId,
// 获取优惠券Id
message.getMessageBody().getVoucherId(),
// 获取用户Id
message.getMessageBody().getUserId(),
// 获取订单Id
message.getMessageBody().getOrderId(),
// 获取回滚需要的库存数
message.getMessageBody().getAfterQty(),
message.getMessageBody().getChangeQty(),
message.getMessageBody().getBeforeQty()
);
// 对账日志:异常-消费失败
try {
String detail = throwable == null ? "consume failed" : ("consume failed: " + throwable.getMessage());
voucherReconcileLogService.saveReconcileLog (
LogType.RESTORE.getCode(),
BusinessType.FAIL.getCode(),
detail,
traceId,
message
);
} catch (Exception e) {
log.warn("保存对账日志失败(消费失败)", e);
}
}

回滚操作rollbackRedisVoucherData

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
public void rollbackRedisVoucherData(
SeckillVoucherOrderOperate seckillVoucherOrderOperate,
Long traceId,
Long voucherId,
Long userId,
Long orderId,
Integer beforeQty,
Integer changeQty,
Integer afterQty) {
// 构建 Lua 执行参数:KEY 列表包含库存、已购用户集合、trace 日志集合
List<String> keys = ListUtil.of(
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_STOCK_TAG_KEY, voucherId).getRelKey(),
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_USER_TAG_KEY, voucherId).getRelKey(),
RedisKeyBuild.createRedisKey(RedisKeyManage.SECKILL_TRACE_LOG_TAG_KEY, voucherId).getRelKey()
);
// Lua ARGV:voucherId、userId、orderId、操作码、traceId、日志类型(恢复)
String[] args = new String[9];
args[0] = String.valueOf(voucherId);
args[1] = String.valueOf(userId);
args[2] = String.valueOf(orderId);
args[3] = String.valueOf(seckillVoucherOrderOperate.getCode());
args[4] = String.valueOf(traceId);
args[5] = String.valueOf(LogType.RESTORE.getCode());
args[6] = String.valueOf(beforeQty);
args[7] = String.valueOf(changeQty);
args[8] = String.valueOf(afterQty);

// 带退避的重试,避免重试风暴,返回最终结果码(成功返回 SUCCESS 码)
Integer finalCode = luaRollbackWithResultCode(keys, args, retryMaxAttempts, initialBackoffMillis, maxBackoffMillis);
boolean ok = finalCode != null && finalCode.equals(BaseCode.SUCCESS.getCode());
if (!ok) {
String reason = BaseCode.getMsg(finalCode == null ? -1 : finalCode);
// 结构化错误日志:包含 voucher/user/order/trace 与失败原因
log.error("Redis回滚最终失败|voucherId={}|userId={}|orderId={}|traceId={} reason={}", voucherId, userId, orderId, traceId, reason);
// 失败日志:异常-恢复失败,记录 Lua 返回码供后续精准路由与统计
saveRollbackFailureLog(voucherId, userId, orderId, traceId, "redis rollback failed after retries: " + reason, finalCode);
// 指标:最终放弃(用尽重试次数)
safeInc("seckill_rollback_retry_give_up", "component", "redis_voucher_data");
}
}

调用Lua脚本回滚luaRollbackWithResultCode

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
private Integer luaRollbackWithResultCode(
List<String> keys,
String[] args,
int maxAttempts,
long initialBackoffMs,
long maxBackoffMs) {
// 当前重试次数(从 0 开始计数,日志展示为 attempt+1)
int attempt = 0;
// 初始退避时间,设置下限避免过小
long backoff = Math.max(50, initialBackoffMs);
Integer lastCode = null;
while (true) {
try {
Integer result = seckillVoucherRollBackOperate.execute(keys, args);
lastCode = result;
if (result != null && result.equals(BaseCode.SUCCESS.getCode())) {
// 指标:重试最终成功(包含首尝即成功和多次后成功)
safeInc("seckill_rollback_retry_success", "component", "redis_voucher_data");
return result;
}
// 失败,打印原因并继续重试
String reason = BaseCode.getMsg(result == null ? -1 : result);
log.warn("Redis回滚失败,准备重试|attempt={} reason={}", attempt + 1, reason);
} catch (Exception e) {
// 异常场景以 -1 兜底作为失败码
lastCode = -1;
log.warn("Redis回滚异常,准备重试|attempt={} error={}", attempt + 1, e.getMessage());
}
attempt++;
if (attempt >= maxAttempts) {
// 用尽重试次数,结束循环
break;
}
// 退避 + 抖动(jitter):避免热点一致重试导致雪崩
sleepQuietly(withJitter(backoff));
backoff = Math.min(backoff * 2, Math.max(backoff, maxBackoffMs));
}
// 返回最后一次的结果码,供失败日志与后续路由使用
return lastCode;
}

Lua脚本内容:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
-- 1.参数列表
-- 单槽位库存key(已格式化,带HashTag)
local stockKey = KEYS[1]
-- 单槽位用户集合key(已格式化,带HashTag)
local seckillUserKey = KEYS[2]
-- 单槽位操作日志key(已格式化,带HashTag)
local traceLogKey = KEYS[3]
-- 优惠券id
local voucherId = ARGV[1]
-- 用户id
local userId = (ARGV[2])
-- 订单id
local orderId = ARGV[3]
-- 操作码
local seckillVoucherOrderOperate = tonumber(ARGV[4])
local traceId = ARGV[5]
local logType = ARGV[6]
local beforeQty = tonumber(ARGV[7])
local changeQty = tonumber(ARGV[8])
local afterQty = tonumber(ARGV[9])
-- 2.脚本业务
local stock = redis.call('get', stockKey);
-- 缓存中的秒杀券库存为空,则直接返回
if not stock then
return 10004
end

-- 把库存删除
redis.call('incrby', stockKey, changeQty)
if seckillVoucherOrderOperate == 1 then
-- 删除下单记录(先判断存在再移除更稳妥)
if (redis.call('sismember', seckillUserKey, userId) == 1) then
redis.call('srem', seckillUserKey, userId)
end
end
-- 记录回滚日志
local timeArr = redis.call('TIME')
local nowMillis = tonumber(timeArr[1]) * 1000 + math.floor(tonumber(timeArr[2]) / 1000)
-- 构建回滚日志信息
local logEntry = cjson.encode({
logType = logType,
ts = nowMillis,
orderId = orderId,
traceId = traceId,
userId = userId,
voucherId = voucherId,
beforeQty = beforeQty,
changeQty = changeQty,
afterQty = afterQty
})
-- 向Redis存放回滚日志
redis.call('hset', traceLogKey, traceId, logEntry)
return 0


总结

放弃了强一致性,换来了“可恢复的一致性”,这才是分布式系统真正的解法。

到这里,这套秒杀系统的核心设计已经完整了。

我们回头看整个演进过程,我们其实做了一件很关键的事情:

不再执着于在分布式环境中“强行维持事务一致性”, 而是通过日志记录、链路追踪和对账补偿,让系统具备“自我修复能力”。

这是一种思维上的转变:从“保证每一步都正确”, 转向“允许过程出错,但结果必须正确”。这也是我们引入一系列中间件后不得不品的一环。

秒杀只是一个场景,但这种“日志驱动一致性”的设计思想,可以应用在更多需要高可靠性的分布式系统中。

由 Hexo 驱动 & 主题 Keep
总字数 28.1k 访客数 访问量