【问题标题】:What is a simple example of a working XSUB / XPUB proxy in zeromq什么是 zeromq 中工作 XSUB / XPUB 代理的简单示例
【发布时间】:2017-07-23 07:33:30
【问题描述】:

我有后续How to implement Pub-Sub Network with a Proxy by using XPUB and XSUB in ZeroMQ(C++)?

该问题要求使用 XSUB 和 XPUB 的 C++ 代理。给出的答案本质上是下面引用的代理 main() 函数。

我将此代理扩展为一个完整的工作示例,包括发布者和订阅者。问题是我的代码仅适用于经销商/路由器选项(如下面的 cmets 所示)。使用下面的实际(未注释)XPUB / XSUB 选项,订阅者不会收到消息。怎么了?是否有任何调整可以让消息到达?

代理不适用于 XPUB/XSUB(在 cmets 中工作的经销商/路由器)

#include <zmq.hpp>

int main(int argc, char* argv[]) {
    zmq::context_t ctx(1);
    zmq::socket_t frontend(ctx, /*ZMQ_ROUTER*/ ZMQ_XSUB);
    zmq::socket_t backend(ctx, /*ZMQ_DEALER*/ ZMQ_XPUB);
    frontend.bind("tcp://*:5570");
    backend.bind("tcp://*:5571");
    zmq::proxy(frontend, backend, nullptr);
    return 0;
}

订阅者不使用 ZMQ_SUB(cmets 中的工作经销商/路由器选项)

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

std::string GetStringFromMessage(const zmq::message_t& msg) {
    char* tmp = new char[msg.size()+1];
    memcpy(tmp,msg.data(),msg.size());
    tmp[msg.size()] = '\0';
    std::string rval(tmp);
    delete[] tmp;
    return rval;
}

int main(int argc, char* argv[]) {
    zmq::context_t ctx(1);
    zmq::socket_t socket(ctx, /*ZMQ_DEALER*/ ZMQ_SUB);
    socket.connect("tcp://localhost:5571");
    while (true) {
        zmq::message_t identity;
        zmq::message_t message;
        socket.recv(&identity);
        socket.recv(&message);
        std::string identityStr(GetStringFromMessage(identity));
        std::string messageStr(GetStringFromMessage(message));
        std::cout << "Identity: " << identityStr << std::endl;
        std::cout << "Message: "  << messageStr  << std::endl;
    }
}

Publisher 不适用于 ZMQ_PUB(cmets 中的工作经销商/路由器选项)

#include <unistd.h>
#include <sstream>
#include <zmq.hpp>

int main (int argc, char* argv[])
{
    // Context
    zmq::context_t ctx(1);

    // Create a socket and set its identity attribute
    zmq::socket_t socket(ctx, /*ZMQ_DEALER*/ ZMQ_PUB);
    char identity[10] = {};
    sprintf(identity, "%d", getpid());
    socket.setsockopt(ZMQ_IDENTITY, identity, strlen(identity));
    socket.connect("tcp://localhost:5570");

    // Send some messages
    unsigned int counter = 0;
    while (true) {
        std::ostringstream ss;
        ss << "Message #" << counter << " from PID " << getpid();
        socket.send(ss.str().c_str(),ss.str().length());
        counter++;
        sleep(1);
    }
    return 0;
}

【问题讨论】:

  • 似乎是一个缓慢的加入者问题。我遇到了同样的问题 - 订阅者无法获取发布者消息,除非我在发送任何内容之前向发布者添加 sleep(1)。我还没有找到可行的解决方案。似乎 XPUB/XSUB 是一个损坏的构造
  • 感谢 RPGillespie 对此进行调查。
  • Dealer / Router 不像 XPUB / XSUB 那样做:每条消息都通过 Dealer / Router 路由到唯一的订阅者。因此,在我提出问题时,我没有向所有订阅者发送工作示例。我通过使用 ZMQ_PUSH 套接字来发布,并使用 ZMQ_PULL 套接字作为代理中的前端来解决 XPUB / XSUB 无法工作的事实。
  • 我能够使用 PUSH/PULL 获得一个不错的 XPUB/XSUB 解决方案,如下所示:stackoverflow.com/questions/43129714/…
  • “不工作”是什么意思?我能够在大约 20 分钟内代理一个有效的 PUB/SUB 示例(将发布者更改为具有多个发布线程,通过 XSUB/proxy/XPUB 联合)

标签: c++ zeromq


【解决方案1】:

代理示例

#include <zmq.hpp>

int main(int argc, char* argv[]) {

void* ctx = zmq_ctx_new();
assert(ctx);

void* frontend = zmq_socket(ctx, ZMQ_XSUB);
assert(frontend);
void* backend = zmq_socket(ctx, ZMQ_XPUB);
assert(backend);

int rc = zmq_bind(frontend, "tcp://*:5570");
assert(rc==0);
rc = zmq_bind(backend, "tcp://*:5571");
assert(rc==0);

zmq_proxy_steerable(frontend, backend, nullptr, nullptr);

zmq_close(frontend);
zmq_close(backend);

rc = zmq_ctx_term(ctx);
return 0;
}

发布示例

#include <zmq.hpp>
#include <bits/stdc++.h>

using namespace std;
using namespace chrono;

int main(int argc, char* argv[]) 
{
void* context = zmq_ctx_new();
assert (context);
/* Create a ZMQ_SUB socket */
void *socket = zmq_socket (context, ZMQ_PUB);
assert (socket);
/* Connect it to the host 

localhost, port 5571 using a TCP transport */
int rc = zmq_connect (socket, "tcp://localhost:5570");
assert (rc == 0);

while (true) 
{
    int len = zmq_send(socket, "hello", 5, 0);
    cout << "pub len = " << len << endl;
    this_thread::sleep_for(milliseconds(1000));
}
}

子示例

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

using namespace std;

int main(int argc, char* argv[]) 
{
void* context = zmq_ctx_new();
assert (context);
/* Create a ZMQ_SUB socket */
void *socket = zmq_socket (context, ZMQ_SUB);
assert (socket);
/* Connect it to the host localhost, port 5571 using a TCP transport */
int rc = zmq_connect (socket, "tcp://localhost:5571");
assert (rc == 0);
rc = zmq_setsockopt(socket, ZMQ_SUBSCRIBE, "", 0);
assert (rc == 0);

while (true) 
{
    char buffer[1024] = {0};
    int len = zmq_recv(socket, buffer, sizeof(buffer), 0);
    cout << "len = " << len << endl;
    cout << "buffer = " << buffer << endl;
}
}

【讨论】:

    【解决方案2】:

    在订阅者代码中,您尚未订阅接收来自发布者的消息。尝试添加该行:

    socket.setsockopt(ZMQ_SUBSCRIBE, "", 0); 
    

    行前/行后:

    socket.connect("tcp://localhost:5571");
    

    在您的订阅者代码中

    【讨论】:

      猜你喜欢
      • 2015-07-03
      • 2015-04-21
      • 1970-01-01
      • 1970-01-01
      • 2023-03-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多