Skip to content

第 14 章 综合实战项目:秒杀系统

学习目标:把前 13 章学到的 String / Hash / Set / ZSet / Lua / Stream / 分布式锁 / 限流 / 缓存设计 等知识串起来,从零设计一套经得起 10 万 QPS、库存绝对不超卖的秒杀系统。看完之后能讲清楚「双 11 茅台秒杀架构怎么搭」「为什么扣库存一定要用 Lua」「限流和兜底要做哪些层」。


14.0 写在最前面

前 13 章我们一直在「拆开」讲 Redis 的某一个能力:第 2 章学了 5 种数据类型、第 6 章学了 Lua、第 7 章学了 Stream、第 11 章学了分布式锁、第 12 章学了缓存三大问题……

这一章我们要做的是「装回去」——把这些零件拼成一个真正能上线的系统:秒杀

为什么选秒杀?因为它几乎用上了 Redis 的每一个特性

┌─────────────────────────────────────────────────────────────┐
│                  秒杀 = Redis 全家桶演练                      │
├─────────────────────────────────────────────────────────────┤
│  String  →  库存计数      (第 2 章)                         │
│  Hash    →  商品多维信息   (第 2 章)                         │
│  Set     →  用户去重      (第 2 章)                         │
│  ZSet    →  滑动窗口限流   (第 2 章 + 第 7 章)               │
│  Lua     →  原子扣减      (第 6 章)                         │
│  Stream  →  异步下单消息  (第 7 章)                         │
│  分布式锁 → 兜底防超卖     (第 11 章)                        │
│  缓存设计 → 静态化+预热    (第 12 章)                        │
│  性能优化 → Pipeline 批量 (第 13 章)                        │
└─────────────────────────────────────────────────────────────┘

🍱 生活类比:秒杀就像「春运抢票」。售票窗口(数据库)一秒只能服务几个人,但站前广场涌进来 10 万人。我们要做的不是把窗口扩到 10 万个(成本爆炸),而是在广场上搭一层层的栅栏:实名验证、扫码进站、安检、排队、闸机……每一层都把人流过滤一遍,最后真正到达窗口的只剩几百个。Redis 在这里就是「站台前最后那道闸机」——又快又准,绝对不放超员。


14.1 业务场景与挑战

14.1.1 典型场景

场景 A:双 11 茅台秒杀
  - 商品:53 度飞天茅台 × 1000 瓶
  - 单价:1499 元(市场价 3000+)
  - 开抢时间:10:00:00 整
  - 预约用户:500 万
  - 预期 QPS:峰值 50 万 / 秒(前 1 秒爆发)

场景 B:演唱会抢票
  - 商品:周杰伦演唱会 × 8 万张
  - 单价:380 ~ 1880 元(多档)
  - 开抢时间:周五 20:00
  - 预约用户:300 万
  - 痛点:还要选座位(更复杂的库存模型)

场景 C:iPhone 新品首发
  - 商品:iPhone 16 Pro × 多种 SKU
  - 复杂度:颜色/容量/运营商组合 × 库存
  - 物流:还要选地区是否有货

14.1.2 核心挑战

┌──────────────────────┬─────────────────────────────────────┐
│  挑战                  │  说明                                │
├──────────────────────┼─────────────────────────────────────┤
│ ① 瞬间高并发           │ 10 万 QPS+,普通 MySQL 单机扛不住    │
│ ② 库存绝对不能超卖      │ 卖出去 1001 瓶 = 公关事故            │
│ ③ 防机器人 / 一人多单   │ 黄牛脚本一秒发 1000 个请求           │
│ ④ 公平性:先到先得      │ 不能让付费用户排不到队               │
│ ⑤ 失败要优雅降级        │ 库存没了要立刻告诉用户,不能让 TA 等  │
│ ⑥ 幂等性               │ 网络抖动重试不能扣两次库存            │
│ ⑦ 数据一致性            │ Redis 扣了 100,MySQL 也得是 100    │
└──────────────────────┴─────────────────────────────────────┘

💡 这 7 个挑战里,②超卖是底线——卖少了顶多挨骂,卖多了可能要赔款。所以本章会花最多笔墨讲「怎么保证不超卖」。

14.1.3 错误示范:直接打 MySQL 会怎样?

            10 万 QPS


         ┌──────────────┐
         │  MySQL 主库   │  ← 实际能扛 ~5000 QPS
         └──────────────┘


         ┌──────────────────┐
         │ 1. 连接池打满       │
         │ 2. 行锁等待 → 慢查询 │
         │ 3. 主从延迟暴涨     │
         │ 4. 雪崩 → 整站挂掉   │
         └──────────────────┘

  💀 真实事故:某电商 2017 年双 11,库存表用悲观锁,
              30 秒内 MySQL CPU 100%,订单系统瘫 8 分钟。

14.2 系统总架构

14.2.1 完整架构图

                            🌍  500 万 用户客户端


┌──────────────────────────────────────────────────────────────────┐
│  ① 边缘层 CDN                                                     │
│   - 商品详情页静态化(HTML/CSS/JS/图片)                            │
│   - 99% 流量被挡在 CDN,根本进不来                                  │
└──────────────────────────────────────────────────────────────────┘
                                       │  ~1% 动态流量

┌──────────────────────────────────────────────────────────────────┐
│  ② 接入层 Nginx / LVS                                             │
│   - 4 层负载均衡 + 7 层路由                                        │
│   - IP 黑名单(黄牛 IP 直接 403)                                  │
│   - 简单限流:单 IP 每秒最多 20 个请求                              │
└──────────────────────────────────────────────────────────────────┘


┌──────────────────────────────────────────────────────────────────┐
│  ③ 网关层 API Gateway                                             │
│   - 用户登录态校验(无 token 直接拒)                                │
│   - 全局 QPS 限流(令牌桶)                                         │
│   - 验证码 / 答题 / 排队页:把瞬时洪峰摊到 30 秒                       │
└──────────────────────────────────────────────────────────────────┘
                                       │  剩 ~50 万 QPS

┌──────────────────────────────────────────────────────────────────┐
│  ④ 应用层 Seckill Service(多机水平扩展)                          │
│   - 用户维度限流(ZSet 滑动窗口)                                    │
│   - 调用 Redis 做核心扣减                                          │
└──────────────────────────────────────────────────────────────────┘


┌──────────────────────────────────────────────────────────────────┐
│  ⑤ Redis 集群 ⭐ 核心战场                                          │
│   ┌──────────────┬──────────────┬──────────────┐                │
│   │ String 库存   │ Set 用户去重   │ Stream 消息  │                │
│   │ stock:1001=1000│ users:1001  │ orders:queue │                │
│   └──────────────┴──────────────┴──────────────┘                │
│   ⚡ 用 Lua 脚本一次完成「判库存 + 防重复 + 扣减 + 推消息」          │
└──────────────────────────────────────────────────────────────────┘
                                       │  Stream 消息

