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

ZeroMQ扩展PUB-SUB拓扑:无sleep避免发布者消息丢失的多节点方案

Great question—this is a super common pain point with ZeroMQ's PUB/SUB pattern, especially when scaling to multiple publishers and subscribers. That sleep(1) hack works for quick demos but is totally unreliable in production (what if connections take longer than a second? Or are faster, and you're wasting time?). Let's break down some robust, sleep-free solutions tailored to your extended topology:


1. Use ROUTER/DEALER Proxy with Connection Handshakes

ZeroMQ's built-in zmq_proxy() works seamlessly with ROUTER (for publishers) and DEALER (for subscribers) sockets, which let you add explicit connection confirmation logic. This is ideal for multi-publisher/subscriber setups because the proxy can uniquely identify each publisher and send targeted ready signals.

Proxy Side (C++)

#include <zmq.hpp>
#include <thread>
#include <iostream>

int main() {
    zmq::context_t context(1);
    zmq::socket_t frontend(context, ZMQ_ROUTER);  // Accepts publisher connections
    zmq::socket_t backend(context, ZMQ_DEALER);   // Connects to subscribers

    frontend.bind("tcp://*:5559");
    backend.bind("tcp://*:5560");

    // Run proxy in a separate thread so we can handle handshakes
    std::thread proxy_thread([&]() {
        zmq::proxy(frontend, backend, nullptr);
    });

    // Handle publisher handshakes
    zmq::message_t identity, empty_frame, ack_msg("READY", 5);
    while (true) {
        // Receive publisher's identity (ROUTER socket requirement)
        frontend.recv(identity);
        // Receive empty separator frame
        frontend.recv(empty_frame);
        // Receive "HELLO" from publisher
        zmq::message_t hello_msg;
        frontend.recv(hello_msg);

        // Send back ready acknowledgment
        frontend.send(identity, ZMQ_SNDMORE);
        frontend.send(empty_frame, ZMQ_SNDMORE);
        frontend.send(ack_msg);
        std::cout << "Sent ready signal to a new publisher" << std::endl;
    }

    proxy_thread.join();
    return 0;
}

Publisher Side (C++)

#include <zmq.hpp>
#include <iostream>

int main() {
    zmq::context_t context(1);
    zmq::socket_t pub_socket(context, ZMQ_PUB);
    zmq::socket_t req_socket(context, ZMQ_REQ);  // For handshake with proxy

    // Connect to proxy's ROUTER port for handshake
    req_socket.connect("tcp://localhost:5559");
    // Connect to proxy's DEALER port (we'll wait before publishing)
    pub_socket.connect("tcp://localhost:5560");

    // Initiate handshake
    zmq::message_t hello_msg("HELLO", 5);
    req_socket.send(hello_msg);

    // Wait for proxy's ready signal
    zmq::message_t ack_msg;
    req_socket.recv(ack_msg);
    std::cout << "Proxy confirmed connection, starting publication..." << std::endl;

    // Publish messages safely without sleep
    for (int i = 0; i < 10; ++i) {
        zmq::message_t msg("Hello from publisher!", 22);
        pub_socket.send(msg);
        std::cout << "Sent message " << i << std::endl;
    }

    return 0;
}

Why this works: The ROUTER socket lets the proxy track individual publishers, ensuring each one only starts sending after the proxy (and its connected subscribers) are fully ready.


2. Enable ZMQ_IMMEDIATE on PUB Sockets

This is a simpler, lightweight fix. The ZMQ_IMMEDIATE option modifies PUB socket behavior: instead of queuing messages when no subscribers are connected, send() will return -1 with errno = EAGAIN. You can use this to retry sending until subscribers are online.

Publisher Code Example

#include <zmq.hpp>
#include <iostream>
#include <errno.h>
#include <unistd.h>

