【发布时间】: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 联合)