WR Blog 加载中...
返回文章
数据库Redis分布式锁消息队列

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams封面

数据库基础体系 · 第 39/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 不只是缓存。它还提供了一组建立在内存数据结构、单命令原子性、Lua 脚本和持久化日志之上的协调能力:

  • 用键值状态实现带租约的分布式锁;
  • 用 Lua 把“读取—判断—写回”合并为不可插入的原子操作;
  • 用计数器、时间戳和令牌桶实现限流;
  • 用 Pub/Sub 传递实时但不持久的通知;
  • 用 Streams 保存消息,并通过消费者组实现可恢复的消费进度。

这些能力的可靠性边界并不相同。锁和限流主要依赖 Redis 主节点上操作的原子性;Pub/Sub 依赖在线连接;Streams 则把消息和消费状态保存为 Redis 数据,因此可以处理消费者短暂故障。将它们混为“Redis 消息队列”或“Redis 事务”会导致错误的设计。


一、共同基础:Redis 的原子性到底保证什么

1. Redis 命令的串行执行

对一个 Redis 实例来说,普通命令在主线程中按顺序执行。假设当前键 counter 的值为 10,两个客户端同时执行:

客户端 A:GET counter
客户端 B:GET counter
客户端 A:SET counter 11
客户端 B:SET counter 11

GETSET 各自是原子的,但整个“读取后加一”不是原子的,最终结果可能是 11 而不是期望的 12

使用单条命令:

INCR counter

或者使用 Lua 脚本,可以让一组操作在执行期间不被其他 Redis 命令插入。

这里的“原子”表示:

  1. 脚本或命令执行期间,其他客户端命令不会交错执行;
  2. 其他客户端不会看到脚本执行过程中的中间状态;
  3. 它不表示业务操作具备跨系统事务;
  4. 它也不表示脚本发生错误后所有已经执行的写入会自动回滚。

Redis 的 Lua 脚本应当短小、确定且不执行阻塞操作。脚本执行时间过长会阻塞整个实例上的其他命令。脚本中的 redis.call 出错时,脚本会报告错误;脚本此前已经产生的写入不能依赖“自动回滚”来恢复,因此应在写入前完成参数校验。

2. MULTI/EXEC 与 Lua 的区别

Redis 事务通常指:

MULTI
SET a 1
SET b 2
EXEC

MULTI 之后命令会先排队,EXEC 时连续执行。事务执行期间,其他客户端命令不能插入这批命令。

但 Redis 事务不是关系数据库意义上的完整 ACID 事务:

  • 没有传统意义的回滚机制;
  • WATCH 提供的是乐观并发控制;
  • 它不自动覆盖外部数据库、文件系统或 HTTP 调用;
  • Redis Cluster 中事务涉及的键通常必须位于同一个哈希槽。

Lua 更适合“根据当前值判断后再写入”的逻辑,例如:

读取锁值
判断是否属于当前持有者
删除锁

如果拆成 GETDEL,中间可能被其他客户端插入;如果放入一个脚本,整个判断和删除不可被插入。


二、分布式锁:带租约的互斥,而不是永久所有权

1. 锁的基本模型

分布式锁通常需要满足三个条件:

  1. 互斥性:同一时刻至多一个客户端被认为持有锁;
  2. 释放安全:客户端只能释放自己持有的锁;
  3. 故障可恢复:持锁客户端崩溃后,锁最终可以再次获取。

在 Redis 中,最基本的加锁命令是:

SET lock:report:daily 9f3c... NX PX 10000

参数含义:

  • NX:仅当键不存在时设置;
  • PX 10000:设置 10 秒毫秒级过期时间;
  • 9f3c...:本次加锁请求生成的随机令牌。

成功时返回:

OK

失败时通常返回:

(nil)

不能使用下面这种两步操作代替:

SETNX lock:report:daily 1
EXPIRE lock:report:daily 10

因为客户端可能在两条命令之间崩溃,留下永不过期的锁。

2. 为什么锁值必须是随机令牌

假设客户端 A 获取锁:

lock:job = token-A

A 执行任务超过租约时间,锁自动过期。之后客户端 B 获取:

lock:job = token-B

如果 A 这时直接执行:

DEL lock:job

A 删除的实际上是 B 的锁。

因此释放锁必须是“比较令牌并删除”的原子操作:

-- unlock.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('DEL', KEYS[1])
end
return 0

执行:

EVAL "$(cat unlock.lua)" 1 lock:job token-A

返回值:

  • 1:比较成功并删除;
  • 0:键不存在,或者当前令牌不是 token-A

不能写成:

GET lock:job
如果等于 token-A,再执行 DEL

因为 GETDEL 之间可能发生过期、重新加锁或其他客户端写入。

3. 租约过期后的真实边界

Redis 锁的过期时间使它成为带租约的锁。它只能保证:

在 Redis 认为租约仍有效、且部署故障模型满足假设时,持有者拥有锁。

它不能保证客户端在业务代码中永远拥有锁。例如:

  1. A 获取锁,租约为 10 秒;
  2. A 因为 GC、进程暂停或网络阻塞,20 秒后才继续执行;
  3. Redis 中的锁早已过期;
  4. B 已经重新获取锁并修改资源;
  5. A 恢复后继续写入。

即使 A 后续调用安全解锁脚本,也只能避免删除 B 的锁,不能阻止 A 已经对业务资源执行过期写入。

因此,对于“过期持有者不能覆盖新持有者”的严格要求,需要栅栏令牌(fencing token)

4. 栅栏令牌:把锁状态传递给资源

获取锁时,同时从 Redis 生成单调递增的令牌:

MULTI
INCR lock:job:fence
SET lock:job token-A NX PX 10000
EXEC

实际应用需要处理 SET 失败时不要误用已经生成的令牌。更常见的方式是把“生成令牌”和“抢锁”封装在 Lua 中:

-- acquire-with-fence.lua
if redis.call('EXISTS', KEYS[1]) == 1 then
    return {0, false}
end

local token = redis.call('INCR', KEYS[2])
redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[2])
return {1, token}

调用:

EVAL "$(cat acquire-with-fence.lua)" 2 lock:job lock:job:fence token-A 10000

成功可能返回:

1
1

客户端把令牌 1 传给真正的资源服务。资源服务保存“最后接受的令牌”,只接受更大的令牌:

请求 A:fence=1,接受
请求 B:fence=2,接受
请求 A:fence=1,拒绝

这样,Redis 只负责发放单调令牌,资源本身负责拒绝旧持有者。若资源是数据库,可以通过条件更新实现:

UPDATE job_state
SET result = :result,
    last_fence = :fence
WHERE job_id = :job_id
  AND last_fence < :fence;

这一步必须和资源更新处于同一个数据库事务中,否则检查和更新之间仍可能被插入。

5. 续租与释放

如果任务可能超过初始租约,应由持有者定期续租。续租也必须比较令牌:

-- renew.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('PEXPIRE', KEYS[1], ARGV[2])
end
return 0

执行:

EVAL "$(cat renew.lua)" 1 lock:job token-A 10000

续租周期不能等到租约最后一刻才执行,否则网络延迟或进程暂停可能造成误过期。与此同时,业务代码仍应设计为:续租失败后停止产生新的副作用,而不是继续假设自己拥有锁。

6. 主从复制和故障转移边界

在单个 Redis 主节点上,SET NX PX 的互斥判断是原子的。但如果主节点宕机,副本尚未复制最新锁,故障转移后副本可能认为锁不存在,另一个客户端便可获取同名锁。

因此:

  • Redis 主从复制默认是异步的;
  • WAIT 可以等待写入传播到若干副本,但不等于共识协议;
  • 网络分区、故障转移和旧主恢复时,不能仅凭普通 Redis 锁得到严格的线性一致互斥;
  • 对高价值资源,栅栏令牌、数据库约束或具备共识语义的协调系统仍然重要。

把多个 Redis 节点拼成所谓“更强的锁”也不会自动消除时钟、网络分区和客户端暂停问题。应先明确业务需要的是“尽量避免重复执行”,还是“在故障条件下也绝不允许旧持有者写入”,两者的系统设计不同。


三、Lua:把条件、读取和写入绑定为一个状态转换

1. 脚本的输入边界

Lua 脚本通过两类参数接收数据:

  • KEYS:脚本访问的 Redis 键;
  • ARGV:普通参数。

例如:

local current = redis.call('GET', KEYS[1])
local limit = tonumber(ARGV[1])

if current and tonumber(current) >= limit then
    return 0
end

redis.call('INCR', KEYS[1])
return 1

调用:

EVAL "$(cat increment-if-below.lua)" 1 quota:user:42 10

在 Redis Cluster 中,脚本访问的键必须满足集群路由要求,通常应位于同一个哈希槽。可以使用哈希标签:

quota:{user-42}:count
quota:{user-42}:timestamp

{user-42} 使两个键使用相同的哈希标签。脚本中应把所有键放入 KEYS,不要把键名藏在 ARGV 中。这样 Redis 才能正确检查和路由脚本涉及的键。

2. EVAL、EVALSHA 与脚本缓存

直接执行:

EVAL "return redis.call('INCR', KEYS[1])" 1 counter

也可以先加载脚本:

SCRIPT LOAD "return redis.call('INCR', KEYS[1])"

Redis 返回一个 SHA-1 摘要,例如:

"7d9..."

之后执行:

EVALSHA 7d9... 1 counter

EVALSHA 只引用当前实例脚本缓存中的脚本。重启、故障转移或连接到另一个节点后,可能出现 NOSCRIPT。客户端必须能够重新 SCRIPT LOAD 或回退到 EVAL。脚本缓存不是持久化的业务配置中心。

对于固定部署的脚本,也可以使用 Redis Functions。它们适合把函数随函数库加载和管理,但仍然受到脚本执行原子性、执行时长和集群键约束等边界限制。不能因为使用 Functions 就获得跨服务事务。

3. 脚本中的时间和随机性

涉及时间的脚本必须考虑时间来源:

  • 使用客户端传入的时间,可能受到不同机器时钟偏差影响;
  • 使用 Redis 时间,可减少客户端时钟差异,但脚本仍应保持短小,并遵守当前 Redis 版本对脚本可调用命令和复制的限制;
  • 不应让不同客户端用完全不一致的时间更新同一个限流状态。

对于需要可复制、可重放的脚本,应避免依赖不可控的随机行为和不确定的外部状态。脚本最好把必要的时间、数量和策略参数显式作为输入。


四、限流:从固定窗口到令牌桶

限流的目标是限制某个主体在一个时间范围内可执行的请求数。主体可以是用户、IP、API、租户或资源键。

需要先区分两个量:

  • 吞吐限制:长期平均每秒允许多少请求;
  • 突发容量:短时间内最多允许多少请求。

固定窗口容易实现,但窗口边界可能产生突发;令牌桶能分别表达这两个量。

1. 固定窗口计数器

目标:每个用户每 60 秒最多请求 3 次。

Lua 脚本:

-- fixed-window.lua
local count = redis.call('INCR', KEYS[1])

if count == 1 then
    redis.call('EXPIRE', KEYS[1], ARGV[1])
end

if count <= tonumber(ARGV[2]) then
    return {1, count, tonumber(ARGV[2]) - count}
end

return {0, count, 0}

调用:

EVAL "$(cat fixed-window.lua)" 1 rate:user:42:window 60 3

可能返回:

1
2
1

含义是:

允许,当前窗口计数为 2,剩余 1 次

脚本中只有第一次计数时设置过期时间,因此不会因为后续请求不断刷新 TTL 而形成永不过期的计数器。

但固定窗口存在边界问题。假设窗口为:

[00:00:00, 00:01:00)
[00:01:00, 00:02:00)

用户在 00:00:59 请求 3 次,又在 00:01:00 请求 3 次。两秒内可以通过 6 次请求,虽然每个窗口分别不超过 3 次。这不是实现错误,而是固定窗口算法的定义结果。

2. 滑动窗口日志

滑动窗口要求检查最近 60 秒内的请求数。可以使用 Sorted Set:

  • member:请求唯一 ID;
  • score:请求时间戳。

每次请求执行:

删除 score <= now - 60000 的成员
统计剩余成员
若小于限制则加入当前请求

对应 Lua:

-- sliding-window.lua
local now = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local limit = tonumber(ARGV[3])
local request_id = ARGV[4]

redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', now - window)
local count = redis.call('ZCARD', KEYS[1])

if count >= limit then
    redis.call('PEXPIRE', KEYS[1], window)
    return {0, count}
end

redis.call('ZADD', KEYS[1], now, request_id)
redis.call('PEXPIRE', KEYS[1], window)
return {1, count + 1}

调用示例:

EVAL "$(cat sliding-window.lua)" 1 rate:user:42:log 1700000000000 60000 3 req-001

它比固定窗口精确,但每个请求都需要维护一个成员,内存和删除成本更高。request_id 必须唯一;如果相同 member 重复使用,ZADD 会更新原成员而不是增加计数。

3. 令牌桶

令牌桶用两个参数表达限制:

  • capacity:桶的最大令牌数,表示允许的突发容量;
  • rate:每秒补充的令牌数,表示长期平均速率。

状态为:

(tokens, last_timestamp)

在时间 now 到来时,先补充:

tokens=min(capacity, tokens+(nowlast)×rate)tokens' = \min(capacity,\ tokens + (now-last)\times rate)

其中:

  • tokens 是上次计算后的剩余令牌;
  • last 是上次更新时间;
  • rate 的单位是“令牌/秒”;
  • now-last 必须换算成秒;
  • capacity 防止空闲期间无限积累。

若请求需要 cost 个令牌:

tokens' >= cost  => 允许,并保存 tokens' - cost
tokens' <  cost  => 拒绝,并保存补充后的 tokens'

一个使用毫秒时间的 Lua 实现如下:

-- token-bucket.lua
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])       -- tokens per second
local now_ms = tonumber(ARGV[3])
local cost = tonumber(ARGV[4])
local ttl_ms = tonumber(ARGV[5])

local values = redis.call('HMGET', KEYS[1], 'tokens', 'ts')
local tokens = tonumber(values[1])
local last_ms = tonumber(values[2])

if tokens == nil then
    tokens = capacity
    last_ms = now_ms
end

-- 客户端时钟回拨时,不让令牌数量因负时间差增加异常
if now_ms < last_ms then
    now_ms = last_ms
end

local elapsed_seconds = (now_ms - last_ms) / 1000.0
tokens = math.min(capacity, tokens + elapsed_seconds * rate)

local allowed = 0
local retry_after_ms = 0

if tokens >= cost then
    tokens = tokens - cost
    allowed = 1
else
    local missing = cost - tokens
    retry_after_ms = math.ceil(missing / rate * 1000)
end

redis.call('HSET', KEYS[1], 'tokens', tokens, 'ts', now_ms)
redis.call('PEXPIRE', KEYS[1], ttl_ms)

return {allowed, tokens, retry_after_ms}

调用:

EVAL "$(cat token-bucket.lua)" 1 rate:user:42:bucket 10 2 1700000000000 1 10000

参数含义是:

  • 容量 10;
  • 每秒补充 2 个令牌;
  • 当前时间戳为 1700000000000 毫秒;
  • 本次请求消耗 1 个令牌;
  • 状态空闲 10 秒后过期。

返回:

1
9
0

表示允许请求,剩余约 9 个令牌,不需要等待。

如果返回:

0
0.5
250

表示拒绝,当前只有 0.5 个令牌,理论上约 250 毫秒后才能补足一个令牌。

这个示例把当前时间作为参数传入,因此调用方必须统一时间来源。若不同应用服务器的时钟偏差很大,同一个桶的状态会出现不一致。可以由 Redis 提供时间来源,或者由网关统一生成时间;关键是不要让多个不受控的时钟共同推进同一个限流状态。

令牌桶脚本中的 rate 不能为 0,否则拒绝请求时计算等待时间会发生除零。生产实现还应限制参数范围,防止客户端传入极大的容量、速率或 TTL。

4. 限流失败路径

限流通常有三种结果:

  1. Redis 判断允许,业务请求继续;
  2. Redis 判断拒绝,返回 HTTP 429 Too Many Requests
  3. Redis 不可用或超时。

第三种不能简单等同于“限流通过”。若放行,可能在 Redis 故障时失去保护;若拒绝,可能把 Redis 故障扩大为全站不可用。应根据接口风险决定:

  • 登录、验证码、支付等安全敏感接口偏向失败关闭;
  • 非关键读接口可能使用本地兜底限流;
  • 本地兜底必须明确它只是每实例限流,不是全局限流。

五、Pub/Sub:在线广播,不是可靠队列

1. 基本数据流

Redis Pub/Sub 有两个角色:

  • 发布者:向频道发送消息;
  • 订阅者:通过长连接订阅频道并接收消息。

订阅:

SUBSCRIBE cache:product:updated

成功后 Redis 会返回订阅确认,例如:

subscribe
cache:product:updated
1

发布:

PUBLISH cache:product:updated '{"product_id":42}'

返回值是当时收到该消息的订阅客户端数量,例如:

(integer) 2

