Redis:平级加权轮询 + Cluster 支持(v1.0,已落地)¶
Phase 2 第四刀:① 平级公平调度 ✅ 与 ② Cluster 支持 ✅(§5 的 C1–C4 按「建议」列拍板后实现)。 实现位置:
taskmq/transport/redis.py、taskmq/redis_client.py、tests/test_redis_cluster.py。
1. 问题¶
原取件规则是「全局最小 score」:score = -priority * 2**40 + seq —— 跨队列严格优先级、同级 FIFO。
这满足插队,但同级队列之间没有公平性:q1 堆 10 万条时 q2 可能饿死。memory/sqlite 早有档内加权轮询,
Redis 之前如实声明了 limitations["queue_weights"] —— 本刀实现掉。
2. 算法(与 memory/sqlite 一致)¶
每个队列取队头(ZRANGE ready:{q} 0 0 WITHSCORES)
band = 所有队头的最高优先级
候选 = 队头处于该 band 的队列
选 served/weight 最小者(平手按队列名,保证确定性)
INCR served:{chosen} -- 计数放 Redis,跨进程仍公平
- 原子性不变:单机整段在
_RESERVE_LUA一次往返里(权重读{prefix}weights、计数走served:{queue}); - 两种模式一致:无 Lua 回退路径(
_peek_candidate)实现同一规则 + 同样 tie-break,一致性套件两种模式都跑; - 权重来源:
Config.queues[name].weight→Transport.set_queue_weights(); - Cluster 下同一规则仍成立,只是「选择」挪到 Python 侧(§5.3)。
3. 权衡¶
- 每次取件多 2 条 Redis 命令(HGET weights / GET served),单机仍在一个 Lua 往返内;
- 公平性只在同一优先级档位内生效,不同优先级仍严格按优先级(方案 D);
- 档内轮询让「同级队列」不再互相饿死,但也不再是全局 FIFO —— 这是取舍点,不是 bug。
4. 验收(✅ 实测)¶
- 2 个同级队列、权重 3:1、各灌 100 条 → 取 40 条比例 30:10(±2)✅(Lua / 无 Lua / Cluster 三种跑法);
- 插队不受影响:低权重队列里的高优先级任务仍然先出 ✅;
limitations里的queue_weights已移除(不再声明降级)。
5. Cluster 支持(✅ 已落地)¶
为什么之前 cluster-safe 不了:Lua 里用前缀拼键(prefix .. "ready:" .. queue)。Cluster 要求一次
命令/Lua 用到的键在同一 slot;而且 Redis 7 不会在脚本里拦截跨槽访问——它会照写,于是数据静默
落在错误的节点上(比 CROSSSLOT 报错更危险)。所以键命名和校验都得自己兜住。
5.1 键布局:按逻辑队列分槽(C1:按逻辑队列 ✅)¶
结构 单机 Cluster 说明
ready taskmq:ready:q taskmq:{q}:ready 队头/取件
delayed taskmq:delayed:q taskmq:{q}:delayed 可见性(score = visible_at)
msg taskmq:msg:<seq> taskmq:{q}:msg:<seq> 消息 hash(Lua 里按 id 拼,必须同槽)
leases taskmq:leases:q taskmq:{q}:leases 租约 ZSET
prio taskmq:prio:q taskmq:{q}:prio 优先级计数
served taskmq:served:q taskmq:{q}:served 档内公平计数
weight taskmq:weights (HASH) taskmq:{q}:weight 每队列权重(Cluster 下拆成单键)
dlq taskmq:dlq:q taskmq:{q}:dlq 死信列表
全局单键(不参与同槽约束,按 slot 路由即可):seq / queues(SET) / jobs(ZSET) /
workers(ZSET) / worker:<id> / job:<id> / key:<idem> / nlease:<name> /
msgq(Cluster 专有:msg_id -> queue 反查索引,消息键带 hash tag 后没法只凭 id 定位)。
?cluster=1时才加 hash tag;默认单机键名一个字节都没变(老数据/老客户端不受影响);- prefix 与队列名在 Cluster 下不能含
{/}(会抢走或破坏 hash tag)→ 构造时报ConfigError; - 只支持 db 0(Cluster 不支持
SELECT),URL 写别的 db 直接报错。
5.2 Lua:显式 KEYS¶
单队列脚本全部改成显式 KEYS(Cluster 也能用,单机共用一份):
| 脚本 | KEYS | 用途 |
|---|---|---|
_MUTATE_LUA |
msg, leases, delayed, prio, dlq | ack/nack/dead/defer/extend/yield(单机与 Cluster 共用) |
_PEEK_QUEUE_LUA |
delayed, ready, served, weight | Cluster:promote 到期 + 返回队头/计数(不改状态) |
_CLAIM_QUEUE_LUA |
ready, leases, prio, served | Cluster:单队列原子取件(队头变了就 lost) |
_REAP_QUEUE_LUA |
leases, delayed, prio | Cluster:逐队列租约回收 |
单机跨队列脚本 _RESERVE_LUA / _REAP_LUA 保留原样(一次往返扫多队列);它们只在非 Cluster 模式使用,
所以允许按前缀拼 msg 键——Cluster 模式永远不会走到这里。
5.3 跨队列 reserve:降到 Python 侧(C2 ✅)¶
跨队列 = 跨 slot,没有单次原子可言。Cluster 下的 reserve([q1, q2, …]) 变成:
for 每个队列: PEEK(1 次 Lua:promote 到期 + 返回队头 + served/weight)
band = max(队头优先级)
候选 = band 内的队列;按 served/weight 升序、平手按队列名
for 候选: CLAIM(1 次 Lua:仍是原子的)→ ok 收下 / expired 清掉 / lost 换下一个
候选全 lost → 重新 PEEK(最多 3 轮,避免空转)
代价:一次取件是 1 + N 次往返(N = 队列数),且跨队列选择期间可能出现瞬时优先级倒挂
(先 peeking 的队列还没被 promote 到时,另一个队列的高优先级消息可能刚好到期)。
因此 cluster_limitations 里显式声明 global_priority 降级,一致性套件据此跳过该场景
(与 AMQP 同样处理)——但同一档位内的加权轮询与跨档严格优先在单进程视角下仍然成立,
另有 test_priority_still_beats_weights 等用例守着这个行为。
job 状态:job:<id> 是全局键(与队列不同槽),Lua 里碰不到 → 取件后由 Python 补 RUNNING、
过期由 Python 补 EXPIRED。崩溃窗口内 job 可能停在 QUEUED,租约回收后会重投(at-least-once 语义内)。
5.4 worker 表 / job 枚举(C3 ✅ 每 worker 一个键 + 维护期枚举)¶
worker:<id> 本来就是每 worker 一个键,指标挂在全局 workers ZSET;job:<id> + 全局
jobs ZSET 索引做 list_jobs(不用 SCAN —— Cluster 的 SCAN 只扫一个节点,会漏键)。
这些键都是单键操作,按 slot 路由即可,不需要维护期全量枚举。
5.5 开关(C4 ✅ ?cluster=1 显式开启)¶
不自动探测 CLUSTER INFO:单机行为保持不变、测试/生产行为不因环境漂移而变。写 cluster=1 才启用。
5.6 客户端:slot 路由(RedisClusterClient)¶
taskmq/redis_client.py 里与 RedisClient 同接口的集群实现:
- CRC16-XMODEM 算 slot + hash tag 规则(与
CLUSTER KEYSLOT实测对齐); CLUSTER SLOTS建拓扑(5s 缓存),单键命令按 slot 路由到对应主节点;MOVED就地更新映射并重试、ASK走一次性ASKING、CLUSTERDOWN/TRYAGAIN/LOADING有界重试;- 发命令前做同槽校验(一命令多键 / 一次 Lua 的 KEYS):不同槽直接报
TransportError,不让它变成"写错节点"; flush_prefix跨所有主节点 SCAN(单节点 SCAN 会漏),DEL 按槽分组(CROSSSLOT与"谁持有"无关);- 无 Lua 回退(
WATCH/MULTI/EXEC)在 Cluster 下同样可用:事务里的键全部同槽,跨槽的job:<id>挪到事务提交后写;MULTI里报错会DISCARD(否则这条连接后面只会拿到QUEUED)。
5.7 影响面(与 PG/AMQP 两刀相当)¶
RedisTransport 键命名加 hash tag(?cluster=1)、Lua 改显式 KEYS、reserve 降级路径 + 声明
global_priority 降级;一致性套件在 cluster 模式下按声明跳过 global_priority。
6. 测试环境与验收(✅)¶
make redis-cluster-up 起一个 3 主容器(7380-7382,--cluster-announce-ip 127.0.0.1 让 MOVED
返回宿主机可达地址),make test-redis-cluster 跑 tests/test_redis_cluster.py(21 个用例):
- 槽位算法(对
CLUSTER KEYSLOT实测值)、hash tag 分槽、两个队列确实落在不同主节点; - 跨槽取件 / ack / 租约回收 / DLQ 重放(含换队列搬槽:整份 hash 搬到目标槽 +
msgq索引 + 搬完可 ack); - 3:1 加权公平、优先级压过权重、4 线程 200 条无重复投递;
- 一致性套件
16个场景全跑,只有global_priority按声明跳过(Lua / 无 Lua 两种模式各一遍); App端到端(?prefix=…&cluster=1):20 个普通任务 + 1 个插队任务跑完。
7. 落地顺序¶
- ✅ 平级加权轮询(Lua / 无 Lua / Cluster 三路一致);
- ✅ Cluster:C1–C4 拍板 → 键命名 + Lua 显式 KEYS + 降级路径 + 客户端 slot 路由 + 集群容器测试。
8. 明确不做(留给后续)¶
- 自动探测
CLUSTER INFO(C4 备选):显式开关更可预测; - 从节点读(
READONLY+replicas参数):当前只用主节点; - 多 DB /
SELECT:Cluster 不支持,直接拒绝; - 跨队列取件的服务端原子化:除非 Redis 支持跨 slot 脚本,否则只能在客户端侧补偿。