┌──────────────────────────────────────────────────────────────────┐
│  ⑥ 异步下单消费者(多实例)                                         │
│   - 从 Stream 取消息                                              │
│   - 写订单表(MySQL)                                              │
│   - 发短信 / 推 App                                                │
│   - ACK 消息(失败自动重试)                                         │
└──────────────────────────────────────────────────────────────────┘
                                       │  ~1000 TPS(库存数)

┌──────────────────────────────────────────────────────────────────┐
│  ⑦ MySQL 主库(订单 / 真实库存账本)                                │
│   - 只承担「真实订单写入」                                           │
│   - QPS 已经被 Redis 削掉 99.99%                                   │
└──────────────────────────────────────────────────────────────────┘

14.2.2 核心思想:把数据库前面所有能挡的都挡住

            漏斗式削峰
                                   
   500w 用户  ████████████████████████████████  100%
                       ↓ CDN
   动态请求  ██████                              5%
                       ↓ Nginx 限流
   合法请求  ████                                2%
                       ↓ 网关 + 验证码
   通过校验  ██                                  1%
                       ↓ 应用层用户限流
   到 Redis  █                                   0.5%
                       ↓ Redis Lua 库存扣减(99% 失败:库存售罄)
   抢到名额  ▏                                   0.001%
                       ↓ Stream → 异步落库
   写 MySQL  ▏                                   0.001%

本章 60% 的篇幅讲第 ⑤ 层(Redis 核心),但你要明白:能上线的秒杀系统前 4 层一层都不能少


14.3 限流层(前置防护)

14.3.1 IP / 用户维度限流:用 ZSet 实现滑动窗口

为什么不用「每秒计数器」(Fixed Window)?

固定窗口(不推荐):
  时间轴   00────01────02────03 (秒)
            │     │     │
  请求    ██████│██████ │
   1 秒到 1 秒 99 ms 来 100 个,
   1 秒 100 ms ~ 1 秒 200 ms 又来 100 个 ←—— 跨窗了,没限住!

滑动窗口(推荐):
  始终看「当前往前 1 秒内」的请求数,无窗口边界问题。

用 ZSet 实现滑动窗口(核心思路)

Key: rate:user:1001
ZSet 结构:
  ┌──────────┬──────────────────┐
  │  member  │  score           │
  ├──────────┼──────────────────┤
  │  uuid_1  │  1234567890.001  │  (毫秒时间戳)
  │  uuid_2  │  1234567890.157  │
  │  uuid_3  │  1234567890.342  │
  └──────────┴──────────────────┘

每次请求:
  1) ZREMRANGEBYSCORE key 0 (now - 1000)   清掉 1 秒前的
  2) ZCARD key                              数现在窗口里有几个
  3) 如果 < 限制:ZADD key now uuid + EXPIRE
  4) 否则:拒绝

Lua 实现(保证原子):

lua
-- KEYS[1] = rate:user:xxx
-- ARGV[1] = now (ms)
-- ARGV[2] = window (ms, e.g. 1000)
-- ARGV[3] = limit (e.g. 10)
-- ARGV[4] = uuid

local now = tonumber(ARGV[1])
local win = tonumber(ARGV[2])
local lim = tonumber(ARGV[3])

redis.call('ZREMRANGEBYSCORE', KEYS[1], 0, now - win)
local cnt = redis.call('ZCARD', KEYS[1])
if cnt >= lim then
  return 0   -- 被限流
end
redis.call('ZADD', KEYS[1], now, ARGV[4])
redis.call('PEXPIRE', KEYS[1], win)
return 1     -- 放行

14.3.2 验证码 / 答题 / 排队页:拉长用户操作时间

🎯 核心目的:把「1 秒内 50 万请求」摊成「30 秒内 50 万请求」。

方案 A:图形验证码
  - 用户输验证码平均 3 秒
  - 50 万 QPS → 1.6 万 QPS

方案B:选择题(如「茅台是哪一年发明的?」)
  - 答题平均 5 秒
  - 同时识别脚本(脚本不会答非标准答案)

方案C:排队页(前端 polling)
  - 用户进入「正在排队,您前面还有 X 人」
  - 服务端每秒放行 1 万人
  - 用户视觉上没感知,体验比直接 503 好得多

实现思路:

python
# 服务端
queue_key = "seckill:queue:1001"
r.rpush(queue_key, user_id)             # 入队
position = r.lpos(queue_key, user_id)   # 查位置

# 放行规则:用 Stream 或定时器,每秒 LPOP 1 万个
allowed_set = "seckill:allowed:1001"
def release_batch(n=10000):
    pipe = r.pipeline()
    for _ in range(n):
        pipe.lpop(queue_key)
    users = pipe.execute()
    if users:
        r.sadd(allowed_set, *[u for u in users if u])

14.3.3 静态资源 CDN:商品页 99% 流量挡在边缘

传统页面:商品详情页每次请求都查数据库
  → 1 万人看页面 = 1 万次 DB 查询

秒杀活动页:
  ┌────────────────────────────────────────────┐
  │ HTML / CSS / JS / 图片  →  全部 CDN 静态化   │
  │ 商品名称 / 价格 / 海报   →  活动开始前编译进 HTML │
  │ 库存数字              →  浏览器轮询 API        │
  │ 抢购按钮 onClick      →  调动态接口            │
  └────────────────────────────────────────────┘
  
  结果:CDN 命中率 99%+,源站只承担那 1% 的动态请求

14.4 库存预热(活动前)

14.4.1 为什么要预热

❌ 不预热:第一个请求来时去 MySQL 查库存
   → 此时 Redis 没缓存,请求穿透到 MySQL
   → 第 2~第 100 个请求都在等第 1 个,都打到 MySQL
   → MySQL 瞬间 100% CPU(缓存击穿)

✅ 预热:活动开始前 1 小时,把所有秒杀商品的库存灌入 Redis
   → 活动开始时 Redis 已就绪,0 穿透

14.4.2 用 String 存库存

bash
# 最简单的库存:Redis String,存一个数字
SET seckill:stock:1001 1000
EXPIRE seckill:stock:1001 7200    # 活动结束 + 1 小时自动清理

# 查库存
GET seckill:stock:1001            # → "1000"
DECR seckill:stock:1001           # 原子扣 1,返回剩余

14.4.3 用 Hash 存商品多维信息

bash
HSET seckill:item:1001 \
  name      "53度飞天茅台" \
  price     1499 \
  stock     1000 \
  sold      0 \
  status    "READY" \
  start_at  1729123200 \
  end_at    1729126800

HGET seckill:item:1001 stock
HINCRBY seckill:item:1001 stock -1
HINCRBY seckill:item:1001 sold +1

🤔 String 还是 Hash?

  • 库存数字单独频繁扣减 → String(INCR/DECR 比 HINCRBY 略快,且 Lua 操作简单)
  • 商品元信息(名称、状态、价格)→ Hash
  • 实战中两个都用:Hash 存元信息(只读),String 存库存(高频写)

14.4.4 预热代码(伪)