这个返回值不是消费成功数,也不是持久化成功数。它只表示消息发布时有多少订阅连接被送达。

发布数据流可以表示为:

发布者
  |
  | PUBLISH
  v
Redis 当前连接集合
  |
  +--> 订阅者 A
  +--> 订阅者 B

Redis 不把 Pub/Sub 消息作为普通键保存。订阅者掉线期间发布的消息不会等待它重新连接后补发。

2. 可靠性语义

Redis Pub/Sub 通常应理解为:

在线连接上的即时、尽力发送、至多一次交付通知。

常见失败路径:

  1. 发布者发布时没有订阅者:消息直接消失;
  2. 订阅者网络断开:断线期间的消息不会补发;
  3. 订阅者收到消息后进程崩溃:没有 Redis 侧确认和重投;
  4. 消费者处理很慢:消息可能在客户端连接缓冲区积压,最终导致连接问题。

因此,以下场景适合 Pub/Sub:

  • 缓存失效通知;
  • 配置刷新提示;
  • 在线 WebSocket 广播;
  • 不要求离线补偿的状态变化通知。

以下场景不应只使用 Pub/Sub:

  • 订单、支付、库存等不可丢失事件;
  • 必须重新消费的任务;
  • 需要确认、重试和消费进度的消息。

如果一个消费者错过通知后可以通过读取当前状态自行修复,Pub/Sub 可以作为低延迟通知通道。例如收到“商品更新”后删除本地缓存;即使通知丢失,下一次 TTL 到期或主动校验仍可修复。若消息本身就是唯一业务事实,则应使用 Streams 或其他持久化消息系统。

3. 频道订阅与模式订阅

普通订阅:

SUBSCRIBE tenant:42:events

模式订阅:

PSUBSCRIBE tenant:*:events

模式订阅更灵活,但匹配和消息分发成本更高。频道命名应包含明确的租户、业务和版本边界,避免所有业务共享一个频道后由客户端自行过滤。

Redis 还提供按分片路由的 Sharded Pub/Sub 命令,适用于 Redis Cluster 中希望减少跨节点广播的场景。它与普通 Pub/Sub 的连接、拓扑和客户端支持要求不同,不能假设所有客户端都自动兼容;使用前应确认客户端对对应命令和集群重定向的支持。


六、Streams:可持久化的追加日志与消费者组

1. Stream 是什么

Redis Stream 是一个按 ID 排序的追加记录结构。每条记录包含:

消息 ID -> field/value 字段

写入一条消息:

XADD orders * order_id 1001 status created

返回类似:

"1700000000000-0"

* 表示由 Redis 生成 ID。ID 通常包含毫秒时间部分和序列部分,不能把它当作业务时间或全局业务排序的替代品。消费者应保存和比较 Stream ID,而不是仅用时间戳。

读取:

XRANGE orders - +

可能返回:

1) 1) "1700000000000-0"
   2) 1) "order_id"
      2) "1001"
      3) "status"
      4) "created"

生产环境通常需要限制 Stream 长度,例如:

XADD orders MAXLEN ~ 100000 * order_id 1001 status created

~ 表示近似裁剪,Redis 可以用更高效的方式接近目标长度;它不是严格保证长度永远不超过 100000。若必须精确限制,可以去掉 ~,但裁剪成本和写入性能边界不同。

2. 普通消费者读取

从头读取:

XREAD COUNT 10 STREAMS orders 0-0

从某个 ID 之后读取:

XREAD COUNT 10 STREAMS orders 1700000000000-0

阻塞读取:

XREAD BLOCK 5000 COUNT 10 STREAMS orders $

这里的 $ 表示只关注执行命令时之后新增的消息。它适合实时监听,但如果客户端启动或重连时使用 $,命令执行前已经写入的消息不会被读取。可靠消费者应保存上次成功处理的 ID,并从该 ID 继续读取,而不是每次重连都使用 $

3. 消费者组和消费确认

消费者组让多个消费者共同处理一个 Stream,并记录组内消费位置。

创建组:

XGROUP CREATE orders order-workers 0-0 MKSTREAM

含义:

  • Stream:orders
  • 消费者组:order-workers
  • 0-0 开始;
  • 如果 Stream 不存在则创建。

消费者以名字加入:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 BLOCK 5000 STREAMS orders >

> 表示读取当前组尚未投递给任何消费者的新消息。

Redis 返回的消息会进入该消费者组的待处理列表,称为 PEL(Pending Entries List)。此时消息已经投递,但还没有确认。

业务处理成功后确认:

XACK orders order-workers 1700000000000-0

XACK 只表示该消费者组不再把该消息视为待确认消息,不会删除 Stream 中的原始消息。消息是否仍能被其他读取方式看到,取决于 Stream 的保留和裁剪策略。

4. 失败、重试和消息转移

考虑下面的过程:

1. Redis 将消息投递给 worker-1
2. worker-1 执行业务操作
3. worker-1 在 XACK 前崩溃

消息仍然在该组的 PEL 中,不会因为消费者进程消失而自动丢失。其他消费者可以检查:

XPENDING orders order-workers

它会报告待处理消息数量、最小和最大 ID,以及按消费者统计的数量。

较新的 Redis 稳定语义中,可以用 XAUTOCLAIM 把长时间未处理的消息转移给当前消费者:

XAUTOCLAIM orders order-workers worker-2 60000 0-0 COUNT 10

这里的 60000 是最小空闲时间,单位为毫秒。只有空闲超过该时间的待处理消息才适合被认领。认领前应考虑 worker-1 只是暂时变慢,避免多个消费者频繁争抢同一消息。

另一种方式是:

XCLAIM orders order-workers worker-2 60000 1700000000000-0

它适合已经明确知道要转移哪些 ID 的场景。

5. Streams 通常是至少一次,而不是恰好一次

Streams 消费者组的典型语义是至少一次投递

投递消息
执行外部业务
消费者在确认前崩溃
消息被重新认领
再次执行外部业务

因此必须允许重复处理。常见做法是使用业务唯一 ID 和幂等约束:

CREATE TABLE processed_event (
    event_id VARCHAR(100) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

处理时,在同一个数据库事务中:

1. 尝试插入 event_id
2. 若主键冲突,说明已经处理过,跳过业务副作用
3. 若插入成功,执行业务更新
4. 数据库事务提交
5. 再执行 XACK

若第 3 步成功但第 5 步前进程崩溃,消息会再次投递;第 1 步的唯一约束会使第二次处理变成安全跳过。

但如果先 XACK 再更新数据库,进程在两者之间崩溃,就可能出现“消息已确认、业务未完成”。所以确认位置不是越早越好,而是应放在业务状态已经可靠提交之后。

XREADGROUPNOACK 选项可以减少 PEL 记录:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 NOACK STREAMS orders >

但它也意味着消息投递后不再等待确认,消费者崩溃时不能依赖 PEL 恢复。它只适合允许丢失或已具备其他可靠来源的场景。

6. Stream 裁剪与 PEL 的关系

裁剪 Stream 不等于自动完成消费确认。若消息从 Stream 主体中被裁剪,消费者组的待处理状态仍可能需要单独管理,具体行为取决于所使用的裁剪和 Redis 版本语义。生产系统需要同时监控:

  • Stream 长度;
  • 消费者组的组内最后 ID;
  • PEL 数量;
  • 最老待处理消息的空闲时间;
  • 认领和重试次数。

仅仅看到 Stream 长度正常,并不能说明消费者没有积压。


七、Pub/Sub 与 Streams 的选择推导

可以从“消息事实是否需要在未来重新获得”开始判断。

情形一:消息只是状态变化提醒

例如:

数据库商品价格发生变化

消费者即使错过事件,也可以重新读取数据库当前价格。因此可以使用:

数据库作为事实来源
Redis Pub/Sub 作为低延迟通知
缓存 TTL 或定期校验作为兜底

这与缓存一致性中的“失效通知”相匹配:通知丢失不会永久丢失业务事实。

情形二:消息本身就是业务事实

例如:

订单已创建
支付成功
库存扣减任务

消费者不能只看当前状态推导出所有历史动作,且必须重试和追踪进度。这时需要 Streams:

XADD 写入事实
XREADGROUP 投递
业务事务提交
XACK 确认
XPENDING/XAUTOCLAIM 恢复

情形三:需要每个订阅者都收到同一消息

Streams 可以通过多个消费者组实现独立消费:

order-workers     -> 订单处理服务
analytics-workers -> 分析服务
audit-workers     -> 审计服务

每个组都有自己的消费进度。组之间互不共享确认状态。

在同一个消费者组内,多个消费者是竞争关系:一条消息通常只投递给组中的一个消费者。若希望每个服务都得到一份,应该创建多个组,而不是把所有服务放进同一个组。


八、部署边界与诊断方法

1. 单机、主从和集群要分开讨论

在单个主节点上,命令和脚本的原子执行最容易理解。加入复制、哨兵或集群后,还要考虑:

  • 读请求是否误读了副本的旧数据;
  • 故障转移是否丢失尚未复制的写入;
  • Lua 脚本中的键是否位于同一哈希槽;
  • 客户端是否能正确处理 MOVEDASKNOSCRIPT
  • Stream 消费者是否在拓扑变化后重新连接并恢复游标。

不要用副本读取“锁是否存在”来决定是否加锁。锁的判断和写入应在当前主节点上完成。

2. 锁的检查

可以检查:

GET lock:job
PTTL lock:job

需要关注:

  • GET 是否为预期令牌;
  • PTTL 是否为正数;
  • 是否频繁出现接近零的剩余时间;
  • 业务执行时间是否长期超过租约;
  • 是否存在大量锁键因释放脚本未执行而等待过期。

不要用管理命令直接 DEL 业务锁,除非已经确认当前令牌和持有者状态。强制删除可能让原持有者和新持有者同时执行。

3. Lua 的检查

脚本异常时检查:

  • 脚本是否包含错误的键数量;
  • Cluster 中是否跨槽;
  • 是否出现 NOSCRIPT
  • 是否有长循环或一次处理过多成员;
  • 是否使用了可能阻塞的命令;
  • SLOWLOG GET 和延迟监控中是否出现长脚本。

一个逻辑正确但执行几百毫秒的脚本,可能让所有普通请求都排队,因此脚本长度和数据规模必须受控。

4. Pub/Sub 的检查

Pub/Sub 不提供历史查询接口。若怀疑消息丢失,无法通过 Redis 事后列出“某频道过去发布过什么”。应在应用侧记录:

  • 发布日志;
  • 连接建立和断开;
  • 订阅确认;
  • 消费者处理延迟;
  • 关键状态的周期性校验结果。

如果系统需要凭 Redis 本身检查积压、重试或历史消息,Pub/Sub 的数据模型就不合适。

5. Streams 的检查

常用检查命令:

XLEN orders
XINFO STREAM orders
XINFO GROUPS orders
XPENDING orders order-workers
XINFO CONSUMERS orders order-workers

这些命令分别帮助确认:

  • Stream 中有多少条消息;
  • 首尾 ID 和长度等信息;
  • 消费者组的消费位置;
  • PEL 中是否有积压;
  • 哪个消费者持有多少待处理消息、空闲多久。

一个典型故障判断过程是:

Stream 长度上升
    -> 检查组的最后消费 ID
    -> 检查消费者是否仍在线
    -> 检查 PEL 是否持续增长
    -> 检查最老消息的空闲时间
    -> 决定重启、认领、重试或转入死信 Stream

Redis Streams 没有自动为业务定义“死信队列”。若消息重试多次仍失败,应用可以将原消息、错误原因、重试次数写入另一个 Stream,并在确认原消息后进行人工或异步处理。


九、几个容易混淆的结论

“Redis 有原子命令,所以分布式锁绝对安全”

不成立。原子命令解决的是同一 Redis 主节点上的并发插入问题;租约过期、客户端暂停、主从切换和外部资源写入仍需单独处理。

“加了过期时间就不会死锁”

过期时间能处理持锁进程崩溃,但会引入租约过期后的旧持有者问题。任务时间不可预测时,续租和栅栏令牌比单纯增大 TTL 更重要。

“Pub/Sub 发布成功就代表消息被处理”

不成立。PUBLISH 的返回值只是当时的订阅连接数量,不代表业务处理成功,更不代表未来可以重放。

“XACK 后消息就从 Redis 删除了”

不成立。XACK 只更新消费者组的确认状态。Stream 中的消息仍可能存在,直到被裁剪或删除。

“Streams 自动提供恰好一次处理”

不成立。它通常提供至少一次投递,业务必须幂等。数据库唯一键、条件更新、业务状态机和去重表都是常见的幂等基础。

“Lua 脚本等同于跨数据库事务”

不成立。Lua 只原子地修改 Redis 内部状态,不能把 Redis 写入和 MySQL、消息网关或外部服务调用绑定为一个事务。


Redis 的协调能力可以归纳为不同的状态模型:

锁:
    一个键 + 随机令牌 + 租约
    目标是互斥和故障后释放

Lua:
    当前状态 + 参数 -> 原子状态转换
    目标是消除读取与写回之间的竞态

限流:
    计数器或令牌桶状态 + 时间
    目标是限制窗口内数量或长期速率

Pub/Sub:
    当前在线订阅连接
    目标是低延迟广播,不保存历史

Streams:
    追加消息 + 消费者组游标 + PEL
    目标是可追踪、可恢复的消息处理

正确使用它们的关键,不是记住某条命令,而是先确定业务需要的状态是否持久、消息是否允许丢失、失败后是否必须重试,以及旧客户端在租约失效后是否仍可能写入资源。只有这些边界清楚后,锁、脚本、限流、广播和消息流才会各自处在适合的位置。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。

评论

0 条讨论
0/1000
还没有评论,来聊聊你的看法
WR Blog 加载中...
返回文章
数据库Redis分布式锁消息队列

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams封面

数据库基础体系 · 第 39/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 不只是缓存。它还提供了一组建立在内存数据结构、单命令原子性、Lua 脚本和持久化日志之上的协调能力:

  • 用键值状态实现带租约的分布式锁;
  • 用 Lua 把“读取—判断—写回”合并为不可插入的原子操作;
  • 用计数器、时间戳和令牌桶实现限流;
  • 用 Pub/Sub 传递实时但不持久的通知;
  • 用 Streams 保存消息,并通过消费者组实现可恢复的消费进度。

这些能力的可靠性边界并不相同。锁和限流主要依赖 Redis 主节点上操作的原子性;Pub/Sub 依赖在线连接;Streams 则把消息和消费状态保存为 Redis 数据,因此可以处理消费者短暂故障。将它们混为“Redis 消息队列”或“Redis 事务”会导致错误的设计。


一、共同基础:Redis 的原子性到底保证什么

1. Redis 命令的串行执行

对一个 Redis 实例来说,普通命令在主线程中按顺序执行。假设当前键 counter 的值为 10,两个客户端同时执行:

客户端 A:GET counter
客户端 B:GET counter
客户端 A:SET counter 11
客户端 B:SET counter 11

GETSET 各自是原子的,但整个“读取后加一”不是原子的,最终结果可能是 11 而不是期望的 12

使用单条命令:

INCR counter

或者使用 Lua 脚本,可以让一组操作在执行期间不被其他 Redis 命令插入。

这里的“原子”表示:

  1. 脚本或命令执行期间,其他客户端命令不会交错执行;
  2. 其他客户端不会看到脚本执行过程中的中间状态;
  3. 它不表示业务操作具备跨系统事务;
  4. 它也不表示脚本发生错误后所有已经执行的写入会自动回滚。

Redis 的 Lua 脚本应当短小、确定且不执行阻塞操作。脚本执行时间过长会阻塞整个实例上的其他命令。脚本中的 redis.call 出错时,脚本会报告错误;脚本此前已经产生的写入不能依赖“自动回滚”来恢复,因此应在写入前完成参数校验。

2. MULTI/EXEC 与 Lua 的区别

Redis 事务通常指:

MULTI
SET a 1
SET b 2
EXEC

MULTI 之后命令会先排队,EXEC 时连续执行。事务执行期间,其他客户端命令不能插入这批命令。

但 Redis 事务不是关系数据库意义上的完整 ACID 事务:

  • 没有传统意义的回滚机制;
  • WATCH 提供的是乐观并发控制;
  • 它不自动覆盖外部数据库、文件系统或 HTTP 调用;
  • Redis Cluster 中事务涉及的键通常必须位于同一个哈希槽。

Lua 更适合“根据当前值判断后再写入”的逻辑,例如:

读取锁值
判断是否属于当前持有者
删除锁

如果拆成 GETDEL,中间可能被其他客户端插入;如果放入一个脚本,整个判断和删除不可被插入。


二、分布式锁:带租约的互斥,而不是永久所有权

1. 锁的基本模型

分布式锁通常需要满足三个条件:

  1. 互斥性:同一时刻至多一个客户端被认为持有锁;
  2. 释放安全:客户端只能释放自己持有的锁;
  3. 故障可恢复:持锁客户端崩溃后,锁最终可以再次获取。

在 Redis 中,最基本的加锁命令是:

SET lock:report:daily 9f3c... NX PX 10000

参数含义:

  • NX:仅当键不存在时设置;
  • PX 10000:设置 10 秒毫秒级过期时间;
  • 9f3c...:本次加锁请求生成的随机令牌。

成功时返回:

OK

失败时通常返回:

(nil)

不能使用下面这种两步操作代替:

SETNX lock:report:daily 1
EXPIRE lock:report:daily 10

因为客户端可能在两条命令之间崩溃,留下永不过期的锁。

2. 为什么锁值必须是随机令牌

假设客户端 A 获取锁:

lock:job = token-A

A 执行任务超过租约时间,锁自动过期。之后客户端 B 获取:

lock:job = token-B

如果 A 这时直接执行:

DEL lock:job

A 删除的实际上是 B 的锁。

因此释放锁必须是“比较令牌并删除”的原子操作:

-- unlock.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('DEL', KEYS[1])
end
return 0

执行:

EVAL "$(cat unlock.lua)" 1 lock:job token-A

返回值:

  • 1:比较成功并删除;
  • 0:键不存在,或者当前令牌不是 token-A

不能写成:

GET lock:job
如果等于 token-A,再执行 DEL

因为 GETDEL 之间可能发生过期、重新加锁或其他客户端写入。

3. 租约过期后的真实边界

Redis 锁的过期时间使它成为带租约的锁。它只能保证:

在 Redis 认为租约仍有效、且部署故障模型满足假设时,持有者拥有锁。

它不能保证客户端在业务代码中永远拥有锁。例如:

  1. A 获取锁,租约为 10 秒;
  2. A 因为 GC、进程暂停或网络阻塞,20 秒后才继续执行;
  3. Redis 中的锁早已过期;
  4. B 已经重新获取锁并修改资源;
  5. A 恢复后继续写入。

即使 A 后续调用安全解锁脚本,也只能避免删除 B 的锁,不能阻止 A 已经对业务资源执行过期写入。

因此,对于“过期持有者不能覆盖新持有者”的严格要求,需要栅栏令牌(fencing token)

4. 栅栏令牌:把锁状态传递给资源

获取锁时,同时从 Redis 生成单调递增的令牌:

MULTI
INCR lock:job:fence
SET lock:job token-A NX PX 10000
EXEC

实际应用需要处理 SET 失败时不要误用已经生成的令牌。更常见的方式是把“生成令牌”和“抢锁”封装在 Lua 中:

-- acquire-with-fence.lua
if redis.call('EXISTS', KEYS[1]) == 1 then
    return {0, false}
end

local token = redis.call('INCR', KEYS[2])
redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[2])
return {1, token}

