You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 09:02:49