python
def preheat(item_ids):
    items = mysql.query("SELECT * FROM seckill_item WHERE id IN %s", item_ids)
    pipe = r.pipeline()
    for it in items:
        pipe.hset(f"seckill:item:{it.id}", mapping={
            "name": it.name, "price": it.price,
            "stock": it.stock, "sold": 0,
            "status": "READY",
        })
        pipe.set(f"seckill:stock:{it.id}", it.stock)
        pipe.expire(f"seckill:item:{it.id}", 7200)
        pipe.expire(f"seckill:stock:{it.id}", 7200)
    pipe.execute()

14.5 库存扣减的多种方案演进

这是本章的核心。我们用「不断翻车再修复」的方式,演化到最终方案。

14.5.1 V1:GET → 判断 → DECR ❌

python
def seckill_v1(item_id, user_id):
    stock = int(r.get(f"seckill:stock:{item_id}"))    # 1
    if stock <= 0:                                     # 2
        return "卖完了"
    r.decr(f"seckill:stock:{item_id}")                 # 3
    return "抢到了!"

会超卖。问题在第 1~3 步不是原子的

线程 A:  GET → 1
线程 B:  GET → 1                  ← 同一时刻读到 1
线程 A:  DECR → 0
线程 B:  DECR → -1               ← 超卖!

💥 高并发下,多线程跑这段会出现库存为负数。本章配套代码 02_seckill_v1_unsafe.py 实测会出现 -3、-7 的负库存。

14.5.2 V2:WATCH + MULTI(乐观锁)⚠️

python
def seckill_v2(item_id):
    key = f"seckill:stock:{item_id}"
    while True:
        try:
            with r.pipeline() as pipe:
                pipe.watch(key)
                stock = int(pipe.get(key))
                if stock <= 0:
                    pipe.unwatch()
                    return "卖完了"
                pipe.multi()
                pipe.decr(key)
                pipe.execute()         # 如果中间被改过,会抛 WatchError
                return "抢到了!"
        except redis.WatchError:
            continue

问题:高并发下 WATCH 几乎都失败 → CPU 空转 → 雪崩。

1000 个并发:
  - 仅 1~2 个能成功 commit
  - 其余 998 个抛 WatchError 重试
  - 重试后又 WATCH 又失败 ……
  
  结果:CPU 100%,Redis QPS 暴涨但成功率 < 1%

14.5.3 V3:Lua 脚本 ✅ 推荐

Redis 单线程执行 Lua 脚本时,整个脚本是原子的——不会有别的命令插进来。

lua
-- KEYS[1] = stock 的 key
-- 返回:1=成功,0=没货
local s = tonumber(redis.call('GET', KEYS[1]))
if not s or s <= 0 then
  return 0
end
redis.call('DECR', KEYS[1])
return 1

Python 调用:

python
LUA = """
local s = tonumber(redis.call('GET', KEYS[1]))
if not s or s <= 0 then return 0 end
redis.call('DECR', KEYS[1])
return 1
"""
seckill_lua = r.register_script(LUA)

def seckill_v3(item_id):
    ok = seckill_lua(keys=[f"seckill:stock:{item_id}"])
    return "抢到了" if ok else "卖完了"

为什么这个方案安全?

Redis 单线程模型:
  ┌─────────────────────────────────────────┐
  │       事件循环 (单线程)                    │
  │   ┌─────────────────────────────────┐   │
  │   │ EVAL "Lua脚本" KEYS ARGV        │ ← 整个脚本期间,
  │   │   - GET stock                    │   其他命令排队等待
  │   │   - 判断                          │   不可能插入!
  │   │   - DECR stock                   │
  │   └─────────────────────────────────┘   │
  └─────────────────────────────────────────┘

14.5.4 V4:DECR 后判断负数 → 回滚

python
def seckill_v4(item_id):
    key = f"seckill:stock:{item_id}"
    new_stock = r.decr(key)
    if new_stock < 0:
        r.incr(key)              # 回滚
        return "卖完了"
    return "抢到了"

优点:最简单,不用 Lua。 缺点:

  1. 中间会出现「库存负数」状态(虽然马上回滚,但如果此时 GET 会读到负数)
  2. 如果回滚 INCR 因网络失败而丢失 → 库存永远短少
  3. 不能附带其他条件(如「检查用户是否已买过」),扩展性差

14.5.5 方案对比

┌─────┬─────────────────────┬───────┬────────┬──────────────┐
│ 版本 │ 方案                │ 安全 │ 性能   │ 适用场景      │
├─────┼─────────────────────┼───────┼────────┼──────────────┤
│ V1  │ GET-判断-DECR        │  ❌  │  ⚡⚡  │ 永远不要用     │
│ V2  │ WATCH + MULTI        │  ✅  │  💀    │ 低并发能用     │
│ V3  │ Lua 脚本             │  ✅  │  ⚡⚡⚡ │ 推荐 ⭐        │
│ V4  │ DECR + 负数回滚       │  ⚠️  │  ⚡⚡  │ 简单场景       │
└─────┴─────────────────────┴───────┴────────┴──────────────┘

14.6 防止一人多单

14.6.1 问题

黄牛脚本:

python
import requests
for i in range(1000):
    requests.post("/seckill", data={"user_id": "myself", "item_id": 1001})

如果只判库存不判用户,1 个黄牛能抢走 1000 瓶。

14.6.2 用 Set 做用户去重

bash
# 每抢到的用户都加进 Set
SADD seckill:users:1001 user_xxx
SISMEMBER seckill:users:1001 user_xxx   # 抢之前先查

但这又回到了「两步操作不原子」的老问题——所以必须和库存扣减一起做进 Lua

14.6.3 完整 Lua(库存 + 去重 + 推消息)

lua
-- KEYS[1] = seckill:stock:{item_id}    库存
-- KEYS[2] = seckill:users:{item_id}    已购用户集合
-- KEYS[3] = seckill:orders              Stream 消息队列
-- ARGV[1] = user_id
-- ARGV[2] = item_id
-- ARGV[3] = order_id (uuid)
--
-- 返回值:
--   1  = 抢购成功
--   0  = 库存不足
--  -1  = 重复购买(已抢过)

-- ① 防重复购买
if redis.call('SISMEMBER', KEYS[2], ARGV[1]) == 1 then
  return -1
end

-- ② 检查库存
local s = tonumber(redis.call('GET', KEYS[1]))
if not s or s <= 0 then
  return 0
end

-- ③ 扣库存 + 记录用户 + 推消息(原子)
redis.call('DECR', KEYS[1])
redis.call('SADD', KEYS[2], ARGV[1])
redis.call('XADD', KEYS[3], '*',
  'order_id', ARGV[3],
  'user_id',  ARGV[1],
  'item_id',  ARGV[2],
  'ts',       redis.call('TIME')[1])

return 1

Python 端:

python
import uuid
SECKILL_LUA = open('lua/seckill.lua').read()
seckill = r.register_script(SECKILL_LUA)

