Dragonfly 集群模式如何启用 Sharded Pub/Sub(SPUBLISH/SSUBSCRIBE)替代全局订阅?
【免费下载链接】dragonflyA modern replacement for Redis and Memcached项目地址: https://gitcode.com/GitHub_Trending/dr/dragonfly
在 Dragonfly 中开启集群模式(--cluster_mode=yes)后,PUBLISH、SUBSCRIBE、PSUBSCRIBE这类全局发布订阅命令会被服务端直接拒绝;集群下唯一可用的发布订阅路径是 Sharded Pub/Sub(SPUBLISH/SSUBSCRIBE/SUNSUBSCRIBE)。本文在一个本地两节点集群上完成完整的操作路径:启动节点、推送 slot 配置、确认频道归属、完成一次订阅—发布—退订往返、检查订阅状态,并说明 slot 迁移时对已有订阅的处理。文中的命令和预期响应都可以对照仓库里的集群 pub/sub 集成测试复核。
集群模式下为什么只有 Sharded Pub/Sub 可用
pub-sub 文档把 Dragonfly 的 pub/sub 分成三种形态,并明确列出它们在集群模式下的行为:
| 形态 | 命令 | 作用域 | 集群模式 |
|---|---|---|---|
| Standard | PUBLISH、SUBSCRIBE、UNSUBSCRIBE | 全局(所有频道) | Blocked — 返回错误 |
| Pattern | PSUBSCRIBE、PUNSUBSCRIBE | 全局(glob 匹配) | Blocked — 返回错误 |
| Sharded | SPUBLISH、SSUBSCRIBE、SUNSUBSCRIBE | 按 slot(频道名决定 slot) | 支持 |
在集群节点上执行全局命令会返回(error) PUBLISH is not supported in cluster mode yet。原因在于全局频道没有 slot 归属,而集群按 slot 做路由(cluster-mode.md §5.1)。Sharded Pub/Sub 把频道名当作 key 参与 slot 计算:slot(channel) = crc16(tag(channel)) & 0x3FFF,hash tag 规则与普通 key 相同(取第一段配平的{...}内容,否则取整个频道名,见 cluster-mode.md §3.1–3.2)。这样同一个频道的SPUBLISH和SSUBSCRIBE总是路由到同一个 slot、同一个属主节点,集群的 slot 归属检查可以直接套用。
命令的元信息(来自 pub-sub.md 的命令注册表):
SPUBLISH:arity 3,参数为频道名 + 消息体;SSUBSCRIBE:arity -2,订阅 1 个或多个频道;SUNSUBSCRIBE:arity -1,不带参数时退订当前连接的全部 sharded 频道。
准备:获取二进制与集群启动参数
二进制可以按 Build From Source 从源码构建:安装构建依赖(Debian/Ubuntu 为ninja-build、libunwind-dev、libboost-context-dev、libssl-dev等),git clone --recursive后执行./helio/blaze.sh -release,再在build-opt目录ninja dragonfly;也可以参考 Quick Start 使用 Docker 镜像。
集群模式下的启动参数(cluster-mode.md §2.2 与 §8):
| 参数 | 说明 |
|---|---|
--cluster_mode=yes | 进入完整集群模式。在收到第一份DFLYCLUSTER CONFIG之前,节点不持有任何 slot,数据面命令一律返回-ERR Cluster is not yet configured |
--admin_port=<p> | 真实集群模式必需。DFLYCLUSTER、DFLYMIGRATE只在这个 admin listener 上接受连接 |
--cluster_node_id=<id> | 可选。节点身份,必须与配置 JSON 中master.id一致 |
--cluster_announce_ip=<ip> | 可选。返回给客户端的 IP,用于CLUSTER SLOTS/SHARDS/NODES和 MOVED 回复 |
第一步:启动本地两节点集群
按仓库测试的端口约定(30001 起顺序递增),启动两个节点:
# 节点 A ./build-opt/dragonfly --port 30001 --admin_port 30002 --cluster_mode=yes # 节点 B ./build-opt/dragonfly --port 30003 --admin_port 30004 --cluster_mode=yes端口可换成任意不冲突的值。启动后在任意 listener(不限 admin 端口)上获取各节点身份:
redis-cli -p 30001 CLUSTER MYID redis-cli -p 30003 CLUSTER MYID第二步:构造并推送 slot 配置
集群拓扑由外部 cluster manager 编写并推送到每个节点。按 §4.1 的 JSON wire format 构造拓扑:每个 shard 包含slot_ranges(闭区间,互不重叠,全部 shard 合起来必须覆盖 0..16383)、master(id为上一步CLUSTER MYID的输出)、replicas(可为空数组)。下面这份示例把 16384 个 slot 对半分给节点 A 和节点 B;JSON 里两处 id 是占位符,请替换为各自CLUSTER MYID的实际输出:
[ { "slot_ranges": [ { "start": 0, "end": 8191 } ], "master": { "id": "<节点 A 的 CLUSTER MYID 输出>", "ip": "127.0.0.1", "port": 30001, "health": "online" }, "replicas": [] }, { "slot_ranges": [ { "start": 8192, "end": 16383 } ], "master": { "id": "<节点 B 的 CLUSTER MYID 输出>", "ip": "127.0.0.1", "port": 30003 }, "replicas": [] } ]同一份 JSON 必须推送到每个节点,集群内部没有 gossip。通过 admin 端口下发:
redis-cli -p 30002 DFLYCLUSTER CONFIG '<上面的 JSON>' redis-cli -p 30004 DFLYCLUSTER CONFIG '<上面的 JSON>'两个节点都返回OK表示配置生效。校验失败时(slot 有缺口或重叠、节点 id 重复等)返回-ERR Invalid cluster configuration.,且节点保留原配置不变,修正后重新推送即可。仓库的集成测试还有一个变体:不设置--admin_port时直接在主端口下发DFLYCLUSTER CONFIG,两种做法二选一即可。
验证方式:
redis-cli -p 30001 CLUSTER SLOTS返回的 slot 拓扑应体现 A 拥有 0..8191、B 拥有 8192..16383。此后,向不拥有某频道 slot 的节点发命令会得到-MOVED <slot> <ip>:<port>,端点是本地配置中该 slot 的属主。
第三步:确认频道的 slot 归属
订阅前先确认目标频道落在哪个节点上:
redis-cli -p 30001 CLUSTER KEYSLOT kostasKEYSLOT <key>是任意 listener 可用的只读查询命令。返回的 slot 号在[0, 8191]归节点 A,在[8192, 16383]归节点 B。本文示例沿用仓库测试使用的频道kostas;如果想控制频道落在哪个节点,可以借助 hash tag,只有{...}内的内容参与散列。
用 MOVED 行为直接验证路由:向不拥有该频道 slot 的节点执行订阅,例如频道归 A 所有时:
redis-cli -p 30003 SSUBSCRIBE kostas -MOVED <slot> 127.0.0.1:30001其中<slot>即CLUSTER KEYSLOT返回的 slot 号。这说明SSUBSCRIBE和SPUBLISH一样先过 slot 路由检查,必须连到属主节点执行。
第四步:SSUBSCRIBE → SPUBLISH → SUNSUBSCRIBE 完整往返
连接到频道属主节点(本例为节点 A)并订阅:
redis-cli -p 30001 SSUBSCRIBE kostasSSUBSCRIBE是阻塞式订阅命令:连接建立后会先收到一条订阅确认推送,类型为ssubscribe,末尾的数字是本次操作后剩余的订阅数。另开一个终端,在同一节点发布:
redis-cli -p 30001 SPUBLISH kostas hello订阅端连接随后收到一条类型为smessage的推送,内容为频道名 + 消息体。仓库测试用 redis-py 集群客户端(RedisCluster+pubsub().ssubscribe(...))验证的就是这两条推送,文档示例输出:
{"type": "ssubscribe", "pattern": None, "channel": b"kostas", "data": 1} {"type": "smessage", "pattern": None, "channel": b"kostas", "data": b"hello"}消息的投递路径见下图(来自 docs/pub-sub.md):SPUBLISH在全局ChannelStore中查找该频道的订阅者,再异步分发到各订阅者所在的 I/O 线程写出。
最后执行退订(不带参数则退订全部 sharded 频道):
redis-cli -p 30001 SUNSUBSCRIBE kostas连接收到一条sunsubscribe推送,测试中此时计数为 0。
两个执行细节:
- 投递是异步的(内部走
DispatchBrief跨线程分发),仓库测试在SPUBLISH后sleep(2)再取消息,用脚本验证时不要假设发布与接收严格同时发生; smessage类型只用于 sharded 频道,与标准频道的message推送类型区分。
第五步:用 PUBSUB SHARDCHANNELS / SHARDNUMSUB 检查订阅状态
PUBSUB SHARDCHANNELS [pattern]列出当前活跃的 sharded 频道,PUBSUB SHARDNUMSUB <channel>...返回各频道的订阅数。仓库测试在订阅了pubsub-shard-channel和shard-channel两个频道后的验证输出(文档示例):
redis-cli -p 30001 PUBSUB SHARDCHANNELS 1) "pubsub-shard-channel" 2) "shard-channel" redis-cli -p 30001 PUBSUB SHARDCHANNELS pubsub* 1) "pubsub-shard-channel" redis-cli -p 30001 PUBSUB SHARDNUMSUB pubsub-shard-channel shard-channel 1) "pubsub-shard-channel" 2) (integer) 1 3) "shard-channel" 4) (integer) 1这两个子命令只在集群模式下支持:非集群节点执行会返回PUBSUB SHARDCHANNELS is not supported in non cluster mode(见 src/server/main_service.cc)。
slot 迁移时:受影响频道的强制退订
当 cluster manager 推送新配置把 slot 迁走(迁移协议见 cluster-mode.md §6),原节点会调用UnsubscribeAfterClusterSlotMigration(src/server/channel_store.cc):收集受影响 slot 内所有频道的订阅者,整体移除,并向每个受影响的连接推送sunsubscribe(计数为 0)。迁移集成测试验证了这条链路:频道的 slot 迁到另一节点后,旧节点上的订阅者收到sunsubscribe;此时再到旧节点SSUBSCRIBE会被 MOVED 重定向到新属主。同样的清理也会发生在DFLYCLUSTER FLUSHSLOTS以及配置使 slot 离开属主时——节点在完成 slot 数据清扫后丢弃这些 slot 的 sharded 订阅。
也就是说:发生 slot 重平衡后,客户端需要自行向新属主节点重新订阅受影响频道,服务端不会代客迁移订阅关系。
边界条件与限制
- 全局发布订阅在集群模式不可用:
PUBLISH、SUBSCRIBE、PSUBSCRIBE、PUNSUBSCRIBE均被拒绝(PUBLISH的确切报错为PUBLISH is not supported in cluster mode yet)。从全局订阅迁移过来的应用,需要把所有调用路径改成 sharded 形态,且频道命名会开始参与 slot 路由。 SSUBSCRIBE一次订阅多个频道时,所有频道必须落在同一 slot,否则触发集群的单 slot 约束被拒绝(-CROSSSLOT,见 cluster-mode.md §5.1)。- 节点重启后不持有 slot,cluster manager 必须重新推送当前配置,节点才能恢复服务。
- 集群模式下存储为单 DB,
SELECT到其他 DB 索引会被拒绝。 - 发布端背压:单个 I/O 线程上订阅者排队字节数达到硬上限(
publish_buffer_limit的 4 倍)才会暂停发布者,细节见 docs/pub-sub.md 的 Backpressure 一节。
延伸阅读
- docs/cluster-mode.md — 集群命令面、配置 JSON 校验规则、迁移协议与失败模式;
- docs/pub-sub.md — ChannelStore、消息分发与背压的内部实现;
- tests/dragonfly/cluster_pubsub_test.py — 本文整条路径可运行的集成测试,包括 slot 迁移场景;
- docs/cluster-node-health.md — 配置中
health字段对客户端拓扑过滤的影响。
【免费下载链接】dragonflyA modern replacement for Redis and Memcached项目地址: https://gitcode.com/GitHub_Trending/dr/dragonfly
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考