Dragonfly 集群模式如何启用 Sharded Pub/Sub(SPUBLISH/SSUBSCRIBE)替代全局订阅?
2026/9/10 19:00:31 网站建设 项目流程

Dragonfly 集群模式如何启用 Sharded Pub/Sub(SPUBLISH/SSUBSCRIBE)替代全局订阅?

【免费下载链接】dragonflyA modern replacement for Redis and Memcached项目地址: https://gitcode.com/GitHub_Trending/dr/dragonfly

在 Dragonfly 中开启集群模式(--cluster_mode=yes)后,PUBLISHSUBSCRIBEPSUBSCRIBE这类全局发布订阅命令会被服务端直接拒绝;集群下唯一可用的发布订阅路径是 Sharded Pub/Sub(SPUBLISH/SSUBSCRIBE/SUNSUBSCRIBE)。本文在一个本地两节点集群上完成完整的操作路径:启动节点、推送 slot 配置、确认频道归属、完成一次订阅—发布—退订往返、检查订阅状态,并说明 slot 迁移时对已有订阅的处理。文中的命令和预期响应都可以对照仓库里的集群 pub/sub 集成测试复核。

集群模式下为什么只有 Sharded Pub/Sub 可用

pub-sub 文档把 Dragonfly 的 pub/sub 分成三种形态,并明确列出它们在集群模式下的行为:

形态命令作用域集群模式
StandardPUBLISHSUBSCRIBEUNSUBSCRIBE全局(所有频道)Blocked — 返回错误
PatternPSUBSCRIBEPUNSUBSCRIBE全局(glob 匹配)Blocked — 返回错误
ShardedSPUBLISHSSUBSCRIBESUNSUBSCRIBE按 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)。这样同一个频道的SPUBLISHSSUBSCRIBE总是路由到同一个 slot、同一个属主节点,集群的 slot 归属检查可以直接套用。

命令的元信息(来自 pub-sub.md 的命令注册表):

  • SPUBLISH:arity 3,参数为频道名 + 消息体;
  • SSUBSCRIBE:arity -2,订阅 1 个或多个频道;
  • SUNSUBSCRIBE:arity -1,不带参数时退订当前连接的全部 sharded 频道。

准备:获取二进制与集群启动参数

二进制可以按 Build From Source 从源码构建:安装构建依赖(Debian/Ubuntu 为ninja-buildlibunwind-devlibboost-context-devlibssl-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>真实集群模式必需。DFLYCLUSTERDFLYMIGRATE只在这个 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)、masterid为上一步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 kostas

KEYSLOT <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 号。这说明SSUBSCRIBESPUBLISH一样先过 slot 路由检查,必须连到属主节点执行。

第四步:SSUBSCRIBE → SPUBLISH → SUNSUBSCRIBE 完整往返

连接到频道属主节点(本例为节点 A)并订阅:

redis-cli -p 30001 SSUBSCRIBE kostas

SSUBSCRIBE是阻塞式订阅命令:连接建立后会先收到一条订阅确认推送,类型为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跨线程分发),仓库测试在SPUBLISHsleep(2)再取消息,用脚本验证时不要假设发布与接收严格同时发生;
  • smessage类型只用于 sharded 频道,与标准频道的message推送类型区分。

第五步:用 PUBSUB SHARDCHANNELS / SHARDNUMSUB 检查订阅状态

PUBSUB SHARDCHANNELS [pattern]列出当前活跃的 sharded 频道,PUBSUB SHARDNUMSUB <channel>...返回各频道的订阅数。仓库测试在订阅了pubsub-shard-channelshard-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 重平衡后,客户端需要自行向新属主节点重新订阅受影响频道,服务端不会代客迁移订阅关系。

边界条件与限制

  • 全局发布订阅在集群模式不可用:PUBLISHSUBSCRIBEPSUBSCRIBEPUNSUBSCRIBE均被拒绝(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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询