def do_seckill(user_id, item_id):
    order_id = uuid.uuid4().hex
    ret = seckill(
        keys=[f"seckill:stock:{item_id}",
              f"seckill:users:{item_id}",
              "seckill:orders"],
        args=[user_id, item_id, order_id])
    return {1: "成功", 0: "卖完", -1: "已购买"}[ret]

14.6.4 这段 Lua 做对了什么?

┌─────────────────────────────────────────────┐
│   单次 EVAL 内同时完成 4 件事 + 全部原子          │
├─────────────────────────────────────────────┤
│  ① 检查用户是否买过                            │
│  ② 检查库存是否够                              │
│  ③ 扣库存 + 记录用户                           │
│  ④ 推消息进 Stream(让消费者去落库)             │
└─────────────────────────────────────────────┘

  ⏱  整个脚本执行时间 < 100 微秒
  💪 单 Redis 节点能扛 ~10 万 QPS 此类脚本

14.7 异步下单

14.7.1 为什么不直接同步写 MySQL?

同步方案:抢成功 → 立即写 MySQL → 返回用户
  问题:MySQL 写入 ~5ms,10 万 QPS 直接打爆数据库

异步方案:抢成功 → Redis 推消息 → 立即返回用户

                            └─ 消费者慢慢写 MySQL(限速 1000 TPS)
  好处:
    - 用户感知延迟 < 50ms
    - MySQL 压力恒定可控
    - 失败可重试(消息持久化)

14.7.2 用 Redis Stream 做消息队列

第 7 章学过 Stream,这里复习核心 API:

bash
# 生产(Lua 内已经 XADD)
XADD seckill:orders * order_id xxx user_id 1 item_id 1001

# 消费组
XGROUP CREATE seckill:orders order_grp $ MKSTREAM
XREADGROUP GROUP order_grp consumer-1 COUNT 100 BLOCK 2000 STREAMS seckill:orders >

# ACK
XACK seckill:orders order_grp <message-id>

14.7.3 消费者代码骨架

python
GROUP, CONSUMER, STREAM = "order_grp", "c-1", "seckill:orders"

def init():
    try:
        r.xgroup_create(STREAM, GROUP, id="0", mkstream=True)
    except redis.ResponseError as e:
        if "BUSYGROUP" not in str(e): raise

def consume_forever():
    while True:
        msgs = r.xreadgroup(GROUP, CONSUMER,
                            {STREAM: ">"}, count=100, block=2000)
        if not msgs: continue
        _, entries = msgs[0]
        for mid, data in entries:
            try:
                handle_order(data)            # 写 MySQL
                r.xack(STREAM, GROUP, mid)
            except Exception:
                pass                          # 不 ACK 即可下次重试

14.7.4 失败重试机制

正常流程:
  ① XREADGROUP 取消息(Pending List 自动登记)
  ② 处理(写 MySQL + 发短信)
  ③ XACK 确认(从 Pending List 移除)

失败处理:
  - 处理过程崩溃 → 消息留在 Pending List
  - 启动巡检任务,定时 XPENDING 查看长时间未 ACK 的消息
  - XCLAIM 重新分配给其他消费者
  - 超过 N 次重试 → 移到死信队列 + 告警

14.7.5 List vs Stream 怎么选?

┌──────────────┬──────────────┬─────────────────┐
│  特性          │  List         │  Stream         │
├──────────────┼──────────────┼─────────────────┤
│  消费组        │  无            │  ✅ 原生支持      │
│  ACK 重试      │  无            │  ✅ Pending List │
│  消费历史      │  POP 即删      │  ✅ 可重放        │
│  顺序保证      │  FIFO         │  FIFO           │
│  内存占用      │  小            │  稍大            │
│  适用          │  简单点对点     │  企业级 MQ        │
└──────────────┴──────────────┴─────────────────┘

5.0 之后新项目首选 Stream。本章实战代码用 Stream。


14.8 兜底方案(容灾)

14.8.1 Redis 挂了怎么办

故障:Redis 集群整个 GG(机房断电、版本升级 bug)
应对:
  ① 快速失败:直接返回「活动太火爆,请稍后再试」
     → 比让用户等 30 秒后超时好得多
  ② 本地内存兜底(鸡肋方案):
     - 应用启动时缓存「商品基本信息」到本地
     - 库存信息只能拒绝,不能就地扣(多机数据不一致)
  ③ 切到备用 Redis 集群(事先准备好的)
     - 库存数据可能略有损失(接受小范围超卖)

14.8.2 MQ 堆积怎么办

现象:Redis 扣减很快,但消费者写 MySQL 慢,Stream 越积越多

根因:生产 > 消费

解决:
  ① 在 Lua 里加「全局已售总量上限」
     - 超过总量直接拒绝,让 Stream 不再增长
  ② 消费者水平扩展
  ③ 消费者批量写(Pipeline 1000 条/批 → 5ms)
  ④ 极端情况:开关把秒杀直接停掉

14.8.3 数据对账(每天必做)

夜间定时任务:
  REDIS:  GET seckill:stock:1001  →  剩余 0
  REDIS:  SCARD seckill:users:1001 → 已售 1000

  MYSQL: SELECT COUNT(*) FROM orders WHERE item_id=1001 → 998

  ❗ Redis 显示 1000 单,MySQL 只有 998 → 2 单丢失!
       → 查 Stream Pending List → XPENDING + XCLAIM 重放

正确性公式:
  原始库存 = MySQL 订单数 + Redis 剩余
  1000    = 998         + 0           ❌ 不对!差 2

14.8.4 兜底方案总览

┌──────────────┬─────────────────────────────────┐
│  故障          │  兜底                             │
├──────────────┼─────────────────────────────────┤
│ CDN 故障       │ 接入层自动回源                     │
│ Nginx 单机故障 │ LB 自动剔除节点                    │
│ 应用层 OOM     │ K8s 重启 + 限流降级                │
│ Redis 主节点挂 │ 哨兵自动故障转移(详见第 9 章)       │
│ Redis 集群挂   │ 切备份集群 / 直接降级               │
│ MQ 堆积        │ Lua 限速 / 暂停秒杀                 │
│ MySQL 写入慢   │ 消费者批量写 + 限流                 │
│ 数据不一致     │ 定时对账 + 人工补录                 │
└──────────────┴─────────────────────────────────┘

14.9 完整代码走读

实战代码见 14_seckill/code/

14_seckill/code/
├── 01_init_stock.py        # 活动初始化:把库存灌入 Redis
├── 02_seckill_v1_unsafe.py # 演示超卖的 V1(多线程跑会出现负库存)
├── 03_seckill_v3_lua.py    # 用 Lua 脚本的安全版本
├── 04_seckill_full.py      # 完整方案(库存+去重+异步消息)
├── 05_consumer.py          # Stream 消费者(写 MySQL)
├── 06_load_test.py         # 压测:QPS / 失败率 / 是否超卖
└── lua/
    └── seckill.lua         # 完整 Lua 脚本

推荐运行顺序

bash
# 1. 初始化库存(设 1000 件)
python 01_init_stock.py