调用:

EVAL "$(cat acquire-with-fence.lua)" 2 lock:job lock:job:fence token-A 10000

成功可能返回:

1
1

客户端把令牌 1 传给真正的资源服务。资源服务保存“最后接受的令牌”,只接受更大的令牌:

请求 A:fence=1,接受
请求 B:fence=2,接受
请求 A:fence=1,拒绝

这样,Redis 只负责发放单调令牌,资源本身负责拒绝旧持有者。若资源是数据库,可以通过条件更新实现:

UPDATE job_state
SET result = :result,
    last_fence = :fence
WHERE job_id = :job_id
  AND last_fence < :fence;

这一步必须和资源更新处于同一个数据库事务中,否则检查和更新之间仍可能被插入。

5. 续租与释放

如果任务可能超过初始租约,应由持有者定期续租。续租也必须比较令牌:

-- renew.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('PEXPIRE', KEYS[1], ARGV[2])
end
return 0

执行:

EVAL "$(cat renew.lua)" 1 lock:job token-A 10000

续租周期不能等到租约最后一刻才执行,否则网络延迟或进程暂停可能造成误过期。与此同时,业务代码仍应设计为:续租失败后停止产生新的副作用,而不是继续假设自己拥有锁。

6. 主从复制和故障转移边界

在单个 Redis 主节点上,SET NX PX 的互斥判断是原子的。但如果主节点宕机,副本尚未复制最新锁,故障转移后副本可能认为锁不存在,另一个客户端便可获取同名锁。

因此:

  • Redis 主从复制默认是异步的;
  • WAIT 可以等待写入传播到若干副本,但不等于共识协议;
  • 网络分区、故障转移和旧主恢复时,不能仅凭普通 Redis 锁得到严格的线性一致互斥;
  • 对高价值资源,栅栏令牌、数据库约束或具备共识语义的协调系统仍然重要。

把多个 Redis 节点拼成所谓“更强的锁”也不会自动消除时钟、网络分区和客户端暂停问题。应先明确业务需要的是“尽量避免重复执行”,还是“在故障条件下也绝不允许旧持有者写入”,两者的系统设计不同。


三、Lua:把条件、读取和写入绑定为一个状态转换

1. 脚本的输入边界

Lua 脚本通过两类参数接收数据:

  • KEYS:脚本访问的 Redis 键;
  • ARGV:普通参数。

例如:

local current = redis.call('GET', KEYS[1])
local limit = tonumber(ARGV[1])

if current and tonumber(current) >= limit then
    return 0
end

redis.call('INCR', KEYS[1])
return 1

调用:

EVAL "$(cat increment-if-below.lua)" 1 quota:user:42 10

在 Redis Cluster 中,脚本访问的键必须满足集群路由要求,通常应位于同一个哈希槽。可以使用哈希标签:

quota:{user-42}:count
quota:{user-42}:timestamp

{user-42} 使两个键使用相同的哈希标签。脚本中应把所有键放入 KEYS,不要把键名藏在 ARGV 中。这样 Redis 才能正确检查和路由脚本涉及的键。

2. EVAL、EVALSHA 与脚本缓存

直接执行:

EVAL "return redis.call('INCR', KEYS[1])" 1 counter

也可以先加载脚本:

SCRIPT LOAD "return redis.call('INCR', KEYS[1])"

Redis 返回一个 SHA-1 摘要,例如:

"7d9..."

之后执行:

EVALSHA 7d9... 1 counter

EVALSHA 只引用当前实例脚本缓存中的脚本。重启、故障转移或连接到另一个节点后,可能出现 NOSCRIPT。客户端必须能够重新 SCRIPT LOAD 或回退到 EVAL。脚本缓存不是持久化的业务配置中心。

对于固定部署的脚本,也可以使用 Redis Functions。它们适合把函数随函数库加载和管理,但仍然受到脚本执行原子性、执行时长和集群键约束等边界限制。不能因为使用 Functions 就获得跨服务事务。

3. 脚本中的时间和随机性

涉及时间的脚本必须考虑时间来源:

  • 使用客户端传入的时间,可能受到不同机器时钟偏差影响;
  • 使用 Redis 时间,可减少客户端时钟差异,但脚本仍应保持短小,并遵守当前 Redis 版本对脚本可调用命令和复制的限制;
  • 不应让不同客户端用完全不一致的时间更新同一个限流状态。

对于需要可复制、可重放的脚本,应避免依赖不可控的随机行为和不确定的外部状态。脚本最好把必要的时间、数量和策略参数显式作为输入。


四、限流:从固定窗口到令牌桶

限流的目标是限制某个主体在一个时间范围内可执行的请求数。主体可以是用户、IP、API、租户或资源键。

需要先区分两个量:

  • 吞吐限制:长期平均每秒允许多少请求;
  • 突发容量:短时间内最多允许多少请求。

固定窗口容易实现,但窗口边界可能产生突发;令牌桶能分别表达这两个量。

1. 固定窗口计数器

目标:每个用户每 60 秒最多请求 3 次。

Lua 脚本:

-- fixed-window.lua
local count = redis.call('INCR', KEYS[1])

if count == 1 then
    redis.call('EXPIRE', KEYS[1], ARGV[1])
end

if count <= tonumber(ARGV[2]) then
    return {1, count, tonumber(ARGV[2]) - count}
end

return {0, count, 0}

调用:

EVAL "$(cat fixed-window.lua)" 1 rate:user:42:window 60 3

可能返回:

1
2
1

含义是:

允许,当前窗口计数为 2,剩余 1 次

脚本中只有第一次计数时设置过期时间,因此不会因为后续请求不断刷新 TTL 而形成永不过期的计数器。

但固定窗口存在边界问题。假设窗口为:

[00:00:00, 00:01:00)
[00:01:00, 00:02:00)

用户在 00:00:59 请求 3 次,又在 00:01:00 请求 3 次。两秒内可以通过 6 次请求,虽然每个窗口分别不超过 3 次。这不是实现错误,而是固定窗口算法的定义结果。

2. 滑动窗口日志

滑动窗口要求检查最近 60 秒内的请求数。可以使用 Sorted Set:

  • member:请求唯一 ID;
  • score:请求时间戳。

每次请求执行:

删除 score <= now - 60000 的成员
统计剩余成员
若小于限制则加入当前请求

对应 Lua:

-- sliding-window.lua
local now = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local limit = tonumber(ARGV[3])
local request_id = ARGV[4]

redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', now - window)
local count = redis.call('ZCARD', KEYS[1])

if count >= limit then
    redis.call('PEXPIRE', KEYS[1], window)
    return {0, count}
end

redis.call('ZADD', KEYS[1], now, request_id)
redis.call('PEXPIRE', KEYS[1], window)
return {1, count + 1}

调用示例:

EVAL "$(cat sliding-window.lua)" 1 rate:user:42:log 1700000000000 60000 3 req-001

它比固定窗口精确,但每个请求都需要维护一个成员,内存和删除成本更高。request_id 必须唯一;如果相同 member 重复使用,ZADD 会更新原成员而不是增加计数。

3. 令牌桶

令牌桶用两个参数表达限制:

  • capacity:桶的最大令牌数,表示允许的突发容量;
  • rate:每秒补充的令牌数,表示长期平均速率。

状态为:

(tokens, last_timestamp)

在时间 now 到来时,先补充:

tokens=min(capacity, tokens+(nowlast)×rate)tokens' = \min(capacity,\ tokens + (now-last)\times rate)

其中:

  • tokens 是上次计算后的剩余令牌;
  • last 是上次更新时间;
  • rate 的单位是“令牌/秒”;
  • now-last 必须换算成秒;
  • capacity 防止空闲期间无限积累。

若请求需要 cost 个令牌:

tokens' >= cost  => 允许,并保存 tokens' - cost
tokens' <  cost  => 拒绝,并保存补充后的 tokens'

一个使用毫秒时间的 Lua 实现如下:

-- token-bucket.lua
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])       -- tokens per second
local now_ms = tonumber(ARGV[3])
local cost = tonumber(ARGV[4])
local ttl_ms = tonumber(ARGV[5])

local values = redis.call('HMGET', KEYS[1], 'tokens', 'ts')
local tokens = tonumber(values[1])
local last_ms = tonumber(values[2])

if tokens == nil then
    tokens = capacity
    last_ms = now_ms
end

-- 客户端时钟回拨时,不让令牌数量因负时间差增加异常
if now_ms < last_ms then
    now_ms = last_ms
end

local elapsed_seconds = (now_ms - last_ms) / 1000.0
tokens = math.min(capacity, tokens + elapsed_seconds * rate)

local allowed = 0
local retry_after_ms = 0

if tokens >= cost then
    tokens = tokens - cost
    allowed = 1
else
    local missing = cost - tokens
    retry_after_ms = math.ceil(missing / rate * 1000)
end

redis.call('HSET', KEYS[1], 'tokens', tokens, 'ts', now_ms)
redis.call('PEXPIRE', KEYS[1], ttl_ms)

return {allowed, tokens, retry_after_ms}

调用:

EVAL "$(cat token-bucket.lua)" 1 rate:user:42:bucket 10 2 1700000000000 1 10000

参数含义是:

  • 容量 10;
  • 每秒补充 2 个令牌;
  • 当前时间戳为 1700000000000 毫秒;
  • 本次请求消耗 1 个令牌;
  • 状态空闲 10 秒后过期。

返回:

1
9
0

表示允许请求,剩余约 9 个令牌,不需要等待。

如果返回:

0
0.5
250

表示拒绝,当前只有 0.5 个令牌,理论上约 250 毫秒后才能补足一个令牌。

这个示例把当前时间作为参数传入,因此调用方必须统一时间来源。若不同应用服务器的时钟偏差很大,同一个桶的状态会出现不一致。可以由 Redis 提供时间来源,或者由网关统一生成时间;关键是不要让多个不受控的时钟共同推进同一个限流状态。

令牌桶脚本中的 rate 不能为 0,否则拒绝请求时计算等待时间会发生除零。生产实现还应限制参数范围,防止客户端传入极大的容量、速率或 TTL。

4. 限流失败路径

限流通常有三种结果:

  1. Redis 判断允许,业务请求继续;
  2. Redis 判断拒绝,返回 HTTP 429 Too Many Requests
  3. Redis 不可用或超时。

第三种不能简单等同于“限流通过”。若放行,可能在 Redis 故障时失去保护;若拒绝,可能把 Redis 故障扩大为全站不可用。应根据接口风险决定:

  • 登录、验证码、支付等安全敏感接口偏向失败关闭;
  • 非关键读接口可能使用本地兜底限流;
  • 本地兜底必须明确它只是每实例限流,不是全局限流。

五、Pub/Sub:在线广播,不是可靠队列

1. 基本数据流

Redis Pub/Sub 有两个角色:

  • 发布者:向频道发送消息;
  • 订阅者:通过长连接订阅频道并接收消息。

订阅:

SUBSCRIBE cache:product:updated

成功后 Redis 会返回订阅确认,例如:

subscribe
cache:product:updated
1

发布:

PUBLISH cache:product:updated '{"product_id":42}'

返回值是当时收到该消息的订阅客户端数量,例如:

(integer) 2

这个返回值不是消费成功数,也不是持久化成功数。它只表示消息发布时有多少订阅连接被送达。

发布数据流可以表示为:

发布者
  |
  | PUBLISH
  v
Redis 当前连接集合
  |
  +--> 订阅者 A
  +--> 订阅者 B

Redis 不把 Pub/Sub 消息作为普通键保存。订阅者掉线期间发布的消息不会等待它重新连接后补发。

2. 可靠性语义

Redis Pub/Sub 通常应理解为:

在线连接上的即时、尽力发送、至多一次交付通知。

常见失败路径:

  1. 发布者发布时没有订阅者:消息直接消失;
  2. 订阅者网络断开:断线期间的消息不会补发;
  3. 订阅者收到消息后进程崩溃:没有 Redis 侧确认和重投;
  4. 消费者处理很慢:消息可能在客户端连接缓冲区积压,最终导致连接问题。

因此,以下场景适合 Pub/Sub:

  • 缓存失效通知;
  • 配置刷新提示;
  • 在线 WebSocket 广播;
  • 不要求离线补偿的状态变化通知。

以下场景不应只使用 Pub/Sub:

  • 订单、支付、库存等不可丢失事件;
  • 必须重新消费的任务;
  • 需要确认、重试和消费进度的消息。

如果一个消费者错过通知后可以通过读取当前状态自行修复,Pub/Sub 可以作为低延迟通知通道。例如收到“商品更新”后删除本地缓存;即使通知丢失,下一次 TTL 到期或主动校验仍可修复。若消息本身就是唯一业务事实,则应使用 Streams 或其他持久化消息系统。

3. 频道订阅与模式订阅

普通订阅:

SUBSCRIBE tenant:42:events

模式订阅:

PSUBSCRIBE tenant:*:events

模式订阅更灵活,但匹配和消息分发成本更高。频道命名应包含明确的租户、业务和版本边界,避免所有业务共享一个频道后由客户端自行过滤。

Redis 还提供按分片路由的 Sharded Pub/Sub 命令,适用于 Redis Cluster 中希望减少跨节点广播的场景。它与普通 Pub/Sub 的连接、拓扑和客户端支持要求不同,不能假设所有客户端都自动兼容;使用前应确认客户端对对应命令和集群重定向的支持。


六、Streams:可持久化的追加日志与消费者组

1. Stream 是什么

Redis Stream 是一个按 ID 排序的追加记录结构。每条记录包含:

消息 ID -> field/value 字段

写入一条消息:

XADD orders * order_id 1001 status created

返回类似:

"1700000000000-0"

* 表示由 Redis 生成 ID。ID 通常包含毫秒时间部分和序列部分,不能把它当作业务时间或全局业务排序的替代品。消费者应保存和比较 Stream ID,而不是仅用时间戳。

读取:

XRANGE orders - +

可能返回:

1) 1) "1700000000000-0"
   2) 1) "order_id"
      2) "1001"
      3) "status"
      4) "created"

生产环境通常需要限制 Stream 长度,例如:

XADD orders MAXLEN ~ 100000 * order_id 1001 status created

~ 表示近似裁剪,Redis 可以用更高效的方式接近目标长度;它不是严格保证长度永远不超过 100000。若必须精确限制,可以去掉 ~,但裁剪成本和写入性能边界不同。

2. 普通消费者读取

从头读取:

XREAD COUNT 10 STREAMS orders 0-0

从某个 ID 之后读取:

XREAD COUNT 10 STREAMS orders 1700000000000-0

