在ZeroMQ的世界里,PUB/SUB是最经典的发布订阅模式,但它有一个明显局限:PUB套接字无法感知订阅者的存在,也无法读取订阅事件。当系统需要做订阅鉴权、订阅统计或者搭建代理转发层时,就要用到XPUB套接字。CZTop是Ruby生态中对libzmq和CZMQ的现代化封装,API简洁且贴近底层,用它来操作XPUB套接字非常顺手。本文将围绕订阅过滤、订阅事件管理以及高水位标记(HWM)三个核心话题,结合实际代码展开分析。

XPUB与PUB的区别:为什么要用XPUB
PUB套接字在ZeroMQ中是“哑巴”式的发布者,它只管往外发消息,完全不知道谁订阅了自己,订阅了什么主题。这在简单场景下没问题,但在稍复杂的系统里就成了短板。比如你想统计当前有多少订阅者订阅了stock-price这个主题,或者想在服务端做订阅白名单校验,PUB套接字无能为力。
XPUB套接字解决的就是这个问题:它把订阅者发出的SUBSCRIBE和UNSUBSCRIBE消息暴露给了应用层。也就是说,订阅消息本身会作为一条普通消息从XPUB套接字的接收管道里读出来,消息内容是一个字节的前缀(1表示订阅,0表示取消订阅)加上订阅过滤器字符串。应用层读到这些消息后,可以自行决定统计、记录甚至拒绝转发。
下面是用CZTop创建XPUB套接字并读取订阅事件的最小示例:
require 'cztop'
# 创建XPUB套接字并绑定端口
xpub = CZTop::Socket::XPUB.new
xpub.bind('tcp://127.0.0.1:5555')
loop do
message = xpub.receive
frame = message.frames.first
# 第一个字节是订阅标志:1=订阅,0=取消订阅
prefix = frame[0]
topic = frame[1..-1]
if prefix == 1
puts "有订阅者订阅了主题: #{topic}"
else
puts "有订阅者取消了订阅: #{topic}"
end
end
需要注意一个细节:默认情况下,XPUB套接字会把收到的订阅消息继续向底层传递(由libzmq自动处理订阅管理),同时也在接收队列里给应用一份拷贝。如果你想完全接管订阅管理,比如实现自定义的过滤逻辑,可以设置ZMQ_XPUB_VERBOSE选项,或者使用ZMQ_XPUB_MANUAL模式,在manual模式下libzmq不会自动应答订阅,你需要自己调用xpub.set_subscribe_message类似机制来确认订阅。CZTop对这两个选项都有封装,通过xpub.observe_socket_options或直接操作Zsock选项即可设置。
发布端订阅过滤与代理模式的实现
ZeroMQ的订阅过滤有两个层面。第一层在SUB端,订阅者用subscribe方法设置前缀过滤器,只有匹配前缀的消息才会被客户端接收。第二层在XPUB端,从3.x版本开始,libzmq支持发布端过滤:PUB/XPUB会根据收到的订阅信息,只为匹配的订阅者发送消息,这被称为上游匹配,能大幅减少网络传输量。
用XPUB搭配XSUB可以搭建一个代理(proxy),这是消息中转层的标准做法。想象这样一个场景:一组数据生产者分布在多台机器上,一群消费者在另一个网段,中间用一台代理机做汇聚分发。生产者连接到XSUB,消费者连接到XPUB,代理内部把两个套接字用管道连起来。
require 'cztop'
frontend = CZTop::Socket::XSUB.new
frontend.bind('tcp://10.0.0.1:5555') # 生产者连接这里
backend = CZTop::Socket::XPUB.new
backend.bind('tcp://10.0.0.1:5556') # 消费者连接这里
# 使用CZMQ内置的代理循环,内部自动转发订阅消息
CZTop::Proxy.new(frontend, backend).run
代理模式里有一个容易被忽视的坑:订阅消息的传递方向。消费者发往XPUB的订阅请求,必须被转发到XSUB端重新发送给上游生产者,否则生产者根本不知道下游订阅了什么,消息会被过滤掉。上面的代理实现已经自动处理了这件事,但如果你手写转发循环,就必须在收到以\x01或\x00开头的消息时,原样从XSUB端发出去。来看一个手写代理的核心逻辑:
require 'cztop'
frontend = CZTop::Socket::XSUB.new
frontend.bind('tcp://127.0.0.1:5555')
backend = CZTop::Socket::XPUB.new
backend.bind('tcp://127.0.0.1:5556')
poller = CZTop::Poller.new(frontend, backend)
loop do
poller.simple_poll do |socket|
message = socket.receive
# XSUB和XPUB之间直接透传,订阅消息会被自动处理
(socket == frontend ? backend : frontend) << message
end
end
这段代码之所以能工作,是因为XSUB套接字在发送以\x01开头的信息帧时,libzmq会把它解释为订阅命令传给上游,这是XPUB/XSUB设计的巧妙之处。同时别忘了给XPUB开启verbose选项,否则重复订阅同一个主题时,上游可能收不到第二次订阅请求。
高水位标记的配置与溢出行为分析
高水位标记(HWM,High Water Mark)是ZeroMQ防止内存爆炸的保护机制。每个套接字都有发送高水位(sndhwm)和接收高水位(rcvhwm),当待发送队列中的消息数超过sndhwm时,PUB/XPUB会直接丢弃新消息;而对于REQ、DEALER这类有可靠语义的套接字,超过HWM后会阻塞或返回EAGAIN。这就是为什么订阅者处理速度跟不上发布速度时,PUB/SUB会静默丢消息的原因。
在CZTop中设置HWM非常直观。CZMQ封装了set_sndhwm和set_rcvhwm这样的高层方法,默认值通常是1000条消息。来看具体用法:
require 'cztop'
xpub = CZTop::Socket::XPUB.new
xpub.options.sndhwm = 100_000 # 发送高水位设为10万条
xpub.options.rcvhwm = 10_000 # 接收高水位(订阅消息队列)
xpub.bind('tcp://127.0.0.1:5555')
# 建议同时在TCP层开启零拷贝缓冲,减少HWM溢出的概率
xpub.options.tcp_keepalive = 1
HWM的设置需要权衡内存与可靠性。假设单条消息1KB,sndhwm为100万意味着最坏情况下要缓存约1GB内存。一个实用经验是:发布频率高但容忍少量丢失的场景(如行情推送、日志收集),保持默认HWM即可,丢消息比崩掉进程好;而对可靠性要求高的场景,不要指望调大HWM解决问题,而应该在应用层引入确认机制,或者改用PUSH/PULL、ROUTER/DEALER这类具备重传基础的组合,再配合磁盘持久化。
判断是否发生了HWM溢出,可以周期性读取套接字的发送队列状态。libzmq提供了zmq_socket_monitor事件流,其中EVENT_MSG_DROPPED事件会在消息被丢弃时触发。CZTop可以通过CZTop::Monitor监听这些事件:
require 'cztop'
xpub = CZTop::Socket::XPUB.new
xpub.bind('tcp://127.0.0.1:5555')
monitor = CZTop::Monitor.new(xpub)
monitor.start
Thread.new do
monitor.each_event do |event|
if event.event == :msg_dropped
puts "消息被丢弃: #{event.address}"
end
end
end
另外一个相关参数是linger,它决定了套接字关闭时等待未发送消息flush的时间。默认值在libzmq中是-1(无限等待),这可能导致代理进程无法退出。实践中建议设置为合理的毫秒数,比如1000,配合HWM一起构成完整的背压策略。
一个完整示例:带订阅统计的发布服务
最后把前面的知识点串起来,实现一个能统计订阅者数量、按主题记录订阅情况的发布服务,并合理配置HWM。这个结构可以直接作为行情推送、事件总线等系统的骨架。
require 'cztop'
xpub = CZTop::Socket::XPUB.new
xpub.options.sndhwm = 50_000
xpub.options.linger = 1_000
xpub.bind('tcp://127.0.0.1:5555')
subscriptions = Hash.new { |h, k| h[k] = 0 }
poller = CZTop::Poller.new(xpub)
# 后台线程负责发布数据
publisher = Thread.new do
loop do
message = CZTop::Message.new(["market.USDJPY #{rand(100..120)}"])
xpub << message
sleep 0.1
end
end
# 主线程处理订阅事件
loop do
poller.simple_poll do |socket|
message = socket.receive
frame = message.frames.first.to_s
topic = frame[1..-1]
if frame[0] == "\x01"
subscriptions[topic] += 1
puts "订阅 #{topic},当前订阅数: #{subscriptions[topic]}"
else
subscriptions[topic] -= 1
puts "退订 #{topic},剩余订阅数: #{subscriptions[topic]}"
end
end
end
运行这个服务后,用任意SUB客户端订阅market前缀,就能在服务端看到订阅事件实时输出。这个例子还演示了一个重要技巧:发布循环和订阅事件处理放在不同线程,通过ZeroMQ线程安全的套接字语义避免锁竞争。需要说明的是,libzmq从4.x开始套接字不是严格线程安全的,更稳妥的做法是给两个线程各建一个套接字,或通过inproc管道在线程间传递订阅事件。
总结一下要点:XPUB的价值在于让订阅事件对应用可见,配合XSUB可以搭建灵活的代理架构;订阅过滤尽量利用发布端匹配来减少无效流量;HWM要结合消息大小、发布频率和业务容忍度综合设定,并配合监控手段及时发现消息丢弃。掌握这些,就能用CZTop构建出既高效又可控的发布订阅系统。