# 2. 跑超卖演示(你会看到负库存)
python 02_seckill_v1_unsafe.py

# 3. 重置库存,跑安全版
python 01_init_stock.py
python 03_seckill_v3_lua.py

# 4. 完整压测
python 01_init_stock.py
python 06_load_test.py

# 5. 启动消费者(另开终端)
python 05_consumer.py

📌 所有脚本都假设 Redis 在 127.0.0.1:6379,不需要密码。


14.10 性能压测

14.10.1 测试环境

硬件:
  - CPU:8 核(虚拟机)
  - 内存:16 GB
  - Redis:单节点 7.2

应用:
  - Python 3.11 + redis-py 5.0
  - 100 个并发线程,每线程 1000 次请求 = 10 万次抢购
  - 库存:1000 件

14.10.2 期望结果

┌─────────────────────┬──────────┬──────────┬──────────┐
│  方案                │  V1 不安全 │  V3 Lua   │  V4 完整  │
├─────────────────────┼──────────┼──────────┼──────────┤
│  总耗时              │  ~3.5 s  │  ~2.8 s  │  ~3.2 s  │
│  实际 QPS            │  ~28000  │  ~35000  │  ~31000  │
│  抢购成功            │  ~1080 ❌│  1000 ✅ │  1000 ✅ │
│  超卖数              │  +80 件  │   0      │   0      │
│  最终库存            │  -80     │   0      │   0      │
│  Stream 消息数        │   N/A    │   N/A    │  1000   │
└─────────────────────┴──────────┴──────────┴──────────┘

14.10.3 极限测试

单 Redis 节点 + 8 客户端机:
  纯 Lua 脚本   → ~10 万 QPS
  库存 1000 件  → 5 秒内售罄,0 超卖

集群(3 主 3 从):
  按商品 ID hash slot 分散到不同主节点
  理论峰值 ~30 万 QPS

14.10.4 性能优化要点

① 客户端批量:用 Pipeline 把多个 EVAL 打包发送
② 脚本缓存:register_script + EVALSHA,省网络带宽
③ Key 设计:保证一次秒杀涉及的 Key 都在同一个 slot(用 {hashtag})
   例:seckill:{1001}:stock / seckill:{1001}:users
④ 连接池:默认 50,按机器数调
⑤ 关闭 AOF 或用 everysec:fsync 是写性能瓶颈

14.11 本章小结:知识地图

┌────────────────────────────────────────────────────────────┐
│                第 14 章 ↔ 前 13 章 知识引用图                │
├────────────────────────────────────────────────────────────┤
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch1 IO 多路 │ → Redis 单线程为何能抗 10 万 QPS              │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch2 String │ → 库存计数 SET / DECR                         │
│   │ Ch2 Hash   │ → 商品多维信息                                │
│   │ Ch2 Set    │ → 用户去重 SADD/SISMEMBER                     │
│   │ Ch2 ZSet   │ → 滑动窗口限流                                │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch4 过期    │ → 活动 Key EXPIRE,自动清理                   │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch6 Lua    │ → 原子扣减 + 防重复 ⭐ 核心                    │
│   │ Ch6 Pipeline│ → 客户端批量发送                              │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch7 Stream │ → 异步下单消息队列 ⭐ 核心                     │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch9 Sentinel│ → Redis 高可用                               │
│   │ Ch10 Cluster│ → 分片承载更高 QPS                            │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch11 分布式锁│ → 兜底防超卖(Lua 主路径,锁是 plan B)        │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch12 缓存   │ → 库存预热(避免击穿)                        │
│   └───────────┘                                             │
│                                                              │
│   ┌───────────┐                                             │
│   │ Ch13 优化   │ → 慢查询 / 热 Key / 大 Key 治理                │
│   └───────────┘                                             │
│                                                              │
└────────────────────────────────────────────────────────────┘

一句话总结

秒杀系统 = 用 Redis 的原子能力把 99% 的请求挡在数据库之外,剩下的 1% 通过异步消息慢慢落库。


14.12 进阶思考题

Q1:如果库存有 1000 万件怎么办?

思考方向:单个 Redis Key 上的所有请求都打到同一个节点(即使是集群也只在一个 slot)。1000 万件意味着可能有 1000 万 QPS 都怼在一个 Key 上 —— 单节点扛不住。

提示:分桶(Sharding)

原始:
  seckill:stock:1001 = 10000000

分 100 桶:
  seckill:stock:1001:bucket:0  = 100000
  seckill:stock:1001:bucket:1  = 100000
  ...
  seckill:stock:1001:bucket:99 = 100000

抢购时随机选一个桶(user_id % 100),扣那个桶的库存。
失败的话再随机选下一个桶(轮询直到都没货)。

效果:QPS 摊到 100 个 slot → 100 个 Redis 主节点 → 总容量 × 100

进一步思考:「有些桶提前空了怎么办?」、「桶之间负载不均怎么办?」、「最后阶段几个零散桶要不要合并?」

Q2:如果要支持「先到先得排队」怎么实现?

思考方向:现有的 Lua 是「谁的请求先到 Redis 谁赢」,但网络延迟会让这个变得不公平。如果要真正按用户点击瞬间的时间戳排序,需要全局有序。

提示

  • 客户端把「点击时间」一并上报
  • 用 ZSet 做排队:ZADD seckill:queue:1001 click_ts user_id
  • 服务端固定速率 ZRANGE 0 N 按顺序放行(释放到「allowed」Set)
  • 抢购时在 Lua 里加判断:必须在 allowed Set 里才能扣

:客户端时间戳可被篡改 → 黄牛把时间改到 -∞ 永远排第一。 解法:服务端给排队接口加上「服务端写入时间戳」做兜底。

Q3:怎么防机器人脚本?

思考方向:脚本最大特征是「反应快、不需要思考」。你能想出多少种区分人类和脚本的方法?

提示几个方向:

  • 行为分析:鼠标轨迹、键盘按键间隔、滚动行为
  • 设备指纹:浏览器 UA、Canvas 指纹、字体列表
  • 挑战题:图形验证码 / 滑块 / 答题
  • 预约机制:必须提前 1 小时预约,限制账号活跃度
  • 风控规则:同 IP 多账号、同设备多账号
  • 延迟惩罚:可疑用户的请求人为加 500ms 延时

Q4:多机房部署怎么保证全局库存正确?

思考方向:北京机房 + 上海机房,库存怎么分?

方案 A:单写双读

  • 库存只在主机房(如北京)的 Redis 上
  • 上海机房的请求通过专线访问北京 Redis
  • 缺点:上海用户延迟高,专线挂了上海全瘫

方案 B:库存分割

  • 1000 件预先切:北京 600 件,上海 400 件
  • 各自独立扣
  • 缺点:北京卖完了上海可能还有,体验不一致

方案 C:动态调度

  • 在 B 的基础上加「调度器」:定时检查各机房库存,把闲置库存挪给紧张的机房
  • 缺点:调度有延迟,复杂度高