阻塞读取:

XREAD BLOCK 5000 COUNT 10 STREAMS orders $

这里的 $ 表示只关注执行命令时之后新增的消息。它适合实时监听,但如果客户端启动或重连时使用 $,命令执行前已经写入的消息不会被读取。可靠消费者应保存上次成功处理的 ID,并从该 ID 继续读取,而不是每次重连都使用 $

3. 消费者组和消费确认

消费者组让多个消费者共同处理一个 Stream,并记录组内消费位置。

创建组:

XGROUP CREATE orders order-workers 0-0 MKSTREAM

含义:

  • Stream:orders
  • 消费者组:order-workers
  • 0-0 开始;
  • 如果 Stream 不存在则创建。

消费者以名字加入:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 BLOCK 5000 STREAMS orders >

> 表示读取当前组尚未投递给任何消费者的新消息。

Redis 返回的消息会进入该消费者组的待处理列表,称为 PEL(Pending Entries List)。此时消息已经投递,但还没有确认。

业务处理成功后确认:

XACK orders order-workers 1700000000000-0

XACK 只表示该消费者组不再把该消息视为待确认消息,不会删除 Stream 中的原始消息。消息是否仍能被其他读取方式看到,取决于 Stream 的保留和裁剪策略。

4. 失败、重试和消息转移

考虑下面的过程:

1. Redis 将消息投递给 worker-1
2. worker-1 执行业务操作
3. worker-1 在 XACK 前崩溃

消息仍然在该组的 PEL 中,不会因为消费者进程消失而自动丢失。其他消费者可以检查:

XPENDING orders order-workers

它会报告待处理消息数量、最小和最大 ID,以及按消费者统计的数量。

较新的 Redis 稳定语义中,可以用 XAUTOCLAIM 把长时间未处理的消息转移给当前消费者:

XAUTOCLAIM orders order-workers worker-2 60000 0-0 COUNT 10

这里的 60000 是最小空闲时间,单位为毫秒。只有空闲超过该时间的待处理消息才适合被认领。认领前应考虑 worker-1 只是暂时变慢,避免多个消费者频繁争抢同一消息。

另一种方式是:

XCLAIM orders order-workers worker-2 60000 1700000000000-0

它适合已经明确知道要转移哪些 ID 的场景。

5. Streams 通常是至少一次,而不是恰好一次

Streams 消费者组的典型语义是至少一次投递

投递消息
执行外部业务
消费者在确认前崩溃
消息被重新认领
再次执行外部业务

因此必须允许重复处理。常见做法是使用业务唯一 ID 和幂等约束:

