Windows下ZeroMQ(libzmq)C语言Pub/Sub模式过滤功能失效问题
ZeroMQ Pub/Sub 订阅过滤失效问题
背景说明
基于ZeroMQ的C语言Pub/Sub示例代码,在Windows环境(Visual Studio 2019编译器、vcpkg安装的zeromq:x64-windows 4.3.4#6)下可正常运行,订阅端能接收发布端发送的所有消息。
原正常运行代码
#include <Windows.h> #include <thread> #include <vector> #include <utility> #include "C:\Source\vcpkg\installed\x64-windows\include\zmq.h" int main(int argc, char const* argv[]) { DWORD thread_id = GetCurrentThreadId(); // 创建所有线程共享的上下文 void* context = zmq_ctx_new(); const char* endpoint = "tcp://127.0.0.1:4040"; std::thread publisher_thread = std::thread([&] { void* publisher = zmq_socket(context, ZMQ_PUB); int rc = zmq_bind(publisher, endpoint); //assert(rc == 0); while (1) { rc = zmq_send(publisher, "Hello World!", 12, 0); //assert(rc == 12); } zmq_close(publisher); }); // 简单延迟确保发布线程启动后订阅端再连接 Sleep(1000); std::thread subscriber_thread_all = std::thread([&] { void* subscriber = zmq_socket(context, ZMQ_SUB); int rc = zmq_connect(subscriber, endpoint); //assert(rc == 0); rc = zmq_setsockopt(subscriber, ZMQ_SUBSCRIBE, "", 0); //assert(rc == 0); char message[12]; while (1) { rc = zmq_recv(subscriber, message, 12, 0); //assert(rc != -1); printf("%s\n", message); } zmq_close(subscriber); }); subscriber_thread_all.join(); publisher_thread.join(); zmq_ctx_destroy(context); return 0; }
问题描述
修改代码后出现订阅过滤失效问题:让发布端交替发送"Hello World!"和"Goodbye World!",同时设置订阅端仅订阅以"H"开头的消息,但订阅端仍能接收所有消息,尝试使用消息信封等方式也无法解决该问题。
修改后过滤失效的代码
#include <Windows.h> #include <thread> #include <vector> #include <utility> #include "C:\Source\vcpkg\installed\x64-windows\include\zmq.h" int main(int argc, char const* argv[]) { DWORD thread_id = GetCurrentThreadId(); // 创建所有线程共享的上下文 void* context = zmq_ctx_new(); const char* endpoint = "tcp://127.0.0.1:4040"; std::thread publisher_thread = std::thread([&] { void* publisher = zmq_socket(context, ZMQ_PUB); int rc = zmq_bind(publisher, endpoint); //assert(rc == 0); int message_counter(0); while (1) { if (message_counter % 2 == 0) { rc = zmq_send(publisher, "Hello World!", 12, 0); } else { rc = zmq_send(publisher, "Goodbye World!", 12, 0); } ++message_counter; //assert(rc == 12); } zmq_close(publisher); }); // 简单延迟确保发布线程启动后订阅端再连接 Sleep(1000); std::thread subscriber_thread_all = std::thread([&] { void* subscriber = zmq_socket(context, ZMQ_SUB); int rc = zmq_connect(subscriber, endpoint); //assert(rc == 0); rc = zmq_setsockopt(subscriber, ZMQ_SUBSCRIBE, "H", 0); //assert(rc == 0); char message[12]; while (1) { rc = zmq_recv(subscriber, message, 12, 0); //assert(rc != -1); printf("%s\n", message); } zmq_close(subscriber); }); subscriber_thread_all.join(); publisher_thread.join(); zmq_ctx_destroy(context); return 0; }
内容的提问来源于stack exchange,提问作者kinar
相关产品推荐
相关产品推荐

