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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 15:39:22