CREATE TABLE processed_event (
    event_id VARCHAR(100) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

处理时,在同一个数据库事务中:

1. 尝试插入 event_id
2. 若主键冲突,说明已经处理过,跳过业务副作用
3. 若插入成功,执行业务更新
4. 数据库事务提交
5. 再执行 XACK

若第 3 步成功但第 5 步前进程崩溃,消息会再次投递;第 1 步的唯一约束会使第二次处理变成安全跳过。

但如果先 XACK 再更新数据库,进程在两者之间崩溃,就可能出现“消息已确认、业务未完成”。所以确认位置不是越早越好,而是应放在业务状态已经可靠提交之后。

XREADGROUPNOACK 选项可以减少 PEL 记录:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 NOACK STREAMS orders >

但它也意味着消息投递后不再等待确认,消费者崩溃时不能依赖 PEL 恢复。它只适合允许丢失或已具备其他可靠来源的场景。

6. Stream 裁剪与 PEL 的关系

裁剪 Stream 不等于自动完成消费确认。若消息从 Stream 主体中被裁剪,消费者组的待处理状态仍可能需要单独管理,具体行为取决于所使用的裁剪和 Redis 版本语义。生产系统需要同时监控:

  • Stream 长度;
  • 消费者组的组内最后 ID;
  • PEL 数量;
  • 最老待处理消息的空闲时间;
  • 认领和重试次数。

仅仅看到 Stream 长度正常,并不能说明消费者没有积压。


七、Pub/Sub 与 Streams 的选择推导

可以从“消息事实是否需要在未来重新获得”开始判断。

情形一:消息只是状态变化提醒

例如:

数据库商品价格发生变化

消费者即使错过事件,也可以重新读取数据库当前价格。因此可以使用:

数据库作为事实来源
Redis Pub/Sub 作为低延迟通知
缓存 TTL 或定期校验作为兜底

这与缓存一致性中的“失效通知”相匹配:通知丢失不会永久丢失业务事实。

情形二:消息本身就是业务事实

例如:

订单已创建
支付成功
库存扣减任务

消费者不能只看当前状态推导出所有历史动作,且必须重试和追踪进度。这时需要 Streams:

XADD 写入事实
XREADGROUP 投递
业务事务提交
XACK 确认
XPENDING/XAUTOCLAIM 恢复

情形三:需要每个订阅者都收到同一消息

Streams 可以通过多个消费者组实现独立消费:

order-workers     -> 订单处理服务
analytics-workers -> 分析服务
audit-workers     -> 审计服务

每个组都有自己的消费进度。组之间互不共享确认状态。

在同一个消费者组内,多个消费者是竞争关系:一条消息通常只投递给组中的一个消费者。若希望每个服务都得到一份,应该创建多个组,而不是把所有服务放进同一个组。


八、部署边界与诊断方法

1. 单机、主从和集群要分开讨论

在单个主节点上,命令和脚本的原子执行最容易理解。加入复制、哨兵或集群后,还要考虑:

  • 读请求是否误读了副本的旧数据;
  • 故障转移是否丢失尚未复制的写入;
  • Lua 脚本中的键是否位于同一哈希槽;
  • 客户端是否能正确处理 MOVEDASKNOSCRIPT
  • Stream 消费者是否在拓扑变化后重新连接并恢复游标。

不要用副本读取“锁是否存在”来决定是否加锁。锁的判断和写入应在当前主节点上完成。

2. 锁的检查

可以检查:

GET lock:job
PTTL lock:job

需要关注:

  • GET 是否为预期令牌;
  • PTTL 是否为正数;
  • 是否频繁出现接近零的剩余时间;
  • 业务执行时间是否长期超过租约;
  • 是否存在大量锁键因释放脚本未执行而等待过期。

不要用管理命令直接 DEL 业务锁,除非已经确认当前令牌和持有者状态。强制删除可能让原持有者和新持有者同时执行。

3. Lua 的检查

脚本异常时检查:

  • 脚本是否包含错误的键数量;
  • Cluster 中是否跨槽;
  • 是否出现 NOSCRIPT
  • 是否有长循环或一次处理过多成员;
  • 是否使用了可能阻塞的命令;
  • SLOWLOG GET 和延迟监控中是否出现长脚本。

一个逻辑正确但执行几百毫秒的脚本,可能让所有普通请求都排队,因此脚本长度和数据规模必须受控。

4. Pub/Sub 的检查

Pub/Sub 不提供历史查询接口。若怀疑消息丢失,无法通过 Redis 事后列出“某频道过去发布过什么”。应在应用侧记录:

  • 发布日志;
  • 连接建立和断开;
  • 订阅确认;
  • 消费者处理延迟;
  • 关键状态的周期性校验结果。

如果系统需要凭 Redis 本身检查积压、重试或历史消息,Pub/Sub 的数据模型就不合适。

5. Streams 的检查

常用检查命令:

XLEN orders
XINFO STREAM orders
XINFO GROUPS orders
XPENDING orders order-workers
XINFO CONSUMERS orders order-workers

这些命令分别帮助确认:

  • Stream 中有多少条消息;
  • 首尾 ID 和长度等信息;
  • 消费者组的消费位置;
  • PEL 中是否有积压;
  • 哪个消费者持有多少待处理消息、空闲多久。

一个典型故障判断过程是:

Stream 长度上升
    -> 检查组的最后消费 ID
    -> 检查消费者是否仍在线
    -> 检查 PEL 是否持续增长
    -> 检查最老消息的空闲时间
    -> 决定重启、认领、重试或转入死信 Stream

Redis Streams 没有自动为业务定义“死信队列”。若消息重试多次仍失败,应用可以将原消息、错误原因、重试次数写入另一个 Stream,并在确认原消息后进行人工或异步处理。


九、几个容易混淆的结论

“Redis 有原子命令,所以分布式锁绝对安全”

不成立。原子命令解决的是同一 Redis 主节点上的并发插入问题;租约过期、客户端暂停、主从切换和外部资源写入仍需单独处理。

“加了过期时间就不会死锁”

过期时间能处理持锁进程崩溃,但会引入租约过期后的旧持有者问题。任务时间不可预测时,续租和栅栏令牌比单纯增大 TTL 更重要。

“Pub/Sub 发布成功就代表消息被处理”

不成立。PUBLISH 的返回值只是当时的订阅连接数量,不代表业务处理成功,更不代表未来可以重放。

“XACK 后消息就从 Redis 删除了”

不成立。XACK 只更新消费者组的确认状态。Stream 中的消息仍可能存在,直到被裁剪或删除。

“Streams 自动提供恰好一次处理”

不成立。它通常提供至少一次投递,业务必须幂等。数据库唯一键、条件更新、业务状态机和去重表都是常见的幂等基础。

“Lua 脚本等同于跨数据库事务”

不成立。Lua 只原子地修改 Redis 内部状态,不能把 Redis 写入和 MySQL、消息网关或外部服务调用绑定为一个事务。


Redis 的协调能力可以归纳为不同的状态模型:

锁:
    一个键 + 随机令牌 + 租约
    目标是互斥和故障后释放

Lua:
    当前状态 + 参数 -> 原子状态转换
    目标是消除读取与写回之间的竞态

限流:
    计数器或令牌桶状态 + 时间
    目标是限制窗口内数量或长期速率

Pub/Sub:
    当前在线订阅连接
    目标是低延迟广播,不保存历史

Streams:
    追加消息 + 消费者组游标 + PEL
    目标是可追踪、可恢复的消息处理

正确使用它们的关键,不是记住某条命令,而是先确定业务需要的状态是否持久、消息是否允许丢失、失败后是否必须重试,以及旧客户端在租约失效后是否仍可能写入资源。只有这些边界清楚后,锁、脚本、限流、广播和消息流才会各自处在适合的位置。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。

评论

0 条讨论
0/1000
还没有评论,来聊聊你的看法
表示只关注执行命令时之后新增的消息。它适合实时监听,但如果客户端启动或重连时使用 ` Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams - WR Blog
WR Blog 加载中...
返回文章
数据库Redis分布式锁消息队列

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams封面

数据库基础体系 · 第 39/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 不只是缓存。它还提供了一组建立在内存数据结构、单命令原子性、Lua 脚本和持久化日志之上的协调能力:

  • 用键值状态实现带租约的分布式锁;
  • 用 Lua 把“读取—判断—写回”合并为不可插入的原子操作;
  • 用计数器、时间戳和令牌桶实现限流;
  • 用 Pub/Sub 传递实时但不持久的通知;
  • 用 Streams 保存消息,并通过消费者组实现可恢复的消费进度。

这些能力的可靠性边界并不相同。锁和限流主要依赖 Redis 主节点上操作的原子性;Pub/Sub 依赖在线连接;Streams 则把消息和消费状态保存为 Redis 数据,因此可以处理消费者短暂故障。将它们混为“Redis 消息队列”或“Redis 事务”会导致错误的设计。


一、共同基础:Redis 的原子性到底保证什么

1. Redis 命令的串行执行

对一个 Redis 实例来说,普通命令在主线程中按顺序执行。假设当前键 counter 的值为 10,两个客户端同时执行:

客户端 A:GET counter
客户端 B:GET counter
客户端 A:SET counter 11
客户端 B:SET counter 11

GETSET 各自是原子的,但整个“读取后加一”不是原子的,最终结果可能是 11 而不是期望的 12

使用单条命令:

INCR counter

或者使用 Lua 脚本,可以让一组操作在执行期间不被其他 Redis 命令插入。

这里的“原子”表示:

  1. 脚本或命令执行期间,其他客户端命令不会交错执行;
  2. 其他客户端不会看到脚本执行过程中的中间状态;
  3. 它不表示业务操作具备跨系统事务;
  4. 它也不表示脚本发生错误后所有已经执行的写入会自动回滚。

Redis 的 Lua 脚本应当短小、确定且不执行阻塞操作。脚本执行时间过长会阻塞整个实例上的其他命令。脚本中的 redis.call 出错时,脚本会报告错误;脚本此前已经产生的写入不能依赖“自动回滚”来恢复,因此应在写入前完成参数校验。

2. MULTI/EXEC 与 Lua 的区别

Redis 事务通常指:

MULTI
SET a 1
SET b 2
EXEC

MULTI 之后命令会先排队,EXEC 时连续执行。事务执行期间,其他客户端命令不能插入这批命令。

但 Redis 事务不是关系数据库意义上的完整 ACID 事务:

  • 没有传统意义的回滚机制;
  • WATCH 提供的是乐观并发控制;
  • 它不自动覆盖外部数据库、文件系统或 HTTP 调用;
  • Redis Cluster 中事务涉及的键通常必须位于同一个哈希槽。

Lua 更适合“根据当前值判断后再写入”的逻辑,例如:

读取锁值
判断是否属于当前持有者
删除锁

如果拆成 GETDEL,中间可能被其他客户端插入;如果放入一个脚本,整个判断和删除不可被插入。


二、分布式锁:带租约的互斥,而不是永久所有权

1. 锁的基本模型

分布式锁通常需要满足三个条件:

  1. 互斥性:同一时刻至多一个客户端被认为持有锁;
  2. 释放安全:客户端只能释放自己持有的锁;
  3. 故障可恢复:持锁客户端崩溃后,锁最终可以再次获取。

在 Redis 中,最基本的加锁命令是:

SET lock:report:daily 9f3c... NX PX 10000

参数含义:

  • NX:仅当键不存在时设置;
  • PX 10000:设置 10 秒毫秒级过期时间;
  • 9f3c...:本次加锁请求生成的随机令牌。

成功时返回:

OK

失败时通常返回:

(nil)

不能使用下面这种两步操作代替:

SETNX lock:report:daily 1
EXPIRE lock:report:daily 10

因为客户端可能在两条命令之间崩溃,留下永不过期的锁。

2. 为什么锁值必须是随机令牌

假设客户端 A 获取锁:

lock:job = token-A

A 执行任务超过租约时间,锁自动过期。之后客户端 B 获取:

lock:job = token-B

如果 A 这时直接执行:

DEL lock:job

A 删除的实际上是 B 的锁。

因此释放锁必须是“比较令牌并删除”的原子操作:

-- unlock.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('DEL', KEYS[1])
end
return 0

执行:

EVAL "$(cat unlock.lua)" 1 lock:job token-A

返回值:

  • 1:比较成功并删除;
  • 0:键不存在,或者当前令牌不是 token-A

不能写成:

GET lock:job
如果等于 token-A,再执行 DEL

因为 GETDEL 之间可能发生过期、重新加锁或其他客户端写入。

3. 租约过期后的真实边界

Redis 锁的过期时间使它成为带租约的锁。它只能保证:

在 Redis 认为租约仍有效、且部署故障模型满足假设时,持有者拥有锁。

它不能保证客户端在业务代码中永远拥有锁。例如:

  1. A 获取锁,租约为 10 秒;
  2. A 因为 GC、进程暂停或网络阻塞,20 秒后才继续执行;
  3. Redis 中的锁早已过期;
  4. B 已经重新获取锁并修改资源;
  5. A 恢复后继续写入。

即使 A 后续调用安全解锁脚本,也只能避免删除 B 的锁,不能阻止 A 已经对业务资源执行过期写入。

因此,对于“过期持有者不能覆盖新持有者”的严格要求,需要栅栏令牌(fencing token)

4. 栅栏令牌:把锁状态传递给资源

获取锁时,同时从 Redis 生成单调递增的令牌:

MULTI
INCR lock:job:fence
SET lock:job token-A NX PX 10000
EXEC

实际应用需要处理 SET 失败时不要误用已经生成的令牌。更常见的方式是把“生成令牌”和“抢锁”封装在 Lua 中:

-- acquire-with-fence.lua
if redis.call('EXISTS', KEYS[1]) == 1 then
    return {0, false}
end

local token = redis.call('INCR', KEYS[2])
redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[2])
return {1, token}

调用:

EVAL "$(cat acquire-with-fence.lua)" 2 lock:job lock:job:fence token-A 10000

成功可能返回:

1
1

客户端把令牌 1 传给真正的资源服务。资源服务保存“最后接受的令牌”,只接受更大的令牌:

请求 A:fence=1,接受
请求 B:fence=2,接受
请求 A:fence=1,拒绝

这样,Redis 只负责发放单调令牌,资源本身负责拒绝旧持有者。若资源是数据库,可以通过条件更新实现:

UPDATE job_state
SET result = :result,
    last_fence = :fence
WHERE job_id = :job_id
  AND last_fence < :fence;

这一步必须和资源更新处于同一个数据库事务中,否则检查和更新之间仍可能被插入。

5. 续租与释放

如果任务可能超过初始租约,应由持有者定期续租。续租也必须比较令牌:

-- renew.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('PEXPIRE', KEYS[1], ARGV[2])
end
return 0

执行:

EVAL "$(cat renew.lua)" 1 lock:job token-A 10000

续租周期不能等到租约最后一刻才执行,否则网络延迟或进程暂停可能造成误过期。与此同时,业务代码仍应设计为:续租失败后停止产生新的副作用,而不是继续假设自己拥有锁。

6. 主从复制和故障转移边界

在单个 Redis 主节点上,SET NX PX 的互斥判断是原子的。但如果主节点宕机,副本尚未复制最新锁,故障转移后副本可能认为锁不存在,另一个客户端便可获取同名锁。

因此:

  • Redis 主从复制默认是异步的;
  • WAIT 可以等待写入传播到若干副本,但不等于共识协议;
  • 网络分区、故障转移和旧主恢复时,不能仅凭普通 Redis 锁得到严格的线性一致互斥;
  • 对高价值资源,栅栏令牌、数据库约束或具备共识语义的协调系统仍然重要。

把多个 Redis 节点拼成所谓“更强的锁”也不会自动消除时钟、网络分区和客户端暂停问题。应先明确业务需要的是“尽量避免重复执行”,还是“在故障条件下也绝不允许旧持有者写入”,两者的系统设计不同。


三、Lua:把条件、读取和写入绑定为一个状态转换

1. 脚本的输入边界

Lua 脚本通过两类参数接收数据:

  • KEYS:脚本访问的 Redis 键;
  • ARGV:普通参数。

例如:

local current = redis.call('GET', KEYS[1])
local limit = tonumber(ARGV[1])

if current and tonumber(current) >= limit then
    return 0
end

redis.call('INCR', KEYS[1])
return 1

调用:

EVAL "$(cat increment-if-below.lua)" 1 quota:user:42 10

在 Redis Cluster 中,脚本访问的键必须满足集群路由要求,通常应位于同一个哈希槽。可以使用哈希标签:

quota:{user-42}:count
quota:{user-42}:timestamp

{user-42} 使两个键使用相同的哈希标签。脚本中应把所有键放入 KEYS,不要把键名藏在 ARGV 中。这样 Redis 才能正确检查和路由脚本涉及的键。

2. EVAL、EVALSHA 与脚本缓存

直接执行:

EVAL "return redis.call('INCR', KEYS[1])" 1 counter

也可以先加载脚本:

SCRIPT LOAD "return redis.call('INCR', KEYS[1])"

Redis 返回一个 SHA-1 摘要,例如:

"7d9..."

之后执行:

EVALSHA 7d9... 1 counter

EVALSHA 只引用当前实例脚本缓存中的脚本。重启、故障转移或连接到另一个节点后,可能出现 NOSCRIPT。客户端必须能够重新 SCRIPT LOAD 或回退到 EVAL。脚本缓存不是持久化的业务配置中心。

对于固定部署的脚本,也可以使用 Redis Functions。它们适合把函数随函数库加载和管理,但仍然受到脚本执行原子性、执行时长和集群键约束等边界限制。不能因为使用 Functions 就获得跨服务事务。

3. 脚本中的时间和随机性

涉及时间的脚本必须考虑时间来源:

  • 使用客户端传入的时间,可能受到不同机器时钟偏差影响;
  • 使用 Redis 时间,可减少客户端时钟差异,但脚本仍应保持短小,并遵守当前 Redis 版本对脚本可调用命令和复制的限制;
  • 不应让不同客户端用完全不一致的时间更新同一个限流状态。

对于需要可复制、可重放的脚本,应避免依赖不可控的随机行为和不确定的外部状态。脚本最好把必要的时间、数量和策略参数显式作为输入。


四、限流:从固定窗口到令牌桶

限流的目标是限制某个主体在一个时间范围内可执行的请求数。主体可以是用户、IP、API、租户或资源键。

需要先区分两个量:

  • 吞吐限制:长期平均每秒允许多少请求;
  • 突发容量:短时间内最多允许多少请求。

固定窗口容易实现,但窗口边界可能产生突发;令牌桶能分别表达这两个量。

1. 固定窗口计数器

目标:每个用户每 60 秒最多请求 3 次。

Lua 脚本:

-- fixed-window.lua
local count = redis.call('INCR', KEYS[1])

if count == 1 then
    redis.call('EXPIRE', KEYS[1], ARGV[1])
end

if count <= tonumber(ARGV[2]) then
    return {1, count, tonumber(ARGV[2]) - count}
end

return {0, count, 0}

调用:

EVAL "$(cat fixed-window.lua)" 1 rate:user:42:window 60 3

可能返回:

1
2
1

含义是:

允许,当前窗口计数为 2,剩余 1 次

脚本中只有第一次计数时设置过期时间,因此不会因为后续请求不断刷新 TTL 而形成永不过期的计数器。

但固定窗口存在边界问题。假设窗口为:

[00:00:00, 00:01:00)
[00:01:00, 00:02:00)

用户在 00:00:59 请求 3 次,又在 00:01:00 请求 3 次。两秒内可以通过 6 次请求,虽然每个窗口分别不超过 3 次。这不是实现错误,而是固定窗口算法的定义结果。

2. 滑动窗口日志

滑动窗口要求检查最近 60 秒内的请求数。可以使用 Sorted Set:

  • member:请求唯一 ID;
  • score:请求时间戳。

每次请求执行:

删除 score <= now - 60000 的成员
统计剩余成员
若小于限制则加入当前请求

对应 Lua:

-- sliding-window.lua
local now = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local limit = tonumber(ARGV[3])
local request_id = ARGV[4]

redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', now - window)
local count = redis.call('ZCARD', KEYS[1])

if count >= limit then
    redis.call('PEXPIRE', KEYS[1], window)
    return {0, count}
end

redis.call('ZADD', KEYS[1], now, request_id)
redis.call('PEXPIRE', KEYS[1], window)
return {1, count + 1}

调用示例:

EVAL "$(cat sliding-window.lua)" 1 rate:user:42:log 1700000000000 60000 3 req-001

它比固定窗口精确,但每个请求都需要维护一个成员,内存和删除成本更高。request_id 必须唯一;如果相同 member 重复使用,ZADD 会更新原成员而不是增加计数。

3. 令牌桶

令牌桶用两个参数表达限制:

  • capacity:桶的最大令牌数,表示允许的突发容量;
  • rate:每秒补充的令牌数,表示长期平均速率。

状态为:

(tokens, last_timestamp)

在时间 now 到来时,先补充:

tokens=min(capacity, tokens+(nowlast)×rate)tokens' = \min(capacity,\ tokens + (now-last)\times rate)

其中:

  • tokens 是上次计算后的剩余令牌;
  • last 是上次更新时间;
  • rate 的单位是“令牌/秒”;
  • now-last 必须换算成秒;
  • capacity 防止空闲期间无限积累。

若请求需要 cost 个令牌:

tokens' >= cost  => 允许,并保存 tokens' - cost
tokens' <  cost  => 拒绝,并保存补充后的 tokens'

一个使用毫秒时间的 Lua 实现如下:

-- token-bucket.lua
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])       -- tokens per second
local now_ms = tonumber(ARGV[3])
local cost = tonumber(ARGV[4])
local ttl_ms = tonumber(ARGV[5])

local values = redis.call('HMGET', KEYS[1], 'tokens', 'ts')
local tokens = tonumber(values[1])
local last_ms = tonumber(values[2])

if tokens == nil then
    tokens = capacity
    last_ms = now_ms
end

-- 客户端时钟回拨时,不让令牌数量因负时间差增加异常
if now_ms < last_ms then
    now_ms = last_ms
end

local elapsed_seconds = (now_ms - last_ms) / 1000.0
tokens = math.min(capacity, tokens + elapsed_seconds * rate)

local allowed = 0
local retry_after_ms = 0

if tokens >= cost then
    tokens = tokens - cost
    allowed = 1
else
    local missing = cost - tokens
    retry_after_ms = math.ceil(missing / rate * 1000)
end

redis.call('HSET', KEYS[1], 'tokens', tokens, 'ts', now_ms)
redis.call('PEXPIRE', KEYS[1], ttl_ms)

return {allowed, tokens, retry_after_ms}

调用:

EVAL "$(cat token-bucket.lua)" 1 rate:user:42:bucket 10 2 1700000000000 1 10000

参数含义是:

  • 容量 10;
  • 每秒补充 2 个令牌;
  • 当前时间戳为 1700000000000 毫秒;
  • 本次请求消耗 1 个令牌;
  • 状态空闲 10 秒后过期。

返回:

1
9
0

表示允许请求,剩余约 9 个令牌,不需要等待。

如果返回:

0
0.5
250

表示拒绝,当前只有 0.5 个令牌,理论上约 250 毫秒后才能补足一个令牌。

这个示例把当前时间作为参数传入,因此调用方必须统一时间来源。若不同应用服务器的时钟偏差很大,同一个桶的状态会出现不一致。可以由 Redis 提供时间来源,或者由网关统一生成时间;关键是不要让多个不受控的时钟共同推进同一个限流状态。

令牌桶脚本中的 rate 不能为 0,否则拒绝请求时计算等待时间会发生除零。生产实现还应限制参数范围,防止客户端传入极大的容量、速率或 TTL。

4. 限流失败路径

限流通常有三种结果:

  1. Redis 判断允许,业务请求继续;
  2. Redis 判断拒绝,返回 HTTP 429 Too Many Requests
  3. Redis 不可用或超时。

第三种不能简单等同于“限流通过”。若放行,可能在 Redis 故障时失去保护;若拒绝,可能把 Redis 故障扩大为全站不可用。应根据接口风险决定:

  • 登录、验证码、支付等安全敏感接口偏向失败关闭;
  • 非关键读接口可能使用本地兜底限流;
  • 本地兜底必须明确它只是每实例限流,不是全局限流。

五、Pub/Sub:在线广播,不是可靠队列

1. 基本数据流

Redis Pub/Sub 有两个角色:

  • 发布者:向频道发送消息;
  • 订阅者:通过长连接订阅频道并接收消息。

订阅:

SUBSCRIBE cache:product:updated

成功后 Redis 会返回订阅确认,例如:

subscribe
cache:product:updated
1

发布:

PUBLISH cache:product:updated '{"product_id":42}'

返回值是当时收到该消息的订阅客户端数量,例如:

(integer) 2

这个返回值不是消费成功数,也不是持久化成功数。它只表示消息发布时有多少订阅连接被送达。

发布数据流可以表示为:

发布者
  |
  | PUBLISH
  v
Redis 当前连接集合
  |
  +--> 订阅者 A
  +--> 订阅者 B

Redis 不把 Pub/Sub 消息作为普通键保存。订阅者掉线期间发布的消息不会等待它重新连接后补发。

2. 可靠性语义

Redis Pub/Sub 通常应理解为:

在线连接上的即时、尽力发送、至多一次交付通知。

常见失败路径:

  1. 发布者发布时没有订阅者:消息直接消失;
  2. 订阅者网络断开:断线期间的消息不会补发;
  3. 订阅者收到消息后进程崩溃:没有 Redis 侧确认和重投;
  4. 消费者处理很慢:消息可能在客户端连接缓冲区积压,最终导致连接问题。

因此,以下场景适合 Pub/Sub:

  • 缓存失效通知;
  • 配置刷新提示;
  • 在线 WebSocket 广播;
  • 不要求离线补偿的状态变化通知。

以下场景不应只使用 Pub/Sub:

  • 订单、支付、库存等不可丢失事件;
  • 必须重新消费的任务;
  • 需要确认、重试和消费进度的消息。

如果一个消费者错过通知后可以通过读取当前状态自行修复,Pub/Sub 可以作为低延迟通知通道。例如收到“商品更新”后删除本地缓存;即使通知丢失,下一次 TTL 到期或主动校验仍可修复。若消息本身就是唯一业务事实,则应使用 Streams 或其他持久化消息系统。

3. 频道订阅与模式订阅

普通订阅:

SUBSCRIBE tenant:42:events

模式订阅:

PSUBSCRIBE tenant:*:events

模式订阅更灵活,但匹配和消息分发成本更高。频道命名应包含明确的租户、业务和版本边界,避免所有业务共享一个频道后由客户端自行过滤。

Redis 还提供按分片路由的 Sharded Pub/Sub 命令,适用于 Redis Cluster 中希望减少跨节点广播的场景。它与普通 Pub/Sub 的连接、拓扑和客户端支持要求不同,不能假设所有客户端都自动兼容;使用前应确认客户端对对应命令和集群重定向的支持。


六、Streams:可持久化的追加日志与消费者组

1. Stream 是什么

Redis Stream 是一个按 ID 排序的追加记录结构。每条记录包含:

消息 ID -> field/value 字段

写入一条消息:

XADD orders * order_id 1001 status created

返回类似:

"1700000000000-0"

* 表示由 Redis 生成 ID。ID 通常包含毫秒时间部分和序列部分,不能把它当作业务时间或全局业务排序的替代品。消费者应保存和比较 Stream ID,而不是仅用时间戳。

读取:

XRANGE orders - +

可能返回:

1) 1) "1700000000000-0"
   2) 1) "order_id"
      2) "1001"
      3) "status"
      4) "created"

生产环境通常需要限制 Stream 长度,例如:

XADD orders MAXLEN ~ 100000 * order_id 1001 status created

~ 表示近似裁剪,Redis 可以用更高效的方式接近目标长度;它不是严格保证长度永远不超过 100000。若必须精确限制,可以去掉 ~,但裁剪成本和写入性能边界不同。

2. 普通消费者读取

从头读取:

XREAD COUNT 10 STREAMS orders 0-0

从某个 ID 之后读取:

XREAD COUNT 10 STREAMS orders 1700000000000-0

阻塞读取:

XREAD BLOCK 5000 COUNT 10 STREAMS orders $

这里的 $ 表示只关注执行命令时之后新增的消息。它适合实时监听,但如果客户端启动或重连时使用 $,命令执行前已经写入的消息不会被读取。可靠消费者应保存上次成功处理的 ID,并从该 ID 继续读取,而不是每次重连都使用 $

3. 消费者组和消费确认

消费者组让多个消费者共同处理一个 Stream,并记录组内消费位置。

创建组:

XGROUP CREATE orders order-workers 0-0 MKSTREAM

含义:

  • Stream:orders
  • 消费者组:order-workers
  • 0-0 开始;
  • 如果 Stream 不存在则创建。

消费者以名字加入:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 BLOCK 5000 STREAMS orders >

> 表示读取当前组尚未投递给任何消费者的新消息。

Redis 返回的消息会进入该消费者组的待处理列表,称为 PEL(Pending Entries List)。此时消息已经投递,但还没有确认。

业务处理成功后确认:

XACK orders order-workers 1700000000000-0

XACK 只表示该消费者组不再把该消息视为待确认消息,不会删除 Stream 中的原始消息。消息是否仍能被其他读取方式看到,取决于 Stream 的保留和裁剪策略。

4. 失败、重试和消息转移

考虑下面的过程:

1. Redis 将消息投递给 worker-1
2. worker-1 执行业务操作
3. worker-1 在 XACK 前崩溃

消息仍然在该组的 PEL 中,不会因为消费者进程消失而自动丢失。其他消费者可以检查:

XPENDING orders order-workers

它会报告待处理消息数量、最小和最大 ID,以及按消费者统计的数量。

较新的 Redis 稳定语义中,可以用 XAUTOCLAIM 把长时间未处理的消息转移给当前消费者:

XAUTOCLAIM orders order-workers worker-2 60000 0-0 COUNT 10

这里的 60000 是最小空闲时间,单位为毫秒。只有空闲超过该时间的待处理消息才适合被认领。认领前应考虑 worker-1 只是暂时变慢,避免多个消费者频繁争抢同一消息。

另一种方式是:

XCLAIM orders order-workers worker-2 60000 1700000000000-0

它适合已经明确知道要转移哪些 ID 的场景。

5. Streams 通常是至少一次,而不是恰好一次

Streams 消费者组的典型语义是至少一次投递

投递消息
执行外部业务
消费者在确认前崩溃
消息被重新认领
再次执行外部业务

因此必须允许重复处理。常见做法是使用业务唯一 ID 和幂等约束:

CREATE TABLE processed_event (
    event_id VARCHAR(100) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

处理时,在同一个数据库事务中:

1. 尝试插入 event_id
2. 若主键冲突,说明已经处理过,跳过业务副作用
3. 若插入成功,执行业务更新
4. 数据库事务提交
5. 再执行 XACK

若第 3 步成功但第 5 步前进程崩溃,消息会再次投递;第 1 步的唯一约束会使第二次处理变成安全跳过。

但如果先 XACK 再更新数据库,进程在两者之间崩溃,就可能出现“消息已确认、业务未完成”。所以确认位置不是越早越好,而是应放在业务状态已经可靠提交之后。

XREADGROUPNOACK 选项可以减少 PEL 记录:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 NOACK STREAMS orders >

但它也意味着消息投递后不再等待确认,消费者崩溃时不能依赖 PEL 恢复。它只适合允许丢失或已具备其他可靠来源的场景。

6. Stream 裁剪与 PEL 的关系

裁剪 Stream 不等于自动完成消费确认。若消息从 Stream 主体中被裁剪,消费者组的待处理状态仍可能需要单独管理,具体行为取决于所使用的裁剪和 Redis 版本语义。生产系统需要同时监控:

  • Stream 长度;
  • 消费者组的组内最后 ID;
  • PEL 数量;
  • 最老待处理消息的空闲时间;
  • 认领和重试次数。

仅仅看到 Stream 长度正常,并不能说明消费者没有积压。


七、Pub/Sub 与 Streams 的选择推导

可以从“消息事实是否需要在未来重新获得”开始判断。

情形一:消息只是状态变化提醒

例如:

数据库商品价格发生变化

消费者即使错过事件,也可以重新读取数据库当前价格。因此可以使用:

数据库作为事实来源
Redis Pub/Sub 作为低延迟通知
缓存 TTL 或定期校验作为兜底

这与缓存一致性中的“失效通知”相匹配:通知丢失不会永久丢失业务事实。

情形二:消息本身就是业务事实

例如:

订单已创建
支付成功
库存扣减任务

消费者不能只看当前状态推导出所有历史动作,且必须重试和追踪进度。这时需要 Streams:

XADD 写入事实
XREADGROUP 投递
业务事务提交
XACK 确认
XPENDING/XAUTOCLAIM 恢复

情形三:需要每个订阅者都收到同一消息

Streams 可以通过多个消费者组实现独立消费:

order-workers     -> 订单处理服务
analytics-workers -> 分析服务
audit-workers     -> 审计服务

每个组都有自己的消费进度。组之间互不共享确认状态。

在同一个消费者组内,多个消费者是竞争关系:一条消息通常只投递给组中的一个消费者。若希望每个服务都得到一份,应该创建多个组,而不是把所有服务放进同一个组。


八、部署边界与诊断方法

1. 单机、主从和集群要分开讨论

在单个主节点上,命令和脚本的原子执行最容易理解。加入复制、哨兵或集群后,还要考虑:

  • 读请求是否误读了副本的旧数据;
  • 故障转移是否丢失尚未复制的写入;
  • Lua 脚本中的键是否位于同一哈希槽;
  • 客户端是否能正确处理 MOVEDASKNOSCRIPT
  • Stream 消费者是否在拓扑变化后重新连接并恢复游标。

不要用副本读取“锁是否存在”来决定是否加锁。锁的判断和写入应在当前主节点上完成。

2. 锁的检查

可以检查:

GET lock:job
PTTL lock:job

需要关注:

  • GET 是否为预期令牌;
  • PTTL 是否为正数;
  • 是否频繁出现接近零的剩余时间;
  • 业务执行时间是否长期超过租约;
  • 是否存在大量锁键因释放脚本未执行而等待过期。

不要用管理命令直接 DEL 业务锁,除非已经确认当前令牌和持有者状态。强制删除可能让原持有者和新持有者同时执行。

3. Lua 的检查

脚本异常时检查:

  • 脚本是否包含错误的键数量;
  • Cluster 中是否跨槽;
  • 是否出现 NOSCRIPT
  • 是否有长循环或一次处理过多成员;
  • 是否使用了可能阻塞的命令;
  • SLOWLOG GET 和延迟监控中是否出现长脚本。

一个逻辑正确但执行几百毫秒的脚本,可能让所有普通请求都排队,因此脚本长度和数据规模必须受控。

4. Pub/Sub 的检查

Pub/Sub 不提供历史查询接口。若怀疑消息丢失,无法通过 Redis 事后列出“某频道过去发布过什么”。应在应用侧记录:

  • 发布日志;
  • 连接建立和断开;
  • 订阅确认;
  • 消费者处理延迟;
  • 关键状态的周期性校验结果。

如果系统需要凭 Redis 本身检查积压、重试或历史消息,Pub/Sub 的数据模型就不合适。

5. Streams 的检查

常用检查命令:

XLEN orders
XINFO STREAM orders
XINFO GROUPS orders
XPENDING orders order-workers
XINFO CONSUMERS orders order-workers

这些命令分别帮助确认:

  • Stream 中有多少条消息;
  • 首尾 ID 和长度等信息;
  • 消费者组的消费位置;
  • PEL 中是否有积压;
  • 哪个消费者持有多少待处理消息、空闲多久。

一个典型故障判断过程是:

Stream 长度上升
    -> 检查组的最后消费 ID
    -> 检查消费者是否仍在线
    -> 检查 PEL 是否持续增长
    -> 检查最老消息的空闲时间
    -> 决定重启、认领、重试或转入死信 Stream

Redis Streams 没有自动为业务定义“死信队列”。若消息重试多次仍失败,应用可以将原消息、错误原因、重试次数写入另一个 Stream,并在确认原消息后进行人工或异步处理。


九、几个容易混淆的结论

“Redis 有原子命令,所以分布式锁绝对安全”

不成立。原子命令解决的是同一 Redis 主节点上的并发插入问题;租约过期、客户端暂停、主从切换和外部资源写入仍需单独处理。

“加了过期时间就不会死锁”

过期时间能处理持锁进程崩溃,但会引入租约过期后的旧持有者问题。任务时间不可预测时,续租和栅栏令牌比单纯增大 TTL 更重要。

“Pub/Sub 发布成功就代表消息被处理”

不成立。PUBLISH 的返回值只是当时的订阅连接数量,不代表业务处理成功,更不代表未来可以重放。

“XACK 后消息就从 Redis 删除了”

不成立。XACK 只更新消费者组的确认状态。Stream 中的消息仍可能存在,直到被裁剪或删除。

“Streams 自动提供恰好一次处理”

不成立。它通常提供至少一次投递,业务必须幂等。数据库唯一键、条件更新、业务状态机和去重表都是常见的幂等基础。

“Lua 脚本等同于跨数据库事务”

不成立。Lua 只原子地修改 Redis 内部状态,不能把 Redis 写入和 MySQL、消息网关或外部服务调用绑定为一个事务。


Redis 的协调能力可以归纳为不同的状态模型:

锁:
    一个键 + 随机令牌 + 租约
    目标是互斥和故障后释放

Lua:
    当前状态 + 参数 -> 原子状态转换
    目标是消除读取与写回之间的竞态

限流:
    计数器或令牌桶状态 + 时间
    目标是限制窗口内数量或长期速率

Pub/Sub:
    当前在线订阅连接
    目标是低延迟广播,不保存历史

Streams:
    追加消息 + 消费者组游标 + PEL
    目标是可追踪、可恢复的消息处理

正确使用它们的关键,不是记住某条命令,而是先确定业务需要的状态是否持久、消息是否允许丢失、失败后是否必须重试,以及旧客户端在租约失效后是否仍可能写入资源。只有这些边界清楚后,锁、脚本、限流、广播和消息流才会各自处在适合的位置。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。

评论

0 条讨论
0/1000
还没有评论,来聊聊你的看法
,命令执行前已经写入的消息不会被读取。可靠消费者应保存上次成功处理的 ID,并从该 ID 继续读取,而不是每次重连都使用 ` Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams - WR Blog
WR Blog 加载中...
返回文章
数据库Redis分布式锁消息队列

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams封面

数据库基础体系 · 第 39/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

Redis 分布式协调:锁、Lua、限流、Pub/Sub 与 Streams

Redis 不只是缓存。它还提供了一组建立在内存数据结构、单命令原子性、Lua 脚本和持久化日志之上的协调能力:

  • 用键值状态实现带租约的分布式锁;
  • 用 Lua 把“读取—判断—写回”合并为不可插入的原子操作;
  • 用计数器、时间戳和令牌桶实现限流;
  • 用 Pub/Sub 传递实时但不持久的通知;
  • 用 Streams 保存消息,并通过消费者组实现可恢复的消费进度。

这些能力的可靠性边界并不相同。锁和限流主要依赖 Redis 主节点上操作的原子性;Pub/Sub 依赖在线连接;Streams 则把消息和消费状态保存为 Redis 数据,因此可以处理消费者短暂故障。将它们混为“Redis 消息队列”或“Redis 事务”会导致错误的设计。


一、共同基础:Redis 的原子性到底保证什么

1. Redis 命令的串行执行

对一个 Redis 实例来说,普通命令在主线程中按顺序执行。假设当前键 counter 的值为 10,两个客户端同时执行:

客户端 A:GET counter
客户端 B:GET counter
客户端 A:SET counter 11
客户端 B:SET counter 11

GETSET 各自是原子的,但整个“读取后加一”不是原子的,最终结果可能是 11 而不是期望的 12

使用单条命令:

INCR counter

或者使用 Lua 脚本,可以让一组操作在执行期间不被其他 Redis 命令插入。

这里的“原子”表示:

  1. 脚本或命令执行期间,其他客户端命令不会交错执行;
  2. 其他客户端不会看到脚本执行过程中的中间状态;
  3. 它不表示业务操作具备跨系统事务;
  4. 它也不表示脚本发生错误后所有已经执行的写入会自动回滚。

Redis 的 Lua 脚本应当短小、确定且不执行阻塞操作。脚本执行时间过长会阻塞整个实例上的其他命令。脚本中的 redis.call 出错时,脚本会报告错误;脚本此前已经产生的写入不能依赖“自动回滚”来恢复,因此应在写入前完成参数校验。

2. MULTI/EXEC 与 Lua 的区别

Redis 事务通常指:

MULTI
SET a 1
SET b 2
EXEC

MULTI 之后命令会先排队,EXEC 时连续执行。事务执行期间,其他客户端命令不能插入这批命令。

但 Redis 事务不是关系数据库意义上的完整 ACID 事务:

  • 没有传统意义的回滚机制;
  • WATCH 提供的是乐观并发控制;
  • 它不自动覆盖外部数据库、文件系统或 HTTP 调用;
  • Redis Cluster 中事务涉及的键通常必须位于同一个哈希槽。

Lua 更适合“根据当前值判断后再写入”的逻辑,例如:

读取锁值
判断是否属于当前持有者
删除锁

如果拆成 GETDEL,中间可能被其他客户端插入;如果放入一个脚本,整个判断和删除不可被插入。


二、分布式锁:带租约的互斥,而不是永久所有权

1. 锁的基本模型

分布式锁通常需要满足三个条件:

  1. 互斥性:同一时刻至多一个客户端被认为持有锁;
  2. 释放安全:客户端只能释放自己持有的锁;
  3. 故障可恢复:持锁客户端崩溃后,锁最终可以再次获取。

在 Redis 中,最基本的加锁命令是:

SET lock:report:daily 9f3c... NX PX 10000

参数含义:

  • NX:仅当键不存在时设置;
  • PX 10000:设置 10 秒毫秒级过期时间;
  • 9f3c...:本次加锁请求生成的随机令牌。

成功时返回:

OK

失败时通常返回:

(nil)

不能使用下面这种两步操作代替:

SETNX lock:report:daily 1
EXPIRE lock:report:daily 10

因为客户端可能在两条命令之间崩溃,留下永不过期的锁。

2. 为什么锁值必须是随机令牌

假设客户端 A 获取锁:

lock:job = token-A

A 执行任务超过租约时间,锁自动过期。之后客户端 B 获取:

lock:job = token-B

如果 A 这时直接执行:

DEL lock:job

A 删除的实际上是 B 的锁。

因此释放锁必须是“比较令牌并删除”的原子操作:

-- unlock.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('DEL', KEYS[1])
end
return 0

执行:

EVAL "$(cat unlock.lua)" 1 lock:job token-A

返回值:

  • 1:比较成功并删除;
  • 0:键不存在,或者当前令牌不是 token-A

不能写成:

GET lock:job
如果等于 token-A,再执行 DEL

因为 GETDEL 之间可能发生过期、重新加锁或其他客户端写入。

3. 租约过期后的真实边界

Redis 锁的过期时间使它成为带租约的锁。它只能保证:

在 Redis 认为租约仍有效、且部署故障模型满足假设时,持有者拥有锁。

它不能保证客户端在业务代码中永远拥有锁。例如:

  1. A 获取锁,租约为 10 秒;
  2. A 因为 GC、进程暂停或网络阻塞,20 秒后才继续执行;
  3. Redis 中的锁早已过期;
  4. B 已经重新获取锁并修改资源;
  5. A 恢复后继续写入。

即使 A 后续调用安全解锁脚本,也只能避免删除 B 的锁,不能阻止 A 已经对业务资源执行过期写入。

因此,对于“过期持有者不能覆盖新持有者”的严格要求,需要栅栏令牌(fencing token)

4. 栅栏令牌:把锁状态传递给资源

获取锁时,同时从 Redis 生成单调递增的令牌:

MULTI
INCR lock:job:fence
SET lock:job token-A NX PX 10000
EXEC

实际应用需要处理 SET 失败时不要误用已经生成的令牌。更常见的方式是把“生成令牌”和“抢锁”封装在 Lua 中:

-- acquire-with-fence.lua
if redis.call('EXISTS', KEYS[1]) == 1 then
    return {0, false}
end

local token = redis.call('INCR', KEYS[2])
redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[2])
return {1, token}