int main() {
    zmq::context_t context(1);
    zmq::socket_t pub_socket(context, ZMQ_PUB);

    // Enable the immediate option
    int immediate_flag = 1;
    pub_socket.setsockopt(ZMQ_IMMEDIATE, &immediate_flag, sizeof(immediate_flag));

    pub_socket.connect("tcp://localhost:5560");

    // Send messages with retry logic
    for (int i = 0; i < 10; ++i) {
        zmq::message_t msg("Hello from publisher!", 22);
        while (true) {
            try {
                if (pub_socket.send(msg, ZMQ_DONTWAIT)) {
                    std::cout << "Sent message " << i << std::endl;
                    break;
                }
            } catch (zmq::error_t& e) {
                if (e.num() == EAGAIN) {
                    // No subscribers yet—wait a short interval and retry
                    usleep(100000);  // 100ms, far more efficient than sleep(1)
                    continue;
                }
                // Handle unexpected errors
                throw;
            }
        }
    }

    return 0;
}

Note: This is great if you don't need strict guarantees that all subscribers get the first message, but just want to avoid losing messages due to initial connection delays.


3. Subscriber-to-Publisher Subscription Acknowledgments

If you need to ensure all subscribers are ready before publishing any messages, have each subscriber send an acknowledgment to the publishers via a separate socket pair (PUSH/PULL works well here).

Subscriber Side (C++)

#include <zmq.hpp>
#include <iostream>

int main() {
    zmq::context_t context(1);
    zmq::socket_t sub_socket(context, ZMQ_SUB);
    zmq::socket_t push_socket(context, ZMQ_PUSH);  // For sending ack to publisher

    sub_socket.connect("tcp://localhost:5560");
    sub_socket.setsockopt(ZMQ_SUBSCRIBE, "", 0);  // Subscribe to all messages

    // Connect to publisher's ack listener
    push_socket.connect("tcp://localhost:5561");

    // Send subscription confirmation
    zmq::message_t ack_msg("SUBSCRIBED", 10);
    push_socket.send(ack_msg);
    std::cout << "Sent subscription ack to publisher" << std::endl;

    // Receive messages as usual
    while (true) {
        zmq::message_t msg;
        sub_socket.recv(msg);
        std::cout << "Received: " << std::string(static_cast<char*>(msg.data()), msg.size()) << std::endl;
    }

    return 0;
}

Publisher Side (C++)

#include <zmq.hpp>
#include <iostream>

int main() {
    zmq::context_t context(1);
    zmq::socket_t pub_socket(context, ZMQ_PUB);
    zmq::socket_t pull_socket(context, ZMQ_PULL);  // For receiving subscriber acks

    pub_socket.connect("tcp://localhost:5560");
    pull_socket.bind("tcp://*:5561");

    const int expected_subscribers = 2;
    int received_acks = 0;

    // Wait for all subscribers to confirm
    std::cout << "Waiting for " << expected_subscribers << " subscribers..." << std::endl;
    while (received_acks < expected_subscribers) {
        zmq::message_t ack_msg;
        pull_socket.recv(ack_msg);
        received_acks++;
        std::cout << "Received ack from subscriber " << received_acks << std::endl;
    }

    // Publish messages knowing all subscribers are ready
    std::cout << "All subscribers connected, starting publication..." << std::endl;
    for (int i = 0; i < 10; ++i) {
        zmq::message_t msg("Hello from publisher!", 22);
        pub_socket.send(msg);
        std::cout << "Sent message " << i << std::endl;
    }

    return 0;
}

Pros/Cons: This gives you strict control over publication timing, but requires tracking the number of expected subscribers and adds an extra socket pair to your architecture.


Which Approach Should You Choose?

  • For scalable multi-publisher/subscriber topologies with a proxy: Go with the ROUTER/DEALER handshake (Option 1) — it's clean and leverages ZeroMQ's native patterns.
  • For a quick, low-complexity fix: Use ZMQ_IMMEDIATE with retry (Option 2) — it avoids sleep without adding much code.
  • For strict delivery guarantees (all subscribers get all messages): Use the subscriber ack mechanism (Option 3) — it's explicit but adds some overhead.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:42:06