方案 D:CRDT / 最终一致

  • 接受短暂超卖(如多卖 0.1%)
  • 各机房独立扣减,定时同步求和对账
  • 缺点:超卖不可避免

生产实践:99% 的公司选 B + C,C 的调度逻辑是核心壁垒。


📌 结业语:恭喜你看完了 14 章。Redis 远远不止「一个缓存」——它是工程师工具箱里最锋利的瑞士军刀之一。秒杀只是它的一个秀场,签到系统、Feed 流、附近的人、限流、链路追踪 ……每一个场景都能用 Redis 写出漂亮的方案。

下一步建议:选一个你身边的真实场景(订单号生成 / 防刷接口 / 简单的搜索建议),用 Redis 重写一遍,体会从「能跑」到「优雅」的过程。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
"""
Ch14 配套代码 1 / 7 —— 活动初始化(库存预热)
================================================================

用途:
  把 MySQL 的秒杀商品库存「灌」到 Redis,活动开始前必跑。
  本脚本会同时清空上一次留下的「已购用户集合」与「订单 Stream」。

运行:
  python 01_init_stock.py            # 默认 item_id=1001 / stock=1000
  python 01_init_stock.py 1001 100   # 商品 1001 设 100 件
"""

import sys
import redis

ITEM_ID = int(sys.argv[1]) if len(sys.argv) > 1 else 1001
STOCK   = int(sys.argv[2]) if len(sys.argv) > 2 else 1000

STOCK_KEY  = f"seckill:stock:{ITEM_ID}"
USERS_KEY  = f"seckill:users:{ITEM_ID}"
ITEM_KEY   = f"seckill:item:{ITEM_ID}"
STREAM_KEY = "seckill:orders"


def main() -> None:
    r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    r.ping()

    pipe = r.pipeline()
    pipe.delete(STOCK_KEY, USERS_KEY, ITEM_KEY)
    pipe.delete(STREAM_KEY)

    pipe.set(STOCK_KEY, STOCK)
    pipe.expire(STOCK_KEY, 7200)

    pipe.hset(ITEM_KEY, mapping={
        "name":   f"秒杀商品-{ITEM_ID}",
        "price":  1499,
        "stock":  STOCK,
        "sold":   0,
        "status": "READY",
    })
    pipe.expire(ITEM_KEY, 7200)

    pipe.execute()

    print("=" * 56)
    print(f"  ✅ 商品 {ITEM_ID} 已预热")
    print(f"     库存 Key  : {STOCK_KEY} = {r.get(STOCK_KEY)}")
    print(f"     元信息 Key: {ITEM_KEY}  = {r.hgetall(ITEM_KEY)}")
    print(f"     用户集合  : {USERS_KEY} (已清空)")
    print(f"     订单 Stream: {STREAM_KEY} (已清空)")
    print("=" * 56)


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
        sys.exit(1)
python
"""
Ch14 配套代码 2 / 7 —— 超卖演示(V1 不安全版)
================================================================

用途:
  演示「GET → 判断 → DECR」三步分开做时,多线程并发会超卖。
  你会看到最终库存变成负数(-3、-7 这种),抢购成功数 > 真实库存。

原因:
  线程 A: GET stock = 1
  线程 B: GET stock = 1     ← 同一时刻读到 1
  线程 A: DECR stock → 0
  线程 B: DECR stock → -1   ← 超卖

运行:
  python 01_init_stock.py 1001 100   # 先把库存设为 100
  python 02_seckill_v1_unsafe.py
"""

import sys
import threading
import redis

ITEM_ID    = 1001
STOCK_KEY  = f"seckill:stock:{ITEM_ID}"
THREADS    = 100
PER_THREAD = 50

success_cnt = 0
lock = threading.Lock()


def seckill_v1_unsafe(r: redis.Redis) -> bool:
    """❌ 不安全:GET-判断-DECR 不原子,会超卖"""
    stock = r.get(STOCK_KEY)
    if stock is None:
        return False
    if int(stock) <= 0:
        return False
    r.decr(STOCK_KEY)
    return True


def worker() -> None:
    global success_cnt
    rr = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    local = 0
    for _ in range(PER_THREAD):
        if seckill_v1_unsafe(rr):
            local += 1
    with lock:
        success_cnt += local


def main() -> None:
    r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    initial = int(r.get(STOCK_KEY) or 0)
    if initial == 0:
        print("❌ 请先运行: python 01_init_stock.py 1001 100")
        sys.exit(1)

    print("=" * 56)
    print(f"  V1 不安全版 · {THREADS} 线程 × {PER_THREAD} 次抢购")
    print(f"  初始库存  : {initial}")
    print("=" * 56)

    ts = [threading.Thread(target=worker) for _ in range(THREADS)]
    for t in ts: t.start()
    for t in ts: t.join()

    final = int(r.get(STOCK_KEY) or 0)
    print(f"  抢购成功  : {success_cnt}")
    print(f"  剩余库存  : {final}")
    print(f"  应剩库存  : {max(0, initial - success_cnt)}")
    if final < 0:
        print(f"  ❌ 超卖 {-final} 件!这就是 V1 的问题")
    elif success_cnt > initial:
        print(f"  ❌ 抢购数 {success_cnt} > 初始 {initial},超卖")
    else:
        print(f"  🎉 这次刚好没翻车(并发不够,多跑几次)")


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch14 配套代码 3 / 7 —— 安全秒杀(V3 Lua 脚本版)
================================================================

用途:
  用 Redis 单线程执行 Lua 的特性,把「判库存 + 扣减」做成原子操作。
  无论多大并发,最终库存都不会变成负数。

核心 Lua(写在脚本字符串里,不依赖外部文件):
  local s = tonumber(redis.call('GET', KEYS[1]))
  if not s or s <= 0 then return 0 end
  redis.call('DECR', KEYS[1])
  return 1

运行:
  python 01_init_stock.py 1001 100
  python 03_seckill_v3_lua.py
"""

import sys
import threading
import redis

ITEM_ID    = 1001
STOCK_KEY  = f"seckill:stock:{ITEM_ID}"
THREADS    = 100
PER_THREAD = 50

LUA = """
local s = tonumber(redis.call('GET', KEYS[1]))
if not s or s <= 0 then
  return 0