调用:

EVAL "$(cat acquire-with-fence.lua)" 2 lock:job lock:job:fence token-A 10000

成功可能返回:

1
1

客户端把令牌 1 传给真正的资源服务。资源服务保存“最后接受的令牌”,只接受更大的令牌:

请求 A:fence=1,接受
请求 B:fence=2,接受
请求 A:fence=1,拒绝

这样,Redis 只负责发放单调令牌,资源本身负责拒绝旧持有者。若资源是数据库,可以通过条件更新实现:

UPDATE job_state
SET result = :result,
    last_fence = :fence
WHERE job_id = :job_id
  AND last_fence < :fence;

这一步必须和资源更新处于同一个数据库事务中,否则检查和更新之间仍可能被插入。

5. 续租与释放

如果任务可能超过初始租约,应由持有者定期续租。续租也必须比较令牌:

-- renew.lua
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('PEXPIRE', KEYS[1], ARGV[2])
end
return 0

执行:

EVAL "$(cat renew.lua)" 1 lock:job token-A 10000

续租周期不能等到租约最后一刻才执行,否则网络延迟或进程暂停可能造成误过期。与此同时,业务代码仍应设计为:续租失败后停止产生新的副作用,而不是继续假设自己拥有锁。

6. 主从复制和故障转移边界

在单个 Redis 主节点上,SET NX PX 的互斥判断是原子的。但如果主节点宕机,副本尚未复制最新锁,故障转移后副本可能认为锁不存在,另一个客户端便可获取同名锁。

因此:

  • Redis 主从复制默认是异步的;
  • WAIT 可以等待写入传播到若干副本,但不等于共识协议;
  • 网络分区、故障转移和旧主恢复时,不能仅凭普通 Redis 锁得到严格的线性一致互斥;
  • 对高价值资源,栅栏令牌、数据库约束或具备共识语义的协调系统仍然重要。

把多个 Redis 节点拼成所谓“更强的锁”也不会自动消除时钟、网络分区和客户端暂停问题。应先明确业务需要的是“尽量避免重复执行”,还是“在故障条件下也绝不允许旧持有者写入”,两者的系统设计不同。


三、Lua:把条件、读取和写入绑定为一个状态转换

1. 脚本的输入边界

Lua 脚本通过两类参数接收数据:

  • KEYS:脚本访问的 Redis 键;
  • ARGV:普通参数。

例如:

local current = redis.call('GET', KEYS[1])
local limit = tonumber(ARGV[1])

if current and tonumber(current) >= limit then
    return 0
end

redis.call('INCR', KEYS[1])
return 1

调用:

EVAL "$(cat increment-if-below.lua)" 1 quota:user:42 10

在 Redis Cluster 中,脚本访问的键必须满足集群路由要求,通常应位于同一个哈希槽。可以使用哈希标签:

quota:{user-42}:count
quota:{user-42}:timestamp

{user-42} 使两个键使用相同的哈希标签。脚本中应把所有键放入 KEYS,不要把键名藏在 ARGV 中。这样 Redis 才能正确检查和路由脚本涉及的键。

2. EVAL、EVALSHA 与脚本缓存

直接执行:

EVAL "return redis.call('INCR', KEYS[1])" 1 counter

也可以先加载脚本:

SCRIPT LOAD "return redis.call('INCR', KEYS[1])"

Redis 返回一个 SHA-1 摘要,例如:

"7d9..."

之后执行:

EVALSHA 7d9... 1 counter

EVALSHA 只引用当前实例脚本缓存中的脚本。重启、故障转移或连接到另一个节点后,可能出现 NOSCRIPT。客户端必须能够重新 SCRIPT LOAD 或回退到 EVAL。脚本缓存不是持久化的业务配置中心。

对于固定部署的脚本,也可以使用 Redis Functions。它们适合把函数随函数库加载和管理,但仍然受到脚本执行原子性、执行时长和集群键约束等边界限制。不能因为使用 Functions 就获得跨服务事务。

