如何在ZeroMQ中无客户端时暂停PUSH消息发送?
解决方案
针对你的需求,这里提供几种可行的实现方式,可根据实际场景选择:
1. 利用PUSH套接字的阻塞特性(最简方案)
ZeroMQ的PUSH套接字默认是阻塞发送模式,当没有下游客户端连接或所有客户端队列满时,zmq_send会自动阻塞直到有可用客户端。如果只是想暂停消息发送(而非主动停止生成消息),直接使用默认的阻塞发送即可:
import zmq context = zmq.Context() controls_socket = context.socket(zmq.PUSH) controls_socket.bind("tcp://127.0.0.1:5556") is_running = True while is_running: # 生成控制命令逻辑... control_commands = "your_command_content" # 无客户端时会自动阻塞,直到客户端重新连接 controls_socket.send_string(control_commands)
如果需要主动暂停生成消息(而非卡在send环节),可以用非阻塞发送检测状态:
import zmq import time context = zmq.Context() controls_socket = context.socket(zmq.PUSH) controls_socket.bind("tcp://127.0.0.1:5556") is_running = True while is_running: try: # 用测试消息非阻塞发送,检测是否有可用客户端 controls_socket.send_string("ping", zmq.NOBLOCK) except zmq.Again: # 发送失败,无客户端/队列满,暂停生成消息 time.sleep(0.1) continue # 有客户端在线,生成并发送实际命令 control_commands = "your_command_content" controls_socket.send_string(control_commands)
2. 用ROUTER套接字主动检测连接数
ROUTER套接字支持获取当前连接的客户端数量,可主动判断是否有客户端在线,进而决定是否发送消息:
import zmq import time context = zmq.Context() # 替换PUSH为ROUTER,需Unity端用DEALER套接字配对 controls_socket = context.socket(zmq.ROUTER) controls_socket.bind("tcp://127.0.0.1:5556") is_running = True while is_running: # 获取当前连接的客户端总数 connection_count = controls_socket.getsockopt(zmq.CONNECTIONS) if connection_count == 0: # 无客户端连接,暂停发送 time.sleep(0.1) continue # 有客户端在线,生成并发送命令 control_commands = "your_command_content" # 单个客户端时可直接发送,多客户端需处理客户端ID controls_socket.send_string(control_commands)
注意:Unity端需使用NetMQ.DealerSocket与ROUTER配对通信,确保收发正常。
3. 心跳机制确认客户端在线状态
通过双向心跳包检测客户端是否存活,适合需要精准判断在线状态的场景:
Python端代码
import zmq import time context = zmq.Context() # 数据传输用ROUTER data_socket = context.socket(zmq.ROUTER) data_socket.bind("tcp://127.0.0.1:5556") # 心跳检测用REP heartbeat_socket = context.socket(zmq.REP) heartbeat_socket.bind("tcp://127.0.0.1:5557") heartbeat_socket.setsockopt(zmq.RCVTIMEO, 1000) # 1秒超时 client_online = False is_running = True while is_running: # 检测心跳 try: msg = heartbeat_socket.recv_string() if msg == "HEARTBEAT": heartbeat_socket.send_string("ACK") client_online = True except zmq.Again: client_online = False if not client_online: time.sleep(0.1) continue # 发送控制命令 control_commands = "your_command_content" data_socket.send_string(control_commands)
Unity端逻辑
定期向tcp://127.0.0.1:5557发送"HEARTBEAT"字符串,并等待ACK回复,确认服务器连接状态。
4. 离线时清空消息队列
如果需要在客户端离线时清空队列,可以通过重启套接字的方式实现:
import zmq import time context = zmq.Context() controls_socket = context.socket(zmq.PUSH) controls_socket.setsockopt(zmq.SNDHWM, 1) controls_socket.bind("tcp://127.0.0.1:5556") is_running = True last_valid_send = time.time() timeout = 3 # 3秒无有效发送则清空队列 while is_running: try: controls_socket.send_string("test", zmq.NOBLOCK) last_valid_send = time.time() except zmq.Again: if time.time() - last_valid_send > timeout: # 清空队列:关闭并重建套接字 controls_socket.close() controls_socket = context.socket(zmq.PUSH) controls_socket.setsockopt(zmq.SNDHWM, 1) controls_socket.bind("tcp://127.0.0.1:5556") last_valid_send = time.time() time.sleep(0.1) continue control_commands = "your_command_content" controls_socket.send_string(control_commands)
内容的提问来源于stack exchange,提问作者dennisklad
相关产品推荐
相关产品推荐

