【问题标题】:ZeroMQ SUB never receives messagesZeroMQ SUB 从不接收消息
【发布时间】:2017-08-24 05:41:26
【问题描述】:

我在使用 ZeroMQ 中的 PUB/SUB 时遇到问题。

连接所有内容后,发布者发布所有消息(套接字的发送消息返回true)但SUB 永远不会收到它们并在.recv() 函数上永远阻塞。

这是我正在使用的代码:

void startPublisher()
{
    zmq::context_t zmq_context(1);
    zmq::socket_t zmq_socket(zmq_context, ZMQ_PUB);
    zmq_socket.bind("tcp://127.0.0.1:58951");

    zmq::message_t msg(3);
    memcpy(msg.data(), "abc", 3);

    for(int i = 0; i < 10; i++)
        zmq_socket.send(msg); // <-- always true
}

void startSubscriber()
{
    zmq::context_t zmq_context(1);
    zmq::socket_t zmq_socket(zmq_context, ZMQ_SUB);

    zmq_socket.connect("tcp://127.0.0.1:58951");
    zmq_socket.setsockopt(ZMQ_SUBSCRIBE, "", 0); // allow all messages

    zmq::message_t msg(3);
    zmq_socket.recv(&msg); // <-- blocks forever (message never received?)
}

请注意,我在两个不同的线程中运行这 2 个函数,首先启动 SUB 线程,等待一段时间然后启动发布者线程(也尝试过其他方式,发布者在无限循环中发送消息,但是没用)。

我在这里做错了什么?

【问题讨论】:

  • 订阅者必须首先作为文档/示例的一部分运行,然后是发布者。在订阅者线程的开始和发布者线程的开始之间添加一个延迟,看看是否有区别。您可以在之后添加一些适当的同步。请像这样设置您的绑定:“tcp://*:58951”
  • 嗨,我刚刚设法解决了这个问题......这被称为“慢连接器”,详细描述了here。简而言之:TCP 握手需要一些时间才能完成,从而建立连接。但你是绝对正确的。
  • 这可以通过在 sub 开始和 pub 之间提供合适的同步机制来解决。
  • .context( nIoThreads ) 实例的底层有更多的任务,而不仅仅是上面提到的与 TCP 相关的任务。应该阅读 API,以及关于正确释放资源(套接字和上下文实例的终止),为什么应该始终为每个套接字实例设置一个 .setsockopt( zmq.LINGER, 0 ) 的预防值。值得一读,如果一个人是认真的去分发。 +并非所有 ZeroMQ 版本都以相同的方式操作主题过滤,因此如果一方订阅了某些内容,则过滤可能发生在 PUB 侧,对此更 "temportisation”需要..

标签: c++ zeromq


【解决方案1】:

根据您的示例,以下代码适用于我。 问题是 PUB / SUB 模式是一个缓慢的加入者,这意味着您需要在绑定 PUB 套接字后等待一段时间才能发送任何消息。

#include <thread>
#include <zmq.hpp>
#include <iostream>
#include <unistd.h>
void startPublisher()
{
    zmq::context_t zmq_context(1);
    zmq::socket_t zmq_socket(zmq_context, ZMQ_PUB);
    zmq_socket.bind("tcp://127.0.0.1:58951");
    usleep(100000); // Sending message too fast after connexion will result in dropped message
    zmq::message_t msg(3);
    for(int i = 0; i < 10; i++) {
        memcpy(msg.data(), "abc", 3);
        zmq_socket.send(msg); // <-- always true
        msg.rebuild(3);
        usleep(1); // Temporisation between message; not necessary
    }
}
volatile bool run = false;
void startSubscriber()
{
    zmq::context_t zmq_context(1);
    zmq::socket_t zmq_socket(zmq_context, ZMQ_SUB);
    zmq_socket.connect("tcp://127.0.0.1:58951");
    std::string TOPIC = "";
    zmq_socket.setsockopt(ZMQ_SUBSCRIBE, TOPIC.c_str(), TOPIC.length()); // allow all messages
    zmq_socket.setsockopt(ZMQ_RCVTIMEO, 1000); // Timeout to get out of the while loop
    while(run) {
        zmq::message_t msg;
        int rc = zmq_socket.recv(&msg);  // Works fine
        if(rc) // Do no print trace when recv return from timeout
            std::cout << std::string(static_cast<char*>(msg.data()), msg.size()) << std::endl;
    }
}
int main() {
    run = true;
    std::thread t_sub(startSubscriber);
    sleep(1); // Slow joiner in ZMQ PUB/SUB pattern
    std::thread t_pub(startPublisher);
    t_pub.join();
    sleep(1);
    run = false;
    t_sub.join();
}

【讨论】:

  • 从调整ZMQ_RCVTIMEO 开始是对适当设计师工作的逃避。有一个明确的非阻塞.recv()-call 语法和一个.poll()-instrumentation,因为不需要使用(主要是不好的)阻塞实践。专业级分布式计算系统不应该阻塞(由于许多明显的原因 - 代码在阻塞状态的整个持续时间内完全失控,可能无限等),所以不要复制任何教科书示例,其中代码设计器很容易,并且对于一些 SLOC 有一些空间限制。
  • 作为一名 RTOS 专家,这听起来应该很清晰,而且您的耳朵也能听到 :o)
  • 是的,分布式 PUB / SUB 系统有更好的模式。我只是试图强调缓慢的连接器属性,同时保持与 OP 的示例一样接近。你是绝对正确的:)。但是对于线程 PUB / SUB 的一个简单示例,它似乎已经足够了。
  • 可能也有兴趣,不仅慢连接是问题所在,还有主题过滤器处理的不确定性(在较新的版本中很难判断哪一方将处理该主题-过滤处理...)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-11-20
  • 1970-01-01
  • 1970-01-01
  • 2019-05-17
  • 1970-01-01
  • 1970-01-01
  • 2019-08-09
相关资源
最近更新 更多