3. 脚本中的时间和随机性

涉及时间的脚本必须考虑时间来源:

  • 使用客户端传入的时间,可能受到不同机器时钟偏差影响;
  • 使用 Redis 时间,可减少客户端时钟差异,但脚本仍应保持短小,并遵守当前 Redis 版本对脚本可调用命令和复制的限制;
  • 不应让不同客户端用完全不一致的时间更新同一个限流状态。

对于需要可复制、可重放的脚本,应避免依赖不可控的随机行为和不确定的外部状态。脚本最好把必要的时间、数量和策略参数显式作为输入。


四、限流:从固定窗口到令牌桶

限流的目标是限制某个主体在一个时间范围内可执行的请求数。主体可以是用户、IP、API、租户或资源键。

需要先区分两个量:

  • 吞吐限制:长期平均每秒允许多少请求;
  • 突发容量:短时间内最多允许多少请求。

固定窗口容易实现,但窗口边界可能产生突发;令牌桶能分别表达这两个量。

1. 固定窗口计数器

目标:每个用户每 60 秒最多请求 3 次。

Lua 脚本:

-- fixed-window.lua
local count = redis.call('INCR', KEYS[1])

if count == 1 then
    redis.call('EXPIRE', KEYS[1], ARGV[1])
end

if count <= tonumber(ARGV[2]) then
    return {1, count, tonumber(ARGV[2]) - count}
end

return {0, count, 0}

调用:

EVAL "$(cat fixed-window.lua)" 1 rate:user:42:window 60 3

可能返回:

1
2
1

含义是:

允许,当前窗口计数为 2,剩余 1 次

脚本中只有第一次计数时设置过期时间,因此不会因为后续请求不断刷新 TTL 而形成永不过期的计数器。

但固定窗口存在边界问题。假设窗口为:

[00:00:00, 00:01:00)
[00:01:00, 00:02:00)

用户在 00:00:59 请求 3 次,又在 00:01:00 请求 3 次。两秒内可以通过 6 次请求,虽然每个窗口分别不超过 3 次。这不是实现错误,而是固定窗口算法的定义结果。

2. 滑动窗口日志

滑动窗口要求检查最近 60 秒内的请求数。可以使用 Sorted Set:

  • member:请求唯一 ID;
  • score:请求时间戳。

每次请求执行:

删除 score <= now - 60000 的成员
统计剩余成员
若小于限制则加入当前请求

对应 Lua:

-- sliding-window.lua
local now = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local limit = tonumber(ARGV[3])
local request_id = ARGV[4]

redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', now - window)
local count = redis.call('ZCARD', KEYS[1])

if count >= limit then
    redis.call('PEXPIRE', KEYS[1], window)
    return {0, count}
end

redis.call('ZADD', KEYS[1], now, request_id)
redis.call('PEXPIRE', KEYS[1], window)
return {1, count + 1}

调用示例:

EVAL "$(cat sliding-window.lua)" 1 rate:user:42:log 1700000000000 60000 3 req-001

它比固定窗口精确,但每个请求都需要维护一个成员,内存和删除成本更高。request_id 必须唯一;如果相同 member 重复使用,ZADD 会更新原成员而不是增加计数。

3. 令牌桶

令牌桶用两个参数表达限制:

  • capacity:桶的最大令牌数,表示允许的突发容量;
  • rate:每秒补充的令牌数,表示长期平均速率。

状态为:

(tokens, last_timestamp)

在时间 now 到来时,先补充:

tokens=min(capacity, tokens+(nowlast)×rate)tokens' = \min(capacity,\ tokens + (now-last)\times rate)

其中:

  • tokens 是上次计算后的剩余令牌;
  • last 是上次更新时间;
  • rate 的单位是“令牌/秒”;
  • now-last 必须换算成秒;
  • capacity 防止空闲期间无限积累。

若请求需要 cost 个令牌:

tokens' >= cost  => 允许,并保存 tokens' - cost
tokens' <  cost  => 拒绝,并保存补充后的 tokens'

一个使用毫秒时间的 Lua 实现如下:

-- token-bucket.lua
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])       -- tokens per second
local now_ms = tonumber(ARGV[3])
local cost = tonumber(ARGV[4])
local ttl_ms = tonumber(ARGV[5])

local values = redis.call('HMGET', KEYS[1], 'tokens', 'ts')
local tokens = tonumber(values[1])
local last_ms = tonumber(values[2])

if tokens == nil then
    tokens = capacity
    last_ms = now_ms
end

-- 客户端时钟回拨时,不让令牌数量因负时间差增加异常
if now_ms < last_ms then
    now_ms = last_ms
end

local elapsed_seconds = (now_ms - last_ms) / 1000.0
tokens = math.min(capacity, tokens + elapsed_seconds * rate)

local allowed = 0
local retry_after_ms = 0

if tokens >= cost then
    tokens = tokens - cost
    allowed = 1
else
    local missing = cost - tokens
    retry_after_ms = math.ceil(missing / rate * 1000)
end

redis.call('HSET', KEYS[1], 'tokens', tokens, 'ts', now_ms)
redis.call('PEXPIRE', KEYS[1], ttl_ms)

return {allowed, tokens, retry_after_ms}

调用:

EVAL "$(cat token-bucket.lua)" 1 rate:user:42:bucket 10 2 1700000000000 1 10000

参数含义是:

  • 容量 10;
  • 每秒补充 2 个令牌;
  • 当前时间戳为 1700000000000 毫秒;
  • 本次请求消耗 1 个令牌;
  • 状态空闲 10 秒后过期。

返回:

1
9
0

表示允许请求,剩余约 9 个令牌,不需要等待。

如果返回:

0
0.5
250

表示拒绝,当前只有 0.5 个令牌,理论上约 250 毫秒后才能补足一个令牌。

这个示例把当前时间作为参数传入,因此调用方必须统一时间来源。若不同应用服务器的时钟偏差很大,同一个桶的状态会出现不一致。可以由 Redis 提供时间来源,或者由网关统一生成时间;关键是不要让多个不受控的时钟共同推进同一个限流状态。

令牌桶脚本中的 rate 不能为 0,否则拒绝请求时计算等待时间会发生除零。生产实现还应限制参数范围,防止客户端传入极大的容量、速率或 TTL。

4. 限流失败路径

限流通常有三种结果:

  1. Redis 判断允许,业务请求继续;
  2. Redis 判断拒绝,返回 HTTP 429 Too Many Requests
  3. Redis 不可用或超时。

第三种不能简单等同于“限流通过”。若放行,可能在 Redis 故障时失去保护;若拒绝,可能把 Redis 故障扩大为全站不可用。应根据接口风险决定:

  • 登录、验证码、支付等安全敏感接口偏向失败关闭;
  • 非关键读接口可能使用本地兜底限流;
  • 本地兜底必须明确它只是每实例限流,不是全局限流。

五、Pub/Sub:在线广播,不是可靠队列

1. 基本数据流

Redis Pub/Sub 有两个角色:

  • 发布者:向频道发送消息;
  • 订阅者:通过长连接订阅频道并接收消息。

订阅:

SUBSCRIBE cache:product:updated

成功后 Redis 会返回订阅确认,例如:

subscribe
cache:product:updated
1

发布:

PUBLISH cache:product:updated '{"product_id":42}'

返回值是当时收到该消息的订阅客户端数量,例如:

(integer) 2

这个返回值不是消费成功数,也不是持久化成功数。它只表示消息发布时有多少订阅连接被送达。

发布数据流可以表示为:

发布者
  |
  | PUBLISH
  v
Redis 当前连接集合
  |
  +--> 订阅者 A
  +--> 订阅者 B

Redis 不把 Pub/Sub 消息作为普通键保存。订阅者掉线期间发布的消息不会等待它重新连接后补发。

2. 可靠性语义

Redis Pub/Sub 通常应理解为:

在线连接上的即时、尽力发送、至多一次交付通知。

常见失败路径:

  1. 发布者发布时没有订阅者:消息直接消失;
  2. 订阅者网络断开:断线期间的消息不会补发;
  3. 订阅者收到消息后进程崩溃:没有 Redis 侧确认和重投;
  4. 消费者处理很慢:消息可能在客户端连接缓冲区积压,最终导致连接问题。

因此,以下场景适合 Pub/Sub:

  • 缓存失效通知;
  • 配置刷新提示;
  • 在线 WebSocket 广播;
  • 不要求离线补偿的状态变化通知。

以下场景不应只使用 Pub/Sub:

  • 订单、支付、库存等不可丢失事件;
  • 必须重新消费的任务;
  • 需要确认、重试和消费进度的消息。

如果一个消费者错过通知后可以通过读取当前状态自行修复,Pub/Sub 可以作为低延迟通知通道。例如收到“商品更新”后删除本地缓存;即使通知丢失,下一次 TTL 到期或主动校验仍可修复。若消息本身就是唯一业务事实,则应使用 Streams 或其他持久化消息系统。

3. 频道订阅与模式订阅

普通订阅:

SUBSCRIBE tenant:42:events

模式订阅:

PSUBSCRIBE tenant:*:events

模式订阅更灵活,但匹配和消息分发成本更高。频道命名应包含明确的租户、业务和版本边界,避免所有业务共享一个频道后由客户端自行过滤。

Redis 还提供按分片路由的 Sharded Pub/Sub 命令,适用于 Redis Cluster 中希望减少跨节点广播的场景。它与普通 Pub/Sub 的连接、拓扑和客户端支持要求不同,不能假设所有客户端都自动兼容;使用前应确认客户端对对应命令和集群重定向的支持。


六、Streams:可持久化的追加日志与消费者组

1. Stream 是什么

Redis Stream 是一个按 ID 排序的追加记录结构。每条记录包含:

消息 ID -> field/value 字段

写入一条消息:

XADD orders * order_id 1001 status created

返回类似:

"1700000000000-0"

* 表示由 Redis 生成 ID。ID 通常包含毫秒时间部分和序列部分,不能把它当作业务时间或全局业务排序的替代品。消费者应保存和比较 Stream ID,而不是仅用时间戳。

读取:

XRANGE orders - +

可能返回:

1) 1) "1700000000000-0"
   2) 1) "order_id"
      2) "1001"
      3) "status"
      4) "created"

生产环境通常需要限制 Stream 长度,例如:

XADD orders MAXLEN ~ 100000 * order_id 1001 status created

~ 表示近似裁剪,Redis 可以用更高效的方式接近目标长度;它不是严格保证长度永远不超过 100000。若必须精确限制,可以去掉 ~,但裁剪成本和写入性能边界不同。

2. 普通消费者读取

从头读取:

XREAD COUNT 10 STREAMS orders 0-0

从某个 ID 之后读取:

XREAD COUNT 10 STREAMS orders 1700000000000-0

阻塞读取:

XREAD BLOCK 5000 COUNT 10 STREAMS orders $

这里的 $ 表示只关注执行命令时之后新增的消息。它适合实时监听,但如果客户端启动或重连时使用 $,命令执行前已经写入的消息不会被读取。可靠消费者应保存上次成功处理的 ID,并从该 ID 继续读取,而不是每次重连都使用 $

3. 消费者组和消费确认

消费者组让多个消费者共同处理一个 Stream,并记录组内消费位置。

创建组:

XGROUP CREATE orders order-workers 0-0 MKSTREAM

含义:

  • Stream:orders
  • 消费者组:order-workers
  • 0-0 开始;
  • 如果 Stream 不存在则创建。

消费者以名字加入:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 BLOCK 5000 STREAMS orders >

> 表示读取当前组尚未投递给任何消费者的新消息。

Redis 返回的消息会进入该消费者组的待处理列表,称为 PEL(Pending Entries List)。此时消息已经投递,但还没有确认。

业务处理成功后确认:

XACK orders order-workers 1700000000000-0

XACK 只表示该消费者组不再把该消息视为待确认消息,不会删除 Stream 中的原始消息。消息是否仍能被其他读取方式看到,取决于 Stream 的保留和裁剪策略。

4. 失败、重试和消息转移

考虑下面的过程:

1. Redis 将消息投递给 worker-1
2. worker-1 执行业务操作
3. worker-1 在 XACK 前崩溃

消息仍然在该组的 PEL 中,不会因为消费者进程消失而自动丢失。其他消费者可以检查:

XPENDING orders order-workers

它会报告待处理消息数量、最小和最大 ID,以及按消费者统计的数量。

较新的 Redis 稳定语义中,可以用 XAUTOCLAIM 把长时间未处理的消息转移给当前消费者:

XAUTOCLAIM orders order-workers worker-2 60000 0-0 COUNT 10

这里的 60000 是最小空闲时间,单位为毫秒。只有空闲超过该时间的待处理消息才适合被认领。认领前应考虑 worker-1 只是暂时变慢,避免多个消费者频繁争抢同一消息。

另一种方式是:

XCLAIM orders order-workers worker-2 60000 1700000000000-0

它适合已经明确知道要转移哪些 ID 的场景。

5. Streams 通常是至少一次,而不是恰好一次

Streams 消费者组的典型语义是至少一次投递

投递消息
执行外部业务
消费者在确认前崩溃
消息被重新认领
再次执行外部业务

因此必须允许重复处理。常见做法是使用业务唯一 ID 和幂等约束:

