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

ZeroMQ反向PUB/SUB模式(SUB绑定、PUB连接)消息丢失求助

问题:ZeroMQ PUB/SUB角色互换后消息丢失问题

我在Python项目中用ZeroMQ实现IPC,需要一个进程接收其他进程的控制命令,因此采用了角色互换的PUB/SUB模式:一个进程用zmq.SUB绑定作为监听端,其他进程用zmq.PUB连接作为发送端。但PUB端发送的消息无法被SUB端全部接收,切换为tcp://传输也无改善,请问如何解决该问题?同时这种PUB作为连接端、SUB作为监听端的方式是否被允许?

相关代码

import zmq
import asyncio

IPC_SOCK = "/tmp/tmp.sock"

class DataObj:
    def __init__(self) -> None:
        self.value = 0

    def __str__(self) -> str:
        return f'DataObj.value: {self.value}'

async def server():
    context = zmq.Context()
    socket = context.socket(zmq.SUB)
    socket.bind(f'ipc://{IPC_SOCK}')
    socket.subscribe("")

    while True:
            try:
                obj = socket.recv_pyobj(flags=zmq.NOBLOCK)
                print(f'<-- {obj}')
                await asyncio.sleep(0.1)
            except zmq.Again:
                pass

            await asyncio.sleep(0.1)

async def client():
    print("Waiting for server to be come up")
    await asyncio.sleep(2)

    context = zmq.Context()
    socket = context.socket(zmq.PUB)

    socket.connect(f'ipc://{IPC_SOCK}')


    data_obj = DataObj()
    data_obj.value = 42

    print("Sending object once")
    socket.send_pyobj(data_obj)
    print(f"--> {data_obj}")
    print("Send object --> Not Received!")

    print("Sending object twice")
    for i in range(2):
        data_obj.value += 1
        socket.send_pyobj(data_obj)
        print(f"--> {data_obj}")
        await asyncio.sleep(0.1)

    print("Send both objects --> Received only once")
    

async def main():
    t_server = asyncio.create_task(server())
    t_client = asyncio.create_task(client())

    await t_client
    await t_server

   
if __name__ == "__main__":
    asyncio.run(main())

运行输出

Waiting for server to be come up
Sending object once
--> DataObj.value: 42
Send object --> Not Received!
Sending object twice
--> DataObj.value: 43
--> DataObj.value: 44
<-- DataObj.value: 44
Send both objects --> Received only once

解决方案

一、消息丢失的原因及修复方法

核心原因

  1. PUB/SUB握手延迟:PUB连接到SUB后,ZeroMQ需要短暂时间完成订阅关系的握手同步,这段时间内发送的消息会被直接丢弃。
  2. 低效的接收逻辑:原代码使用非阻塞接收+两次固定睡眠,导致SUB端每0.2秒才检查一次消息,极易错过PUB发送的消息。
  3. 未使用asyncio适配的ZeroMQ上下文:普通zmq.Context无法和asyncio的事件循环高效配合,影响消息接收的及时性。

修复后的代码

import zmq
import asyncio

IPC_SOCK = "/tmp/tmp.sock"

class DataObj:
    def __init__(self) -> None:
        self.value = 0

    def __str__(self) -> str:
        return f'DataObj.value: {self.value}'

async def server():
    # 使用适配asyncio的上下文
    context = zmq.asyncio.Context()
    socket = context.socket(zmq.SUB)
    socket.bind(f'ipc://{IPC_SOCK}')
    socket.subscribe("")

    while True:
        # 异步接收,避免轮询和遗漏
        obj = await socket.recv_pyobj()
        print(f'<-- {obj}')
        await asyncio.sleep(0.1)

async def client():
    print("Waiting for server to be come up")
    await asyncio.sleep(2)

    # 使用适配asyncio的上下文
    context = zmq.asyncio.Context()
    socket = context.socket(zmq.PUB)

    socket.connect(f'ipc://{IPC_SOCK}')
    # 等待PUB/SUB握手完成
    await asyncio.sleep(0.1)

    data_obj = DataObj()
    data_obj.value = 42

    print("Sending object once")
    socket.send_pyobj(data_obj)
    print(f"--> {data_obj}")

    print("Sending object twice")
    for i in range(2):
        data_obj.value += 1
        socket.send_pyobj(data_obj)
        print(f"--> {data_obj}")
        await asyncio.sleep(0.1)

async def main():
    t_server = asyncio.create_task(server())
    t_client = asyncio.create_task(client())

    await t_client
    await t_server

if __name__ == "__main__":
    asyncio.run(main())

关键修改点

  • 替换为zmq.asyncio.Context,让ZeroMQ和asyncio事件循环原生适配。
  • SUB端改用异步recv_pyobj(),无需非阻塞轮询,消息到达立即处理。
  • PUB连接后添加0.1秒延迟,确保订阅握手完成后再发送消息。

二、PUB连接、SUB绑定的方式是否允许?

完全允许。ZeroMQ的套接字角色(PUB/SUB等)与绑定/连接行为是解耦的,无论哪种角色都可以选择绑定或连接地址。你这种“中心SUB进程绑定地址,多个PUB进程连接发送”的架构,完全符合ZeroMQ的设计逻辑,是合理且被支持的用法。

内容的提问来源于stack exchange,提问作者seho85

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 00:00:06