end
redis.call('DECR', KEYS[1])
return 1
"""

success_cnt = 0
lock = threading.Lock()


def worker(script) -> None:
    global success_cnt
    rr = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    local = 0
    for _ in range(PER_THREAD):
        if script(keys=[STOCK_KEY], client=rr) == 1:
            local += 1
    with lock:
        success_cnt += local


def main() -> None:
    r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    initial = int(r.get(STOCK_KEY) or 0)
    if initial == 0:
        print("❌ 请先运行: python 01_init_stock.py 1001 100")
        sys.exit(1)

    seckill = r.register_script(LUA)

    print("=" * 56)
    print(f"  V3 Lua 安全版 · {THREADS} 线程 × {PER_THREAD} 次抢购")
    print(f"  初始库存  : {initial}")
    print("=" * 56)

    ts = [threading.Thread(target=worker, args=(seckill,)) for _ in range(THREADS)]
    for t in ts: t.start()
    for t in ts: t.join()

    final = int(r.get(STOCK_KEY) or 0)
    print(f"  抢购成功  : {success_cnt}")
    print(f"  剩余库存  : {final}")
    if final < 0:
        print(f"  ❌ 不该出现的负库存,请检查脚本")
    elif success_cnt == initial and final == 0:
        print(f"  ✅ 完美:恰好卖完 {initial} 件,0 超卖")
    else:
        print(f"  ✅ 安全:成功 {success_cnt} ≤ 初始 {initial}")


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch14 配套代码 4 / 7 —— 完整秒杀方案
================================================================

包含:
  ① 用户维度滑动窗口限流(ZSet)
  ② 库存原子扣减(Lua)
  ③ 用户去重(Set)
  ④ 异步下单消息(Stream)

整套核心逻辑只用 2 个 Lua 脚本:
  - rate_limit.lua  限流
  - seckill.lua     扣库存 + 去重 + 推消息

运行:
  python 01_init_stock.py 1001 1000
  python 04_seckill_full.py
  # 另开终端运行 python 05_consumer.py 看消费
"""

import os
import sys
import time
import uuid
import random
import threading
import redis

ITEM_ID    = 1001
STOCK_KEY  = f"seckill:stock:{ITEM_ID}"
USERS_KEY  = f"seckill:users:{ITEM_ID}"
STREAM_KEY = "seckill:orders"
RATE_LIMIT = 10        # 每用户每秒最多 10 次抢购请求
RATE_WINDOW_MS = 1000

THREADS    = 100
PER_THREAD = 20

LUA_RATE = """
-- KEYS[1]=rate:user:xxx ARGV[1]=now_ms ARGV[2]=window_ms
-- ARGV[3]=limit       ARGV[4]=uuid
local now = tonumber(ARGV[1])
local win = tonumber(ARGV[2])
local lim = tonumber(ARGV[3])
redis.call('ZREMRANGEBYSCORE', KEYS[1], 0, now - win)
local cnt = redis.call('ZCARD', KEYS[1])
if cnt >= lim then return 0 end
redis.call('ZADD', KEYS[1], now, ARGV[4])
redis.call('PEXPIRE', KEYS[1], win)
return 1
"""

LUA_FILE = os.path.join(os.path.dirname(__file__), "lua", "seckill.lua")

stats = {"ok": 0, "no_stock": 0, "duplicate": 0, "rate_limited": 0, "no_item": 0}
stats_lock = threading.Lock()


def do_seckill(r, rate_script, sk_script, user_id):
    now_ms = int(time.time() * 1000)
    rid = uuid.uuid4().hex
    if rate_script(keys=[f"rate:user:{user_id}"],
                   args=[now_ms, RATE_WINDOW_MS, RATE_LIMIT, rid],
                   client=r) == 0:
        return "rate_limited"

    order_id = uuid.uuid4().hex
    ret = sk_script(keys=[STOCK_KEY, USERS_KEY, STREAM_KEY],
                    args=[user_id, ITEM_ID, order_id],
                    client=r)
    return {1: "ok", 0: "no_stock", -1: "duplicate", -2: "no_item"}[ret]


def worker(rate_script, sk_script):
    rr = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    local = {"ok": 0, "no_stock": 0, "duplicate": 0, "rate_limited": 0, "no_item": 0}
    for _ in range(PER_THREAD):
        user_id = f"user_{random.randint(1, THREADS * 5)}"
        local[do_seckill(rr, rate_script, sk_script, user_id)] += 1
    with stats_lock:
        for k, v in local.items():
            stats[k] += v


def main():
    r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    initial = int(r.get(STOCK_KEY) or 0)
    if initial == 0:
        print("❌ 请先运行: python 01_init_stock.py 1001 1000")
        sys.exit(1)

    rate_script = r.register_script(LUA_RATE)
    sk_script   = r.register_script(open(LUA_FILE).read())

    print("=" * 60)
    print(f"  完整秒杀方案 · {THREADS} 线程 × {PER_THREAD} 次")
    print(f"  初始库存 = {initial}")
    print("=" * 60)

    t0 = time.time()
    ts = [threading.Thread(target=worker, args=(rate_script, sk_script))
          for _ in range(THREADS)]
    for t in ts: t.start()
    for t in ts: t.join()
    cost = time.time() - t0

    final = int(r.get(STOCK_KEY) or 0)
    sold = initial - final
    total = sum(stats.values())
    qps = total / cost if cost > 0 else 0
    queue_len = r.xlen(STREAM_KEY)

    print(f"  总请求    : {total}")
    print(f"  耗时      : {cost:.2f}s   QPS ≈ {qps:.0f}")
    print(f"  抢购成功  : {stats['ok']}")
    print(f"  库存不足  : {stats['no_stock']}")
    print(f"  重复购买  : {stats['duplicate']}")
    print(f"  被限流    : {stats['rate_limited']}")
    print(f"  剩余库存  : {final}  (卖出 {sold})")
    print(f"  Stream 消息 : {queue_len}")
    print(f"  ✅ 一致性   : 卖出={sold} == 成功={stats['ok']} == 消息={queue_len}"
          if sold == stats['ok'] == queue_len else
          "  ❌ 数据不一致,请检查")


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
python
"""
Ch14 配套代码 5 / 7 —— 异步下单消费者
================================================================

用途:
  从 Stream `seckill:orders` 消费秒杀消息,模拟「写订单到 MySQL」。
  - 用消费组保证「一个消息只被一个消费者处理」
  - 处理完才 XACK,崩溃重启后自动重做未 ACK 的消息

运行:
  python 05_consumer.py
  Ctrl+C 退出
"""

import sys
import time
import signal
import redis

STREAM_KEY = "seckill:orders"
GROUP      = "order_grp"
CONSUMER   = "consumer-1"

stop = False


def init_group(r: redis.Redis) -> None:
    try:
        r.xgroup_create(STREAM_KEY, GROUP, id="0", mkstream=True)
        print(f"✅ 创建消费组 {GROUP}")
    except redis.ResponseError as e:
        if "BUSYGROUP" in str(e):
            print(f"ℹ️  消费组 {GROUP} 已存在")
        else:
            raise


def write_to_mysql(order_id: str, user_id: str, item_id: str) -> None:
    """模拟写 MySQL:实际项目这里是 INSERT INTO orders ..."""
    time.sleep(0.001)


def send_sms(user_id: str, item_id: str) -> None:
    """模拟发短信"""
    time.sleep(0.0005)


def handle_signal(signum, frame):
    global stop
    stop = True
    print("\n收到退出信号,优雅退出中…")