CREATE TABLE processed_event (
    event_id VARCHAR(100) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

处理时,在同一个数据库事务中:

1. 尝试插入 event_id
2. 若主键冲突,说明已经处理过,跳过业务副作用
3. 若插入成功,执行业务更新
4. 数据库事务提交
5. 再执行 XACK

若第 3 步成功但第 5 步前进程崩溃,消息会再次投递;第 1 步的唯一约束会使第二次处理变成安全跳过。

但如果先 XACK 再更新数据库,进程在两者之间崩溃,就可能出现“消息已确认、业务未完成”。所以确认位置不是越早越好,而是应放在业务状态已经可靠提交之后。

XREADGROUPNOACK 选项可以减少 PEL 记录:

XREADGROUP GROUP order-workers worker-1 \
  COUNT 10 NOACK STREAMS orders >

但它也意味着消息投递后不再等待确认,消费者崩溃时不能依赖 PEL 恢复。它只适合允许丢失或已具备其他可靠来源的场景。

6. Stream 裁剪与 PEL 的关系

裁剪 Stream 不等于自动完成消费确认。若消息从 Stream 主体中被裁剪,消费者组的待处理状态仍可能需要单独管理,具体行为取决于所使用的裁剪和 Redis 版本语义。生产系统需要同时监控:

  • Stream 长度;
  • 消费者组的组内最后 ID;
  • PEL 数量;
  • 最老待处理消息的空闲时间;
  • 认领和重试次数。

仅仅看到 Stream 长度正常,并不能说明消费者没有积压。


七、Pub/Sub 与 Streams 的选择推导

可以从“消息事实是否需要在未来重新获得”开始判断。

情形一:消息只是状态变化提醒

例如:

数据库商品价格发生变化

消费者即使错过事件,也可以重新读取数据库当前价格。因此可以使用:

数据库作为事实来源
Redis Pub/Sub 作为低延迟通知
缓存 TTL 或定期校验作为兜底

这与缓存一致性中的“失效通知”相匹配:通知丢失不会永久丢失业务事实。

情形二:消息本身就是业务事实

例如:

订单已创建
支付成功
库存扣减任务

消费者不能只看当前状态推导出所有历史动作,且必须重试和追踪进度。这时需要 Streams:

XADD 写入事实
XREADGROUP 投递
业务事务提交
XACK 确认
XPENDING/XAUTOCLAIM 恢复

情形三:需要每个订阅者都收到同一消息

Streams 可以通过多个消费者组实现独立消费:

order-workers     -> 订单处理服务
analytics-workers -> 分析服务
audit-workers     -> 审计服务

每个组都有自己的消费进度。组之间互不共享确认状态。

在同一个消费者组内,多个消费者是竞争关系:一条消息通常只投递给组中的一个消费者。若希望每个服务都得到一份,应该创建多个组,而不是把所有服务放进同一个组。


八、部署边界与诊断方法

1. 单机、主从和集群要分开讨论

在单个主节点上,命令和脚本的原子执行最容易理解。加入复制、哨兵或集群后,还要考虑:

  • 读请求是否误读了副本的旧数据;
  • 故障转移是否丢失尚未复制的写入;
  • Lua 脚本中的键是否位于同一哈希槽;
  • 客户端是否能正确处理 MOVEDASKNOSCRIPT
  • Stream 消费者是否在拓扑变化后重新连接并恢复游标。

不要用副本读取“锁是否存在”来决定是否加锁。锁的判断和写入应在当前主节点上完成。

2. 锁的检查

可以检查:

GET lock:job
PTTL lock:job

需要关注:

  • GET 是否为预期令牌;
  • PTTL 是否为正数;
  • 是否频繁出现接近零的剩余时间;
  • 业务执行时间是否长期超过租约;
  • 是否存在大量锁键因释放脚本未执行而等待过期。

不要用管理命令直接 DEL 业务锁,除非已经确认当前令牌和持有者状态。强制删除可能让原持有者和新持有者同时执行。

3. Lua 的检查

脚本异常时检查:

  • 脚本是否包含错误的键数量;
  • Cluster 中是否跨槽;
  • 是否出现 NOSCRIPT
  • 是否有长循环或一次处理过多成员;
  • 是否使用了可能阻塞的命令;
  • SLOWLOG GET 和延迟监控中是否出现长脚本。

一个逻辑正确但执行几百毫秒的脚本,可能让所有普通请求都排队,因此脚本长度和数据规模必须受控。

4. Pub/Sub 的检查

Pub/Sub 不提供历史查询接口。若怀疑消息丢失,无法通过 Redis 事后列出“某频道过去发布过什么”。应在应用侧记录:

  • 发布日志;
  • 连接建立和断开;
  • 订阅确认;
  • 消费者处理延迟;
  • 关键状态的周期性校验结果。

如果系统需要凭 Redis 本身检查积压、重试或历史消息,Pub/Sub 的数据模型就不合适。

5. Streams 的检查

常用检查命令:

XLEN orders
XINFO STREAM orders
XINFO GROUPS orders
XPENDING orders order-workers
XINFO CONSUMERS orders order-workers

这些命令分别帮助确认:

  • Stream 中有多少条消息;
  • 首尾 ID 和长度等信息;
  • 消费者组的消费位置;
  • PEL 中是否有积压;
  • 哪个消费者持有多少待处理消息、空闲多久。

一个典型故障判断过程是:

Stream 长度上升
    -> 检查组的最后消费 ID
    -> 检查消费者是否仍在线
    -> 检查 PEL 是否持续增长
    -> 检查最老消息的空闲时间
    -> 决定重启、认领、重试或转入死信 Stream

Redis Streams 没有自动为业务定义“死信队列”。若消息重试多次仍失败,应用可以将原消息、错误原因、重试次数写入另一个 Stream,并在确认原消息后进行人工或异步处理。


九、几个容易混淆的结论

“Redis 有原子命令,所以分布式锁绝对安全”

不成立。原子命令解决的是同一 Redis 主节点上的并发插入问题;租约过期、客户端暂停、主从切换和外部资源写入仍需单独处理。

“加了过期时间就不会死锁”

过期时间能处理持锁进程崩溃,但会引入租约过期后的旧持有者问题。任务时间不可预测时,续租和栅栏令牌比单纯增大 TTL 更重要。

“Pub/Sub 发布成功就代表消息被处理”

不成立。PUBLISH 的返回值只是当时的订阅连接数量,不代表业务处理成功,更不代表未来可以重放。

“XACK 后消息就从 Redis 删除了”

不成立。XACK 只更新消费者组的确认状态。Stream 中的消息仍可能存在,直到被裁剪或删除。

“Streams 自动提供恰好一次处理”

不成立。它通常提供至少一次投递,业务必须幂等。数据库唯一键、条件更新、业务状态机和去重表都是常见的幂等基础。

“Lua 脚本等同于跨数据库事务”

不成立。Lua 只原子地修改 Redis 内部状态,不能把 Redis 写入和 MySQL、消息网关或外部服务调用绑定为一个事务。


Redis 的协调能力可以归纳为不同的状态模型:

锁:
    一个键 + 随机令牌 + 租约
    目标是互斥和故障后释放

Lua:
    当前状态 + 参数 -> 原子状态转换
    目标是消除读取与写回之间的竞态

限流:
    计数器或令牌桶状态 + 时间
    目标是限制窗口内数量或长期速率

Pub/Sub:
    当前在线订阅连接
    目标是低延迟广播,不保存历史

Streams:
    追加消息 + 消费者组游标 + PEL
    目标是可追踪、可恢复的消息处理

正确使用它们的关键,不是记住某条命令,而是先确定业务需要的状态是否持久、消息是否允许丢失、失败后是否必须重试,以及旧客户端在租约失效后是否仍可能写入资源。只有这些边界清楚后,锁、脚本、限流、广播和消息流才会各自处在适合的位置。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。

评论

0 条讨论
0/1000
还没有评论,来聊聊你的看法
。\n\n### 3. 消费者组和消费确认\n\n消费者组让多个消费者共同处理一个 Stream,并记录组内消费位置。\n\n创建组:\n\n```redis\nXGROUP CREATE orders order-workers 0-0 MKSTREAM\n```\n\n含义:\n\n- Stream:`orders`;\n- 消费者组:`order-workers`;\n- 从 `0-0` 开始;\n- 如果 Stream 不存在则创建。\n\n消费者以名字加入:\n\n```redis\nXREADGROUP GROUP order-workers worker-1 \\\n COUNT 10 BLOCK 5000 STREAMS orders >\n```\n\n`>` 表示读取当前组尚未投递给任何消费者的新消息。\n\nRedis 返回的消息会进入该消费者组的待处理列表,称为 **PEL(Pending Entries List)**。此时消息已经投递,但还没有确认。\n\n业务处理成功后确认:\n\n```redis\nXACK orders order-workers 1700000000000-0\n```\n\n`XACK` 只表示该消费者组不再把该消息视为待确认消息,不会删除 Stream 中的原始消息。消息是否仍能被其他读取方式看到,取决于 Stream 的保留和裁剪策略。\n\n### 4. 失败、重试和消息转移\n\n考虑下面的过程:\n\n```text\n1. Redis 将消息投递给 worker-1\n2. worker-1 执行业务操作\n3. worker-1 在 XACK 前崩溃\n```\n\n消息仍然在该组的 PEL 中,不会因为消费者进程消失而自动丢失。其他消费者可以检查:\n\n```redis\nXPENDING orders order-workers\n```\n\n它会报告待处理消息数量、最小和最大 ID,以及按消费者统计的数量。\n\n较新的 Redis 稳定语义中,可以用 `XAUTOCLAIM` 把长时间未处理的消息转移给当前消费者:\n\n```redis\nXAUTOCLAIM orders order-workers worker-2 60000 0-0 COUNT 10\n```\n\n这里的 `60000` 是最小空闲时间,单位为毫秒。只有空闲超过该时间的待处理消息才适合被认领。认领前应考虑 worker-1 只是暂时变慢,避免多个消费者频繁争抢同一消息。\n\n另一种方式是:\n\n```redis\nXCLAIM orders order-workers worker-2 60000 1700000000000-0\n```\n\n它适合已经明确知道要转移哪些 ID 的场景。\n\n### 5. Streams 通常是至少一次,而不是恰好一次\n\nStreams 消费者组的典型语义是**至少一次投递**:\n\n```text\n投递消息\n执行外部业务\n消费者在确认前崩溃\n消息被重新认领\n再次执行外部业务\n```\n\n因此必须允许重复处理。常见做法是使用业务唯一 ID 和幂等约束:\n\n```sql\nCREATE TABLE processed_event (\n event_id VARCHAR(100) PRIMARY KEY,\n processed_at TIMESTAMP NOT NULL\n);\n```\n\n处理时,在同一个数据库事务中:\n\n```text\n1. 尝试插入 event_id\n2. 若主键冲突,说明已经处理过,跳过业务副作用\n3. 若插入成功,执行业务更新\n4. 数据库事务提交\n5. 再执行 XACK\n```\n\n若第 3 步成功但第 5 步前进程崩溃,消息会再次投递;第 1 步的唯一约束会使第二次处理变成安全跳过。\n\n但如果先 `XACK` 再更新数据库,进程在两者之间崩溃,就可能出现“消息已确认、业务未完成”。所以确认位置不是越早越好,而是应放在业务状态已经可靠提交之后。\n\n`XREADGROUP` 的 `NOACK` 选项可以减少 PEL 记录:\n\n```redis\nXREADGROUP GROUP order-workers worker-1 \\\n COUNT 10 NOACK STREAMS orders >\n```\n\n但它也意味着消息投递后不再等待确认,消费者崩溃时不能依赖 PEL 恢复。它只适合允许丢失或已具备其他可靠来源的场景。\n\n### 6. Stream 裁剪与 PEL 的关系\n\n裁剪 Stream 不等于自动完成消费确认。若消息从 Stream 主体中被裁剪,消费者组的待处理状态仍可能需要单独管理,具体行为取决于所使用的裁剪和 Redis 版本语义。生产系统需要同时监控:\n\n- Stream 长度;\n- 消费者组的组内最后 ID;\n- PEL 数量;\n- 最老待处理消息的空闲时间;\n- 认领和重试次数。\n\n仅仅看到 Stream 长度正常,并不能说明消费者没有积压。\n\n---\n\n## 七、Pub/Sub 与 Streams 的选择推导\n\n可以从“消息事实是否需要在未来重新获得”开始判断。\n\n### 情形一:消息只是状态变化提醒\n\n例如:\n\n```text\n数据库商品价格发生变化\n```\n\n消费者即使错过事件,也可以重新读取数据库当前价格。因此可以使用:\n\n```text\n数据库作为事实来源\nRedis Pub/Sub 作为低延迟通知\n缓存 TTL 或定期校验作为兜底\n```\n\n这与缓存一致性中的“失效通知”相匹配:通知丢失不会永久丢失业务事实。\n\n### 情形二:消息本身就是业务事实\n\n例如:\n\n```text\n订单已创建\n支付成功\n库存扣减任务\n```\n\n消费者不能只看当前状态推导出所有历史动作,且必须重试和追踪进度。这时需要 Streams:\n\n```text\nXADD 写入事实\nXREADGROUP 投递\n业务事务提交\nXACK 确认\nXPENDING/XAUTOCLAIM 恢复\n```\n\n### 情形三:需要每个订阅者都收到同一消息\n\nStreams 可以通过多个消费者组实现独立消费:\n\n```text\norder-workers -> 订单处理服务\nanalytics-workers -> 分析服务\naudit-workers -> 审计服务\n```\n\n每个组都有自己的消费进度。组之间互不共享确认状态。\n\n在同一个消费者组内,多个消费者是竞争关系:一条消息通常只投递给组中的一个消费者。若希望每个服务都得到一份,应该创建多个组,而不是把所有服务放进同一个组。\n\n---\n\n## 八、部署边界与诊断方法\n\n### 1. 单机、主从和集群要分开讨论\n\n在单个主节点上,命令和脚本的原子执行最容易理解。加入复制、哨兵或集群后,还要考虑:\n\n- 读请求是否误读了副本的旧数据;\n- 故障转移是否丢失尚未复制的写入;\n- Lua 脚本中的键是否位于同一哈希槽;\n- 客户端是否能正确处理 `MOVED`、`ASK` 和 `NOSCRIPT`;\n- Stream 消费者是否在拓扑变化后重新连接并恢复游标。\n\n不要用副本读取“锁是否存在”来决定是否加锁。锁的判断和写入应在当前主节点上完成。\n\n### 2. 锁的检查\n\n可以检查:\n\n```redis\nGET lock:job\nPTTL lock:job\n```\n\n需要关注:\n\n- `GET` 是否为预期令牌;\n- `PTTL` 是否为正数;\n- 是否频繁出现接近零的剩余时间;\n- 业务执行时间是否长期超过租约;\n- 是否存在大量锁键因释放脚本未执行而等待过期。\n\n不要用管理命令直接 `DEL` 业务锁,除非已经确认当前令牌和持有者状态。强制删除可能让原持有者和新持有者同时执行。\n\n### 3. Lua 的检查\n\n脚本异常时检查:\n\n- 脚本是否包含错误的键数量;\n- Cluster 中是否跨槽;\n- 是否出现 `NOSCRIPT`;\n- 是否有长循环或一次处理过多成员;\n- 是否使用了可能阻塞的命令;\n- `SLOWLOG GET` 和延迟监控中是否出现长脚本。\n\n一个逻辑正确但执行几百毫秒的脚本,可能让所有普通请求都排队,因此脚本长度和数据规模必须受控。\n\n### 4. Pub/Sub 的检查\n\nPub/Sub 不提供历史查询接口。若怀疑消息丢失,无法通过 Redis 事后列出“某频道过去发布过什么”。应在应用侧记录:\n\n- 发布日志;\n- 连接建立和断开;\n- 订阅确认;\n- 消费者处理延迟;\n- 关键状态的周期性校验结果。\n\n如果系统需要凭 Redis 本身检查积压、重试或历史消息,Pub/Sub 的数据模型就不合适。\n\n### 5. Streams 的检查\n\n常用检查命令:\n\n```redis\nXLEN orders\nXINFO STREAM orders\nXINFO GROUPS orders\nXPENDING orders order-workers\nXINFO CONSUMERS orders order-workers\n```\n\n这些命令分别帮助确认:\n\n- Stream 中有多少条消息;\n- 首尾 ID 和长度等信息;\n- 消费者组的消费位置;\n- PEL 中是否有积压;\n- 哪个消费者持有多少待处理消息、空闲多久。\n\n一个典型故障判断过程是:\n\n```text\nStream 长度上升\n -> 检查组的最后消费 ID\n -> 检查消费者是否仍在线\n -> 检查 PEL 是否持续增长\n -> 检查最老消息的空闲时间\n -> 决定重启、认领、重试或转入死信 Stream\n```\n\nRedis Streams 没有自动为业务定义“死信队列”。若消息重试多次仍失败,应用可以将原消息、错误原因、重试次数写入另一个 Stream,并在确认原消息后进行人工或异步处理。\n\n---\n\n## 九、几个容易混淆的结论\n\n### “Redis 有原子命令,所以分布式锁绝对安全”\n\n不成立。原子命令解决的是同一 Redis 主节点上的并发插入问题;租约过期、客户端暂停、主从切换和外部资源写入仍需单独处理。\n\n### “加了过期时间就不会死锁”\n\n过期时间能处理持锁进程崩溃,但会引入租约过期后的旧持有者问题。任务时间不可预测时,续租和栅栏令牌比单纯增大 TTL 更重要。\n\n### “Pub/Sub 发布成功就代表消息被处理”\n\n不成立。`PUBLISH` 的返回值只是当时的订阅连接数量,不代表业务处理成功,更不代表未来可以重放。\n\n### “XACK 后消息就从 Redis 删除了”\n\n不成立。`XACK` 只更新消费者组的确认状态。Stream 中的消息仍可能存在,直到被裁剪或删除。\n\n### “Streams 自动提供恰好一次处理”\n\n不成立。它通常提供至少一次投递,业务必须幂等。数据库唯一键、条件更新、业务状态机和去重表都是常见的幂等基础。\n\n### “Lua 脚本等同于跨数据库事务”\n\n不成立。Lua 只原子地修改 Redis 内部状态,不能把 Redis 写入和 MySQL、消息网关或外部服务调用绑定为一个事务。\n\n---\n\nRedis 的协调能力可以归纳为不同的状态模型:\n\n```text\n锁:\n 一个键 + 随机令牌 + 租约\n 目标是互斥和故障后释放\n\nLua:\n 当前状态 + 参数 -> 原子状态转换\n 目标是消除读取与写回之间的竞态\n\n限流:\n 计数器或令牌桶状态 + 时间\n 目标是限制窗口内数量或长期速率\n\nPub/Sub:\n 当前在线订阅连接\n 目标是低延迟广播,不保存历史\n\nStreams:\n 追加消息 + 消费者组游标 + PEL\n 目标是可追踪、可恢复的消息处理\n```\n\n正确使用它们的关键,不是记住某条命令,而是先确定业务需要的状态是否持久、消息是否允许丢失、失败后是否必须重试,以及旧客户端在租约失效后是否仍可能写入资源。只有这些边界清楚后,锁、脚本、限流、广播和消息流才会各自处在适合的位置。\n\n---\n\n## 系列导航与关联阅读\n\n- 系列入口:[数据库完整学习路线:从关系模型、事务索引到分布式与向量检索](https://wrblog.cn/articles/e04c40d6-ba22-5c0c-8442-2252df05d216)\n- 上一篇:[Redis 缓存体系:一致性、穿透、击穿、雪崩和多级缓存](https://wrblog.cn/articles/02a05981-d76b-5aa8-9e41-fd6af20d978a)\n- 下一篇:[Elasticsearch Mapping 与分词:字段类型、Analyzer 和索引设计](https://wrblog.cn/articles/538f80fb-d74b-5495-943c-aa3fe5233bdf)\n- 延伸:[数据库事务完整指南:ACID、隔离级别、异常现象与正确边界](https://wrblog.cn/articles/edbe233d-f1b6-51b7-911d-4ce5ba6d6323)\n\n## 官方资料\n\n- [Redis Documentation](https://redis.io/docs/latest/)\n- [Redis Commands](https://redis.io/docs/latest/commands/)\n\n> 本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。\n","tags":["数据库","Redis","分布式锁","消息队列"],"likeCount":0,"commentCount":0,"createdByUserId":"10000000000","createdByDisplayName":"小郝","createdByAvatar":"/public/profile/10000000000/avatar/2026/08/04/db02b81c-42f2-441b-8a80-61370cdbb581.webp","publishTime":"2026-09-01 14:00:19","updateTime":"2026-09-01 14:00:19"}},"status":200,"locale":"zh-CN","theme":"light"}