def main():
    signal.signal(signal.SIGINT, handle_signal)
    signal.signal(signal.SIGTERM, handle_signal)

    r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    r.ping()
    init_group(r)

    handled = 0
    failed = 0
    print("=" * 56)
    print(f"  消费者 {CONSUMER} 开始消费 (Ctrl+C 退出)")
    print("=" * 56)

    while not stop:
        try:
            msgs = r.xreadgroup(GROUP, CONSUMER,
                                {STREAM_KEY: ">"},
                                count=200, block=2000)
        except redis.ConnectionError:
            time.sleep(1)
            continue

        if not msgs:
            print(f"  …等待新消息  (已处理 {handled})")
            continue

        _, entries = msgs[0]
        for mid, data in entries:
            try:
                write_to_mysql(data["order_id"], data["user_id"], data["item_id"])
                send_sms(data["user_id"], data["item_id"])
                r.xack(STREAM_KEY, GROUP, mid)
                handled += 1
                if handled % 100 == 0:
                    print(f"  ✓ 已处理 {handled} 单, 最新: order={data['order_id'][:8]} user={data['user_id']}")
            except Exception as e:
                failed += 1
                print(f"  ❌ 处理失败 {mid}: {e}(不 ACK,下次重试)")

    pending = r.xpending(STREAM_KEY, GROUP).get("pending", 0) if False else 0
    print(f"\n退出。处理成功 {handled} 单,失败 {failed} 单")


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
        sys.exit(1)
python
"""
Ch14 配套代码 6 / 7 —— 压测脚本
================================================================

用途:
  对完整秒杀方案做压测,输出:
    - 总耗时 / QPS
    - 成功 / 失败分类计数
    - 是否超卖
    - Stream 消息数 vs 抢购成功数(一致性校验)

运行:
  python 01_init_stock.py 1001 1000
  python 06_load_test.py [threads] [per_thread]
  # 默认 200 线程 × 100 次 = 2 万请求
"""

import os
import sys
import time
import uuid
import random
import threading
from collections import Counter
import redis

ITEM_ID    = 1001
STOCK_KEY  = f"seckill:stock:{ITEM_ID}"
USERS_KEY  = f"seckill:users:{ITEM_ID}"
STREAM_KEY = "seckill:orders"

THREADS    = int(sys.argv[1]) if len(sys.argv) > 1 else 200
PER_THREAD = int(sys.argv[2]) if len(sys.argv) > 2 else 100
USER_POOL  = THREADS * 5

LUA_FILE = os.path.join(os.path.dirname(__file__), "lua", "seckill.lua")

results = Counter()
latencies = []
lat_lock = threading.Lock()


def worker(script):
    rr = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    local = Counter()
    local_lat = []
    for _ in range(PER_THREAD):
        user_id = f"u_{random.randint(1, USER_POOL)}"
        order_id = uuid.uuid4().hex
        t = time.perf_counter()
        ret = script(keys=[STOCK_KEY, USERS_KEY, STREAM_KEY],
                     args=[user_id, ITEM_ID, order_id],
                     client=rr)
        local_lat.append((time.perf_counter() - t) * 1000)
        local[{1: "ok", 0: "no_stock", -1: "dup", -2: "no_item"}[ret]] += 1
    with lat_lock:
        results.update(local)
        latencies.extend(local_lat)


def percentile(data, p):
    if not data: return 0
    data = sorted(data)
    idx = int(len(data) * p)
    return data[min(idx, len(data) - 1)]


def main():
    r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    initial = int(r.get(STOCK_KEY) or 0)
    if initial == 0:
        print("❌ 请先运行: python 01_init_stock.py 1001 1000")
        sys.exit(1)

    script = r.register_script(open(LUA_FILE).read())

    print("=" * 64)
    print(f"  压测开始 · {THREADS} 线程 × {PER_THREAD} 请求 = {THREADS*PER_THREAD}")
    print(f"  初始库存  : {initial}")
    print(f"  用户池    : {USER_POOL}")
    print("=" * 64)

    t0 = time.time()
    ts = [threading.Thread(target=worker, args=(script,)) for _ in range(THREADS)]
    for t in ts: t.start()
    for t in ts: t.join()
    cost = time.time() - t0

    total = sum(results.values())
    final = int(r.get(STOCK_KEY) or 0)
    sold = initial - final
    msgs = r.xlen(STREAM_KEY)

    print()
    print("--- 性能 ---")
    print(f"  总请求    : {total}")
    print(f"  总耗时    : {cost:.2f} s")
    print(f"  QPS       : {total/cost:.0f}")
    print(f"  延迟 P50  : {percentile(latencies, 0.50):.2f} ms")
    print(f"  延迟 P95  : {percentile(latencies, 0.95):.2f} ms")
    print(f"  延迟 P99  : {percentile(latencies, 0.99):.2f} ms")

    print()
    print("--- 结果分类 ---")
    print(f"  抢购成功  : {results['ok']}")
    print(f"  库存不足  : {results['no_stock']}")
    print(f"  重复购买  : {results['dup']}")
    print(f"  商品不存在 : {results['no_item']}")

    print()
    print("--- 库存校验 ---")
    print(f"  剩余库存   : {final}")
    print(f"  实际卖出   : {sold}")
    print(f"  Stream 消息: {msgs}")
    if final < 0:
        print(f"  ❌ 出现负库存(超卖 {-final})!")
    elif sold == results['ok'] == msgs:
        print(f"  ✅ 完美一致:库存 / 成功 / 消息 三方对齐")
    else:
        print(f"  ⚠️  数据有偏差,请检查")


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
txt
--[[
  完整秒杀 Lua 脚本
  ===========================================================
  KEYS[1] = seckill:stock:{item_id}     库存计数 (String)
  KEYS[2] = seckill:users:{item_id}     已购用户集合 (Set)
  KEYS[3] = seckill:orders              订单消息流 (Stream)
  ARGV[1] = user_id
  ARGV[2] = item_id
  ARGV[3] = order_id  (uuid)

  返回值:
     1   抢购成功
     0   库存不足
    -1   重复购买
    -2   商品不存在 / 未预热
  ===========================================================
  说明:
    - Redis 单线程执行 Lua,整个脚本期间不会被打断
    - 一次脚本同时完成「防重 + 判库存 + 扣减 + 推消息」4 件事
    - 失败用负数区分原因,方便客户端做不同提示
]]

if redis.call('EXISTS', KEYS[1]) == 0 then
  return -2
end

if redis.call('SISMEMBER', KEYS[2], ARGV[1]) == 1 then
  return -1
end

local s = tonumber(redis.call('GET', KEYS[1]))
if not s or s <= 0 then
  return 0
end

redis.call('DECR', KEYS[1])
redis.call('SADD', KEYS[2], ARGV[1])
redis.call('XADD', KEYS[3], '*',
  'order_id', ARGV[3],
  'user_id',  ARGV[1],
  'item_id',  ARGV[2],
  'ts',       redis.call('TIME')[1])

return 1

01_init_stock.py ↗ · 02_seckill_v1_unsafe.py ↗ · 03_seckill_v3_lua.py ↗ · 04_seckill_full.py ↗ · 05_consumer.py ↗ · 06_load_test.py ↗ · lua/seckill.lua ↗