代码之家  ›  专栏  ›  技术社区  ›  ravi

ZeroMQ C++中的多发布服务器这是一个好选择吗?

  •  1
  • ravi  · 技术社区  · 8 年前

    我是ZeroMQ的新手。我想创建多个发布者,其中每个发布者发布特定数据,例如:

    1. Publisher 1:发布图像数据
    2. Publisher 2:发布音频数据
    3. Publisher 3:发布文本数据

    基本上,我的要求是从多个发布者发布数据,并在另一侧使用多个接收器接收。

    请参见下面的示例代码:

    data_发布者。清洁石油产品

    //  Prepare our context and all publishers
    zmq::context_t context(1);
    zmq::socket_t publisher1(context, ZMQ_PUB);
    zmq::socket_t publisher2(context, ZMQ_PUB);
    zmq::socket_t publisher3(context, ZMQ_PUB);
    zmq::socket_t publisher4(context, ZMQ_PUB);
    
    publisher1.bind("tcp://*:5556");
    publisher2.bind("tcp://*:5557");
    publisher3.bind("tcp://*:5558");
    publisher4.bind("tcp://*:5559");
    
    //  Initialize random number generator
    srandom((unsigned)time(NULL));
    while (1) {
        // sample data
        int zipcode1 = within(100000);
        int zipcode2 = within(100000);
        int zipcode3 = within(100000);
        int zipcode4 = within(100000);
    
        int temperature1 = within(215) - 80;
        int temperature2 = within(215) - 80;
        int temperature3 = within(215) - 80;
        int temperature4 = within(215) - 80;
    
        int relhumidity1 = within(50) + 10;
        int relhumidity2 = within(50) + 10;
        int relhumidity3 = within(50) + 10;
        int relhumidity4 = within(50) + 10;
    
        zmq::message_t message1(20);
        zmq::message_t message2(20);
        zmq::message_t message3(20);
        zmq::message_t message4(20);
    
        snprintf((char*)message1.data(), 20, "%05d %d %d", zipcode1, temperature1, relhumidity1);
        snprintf((char*)message2.data(), 20, "%05d %d %d", zipcode2, temperature2, relhumidity2);
        snprintf((char*)message3.data(), 20, "%05d %d %d", zipcode3, temperature3, relhumidity3);
        snprintf((char*)message4.data(), 20, "%05d %d %d", zipcode4, temperature4, relhumidity4);
    
        publisher1.send(message1);
        publisher2.send(message2);
        publisher3.send(message3);
        publisher4.send(message4);
    }
    

    data_接收器。清洁石油产品

    zmq::context_t context(1);
    
    //  Socket to talk to server
    zmq::socket_t subscriber1(context, ZMQ_SUB);
    zmq::socket_t subscriber2(context, ZMQ_SUB);
    zmq::socket_t subscriber3(context, ZMQ_SUB);
    zmq::socket_t subscriber4(context, ZMQ_SUB);
    
    subscriber1.connect("tcp://localhost:5556");
    subscriber2.connect("tcp://localhost:5557");
    subscriber3.connect("tcp://localhost:5558");
    subscriber4.connect("tcp://localhost:5559");
    
    const char* filter = (argc > 1) ? argv[1] : "10001 ";
    subscriber1.setsockopt(ZMQ_SUBSCRIBE, filter, strlen(filter));
    subscriber2.setsockopt(ZMQ_SUBSCRIBE, filter, strlen(filter));
    subscriber3.setsockopt(ZMQ_SUBSCRIBE, filter, strlen(filter));
    subscriber4.setsockopt(ZMQ_SUBSCRIBE, filter, strlen(filter));
    
    //  Process 100 updates
    int update_nbr;
    long total_temp1 = 0;
    long total_temp2 = 0;
    long total_temp3 = 0;
    long total_temp4 = 0;
    
    for (update_nbr = 0; update_nbr < 100; update_nbr++)
    {
        zmq::message_t update1;
        zmq::message_t update2;
        zmq::message_t update3;
        zmq::message_t update4;
        int zipcode1, temperature1, relhumidity1;
        int zipcode2, temperature2, relhumidity2;
        int zipcode3, temperature3, relhumidity3;
        int zipcode4, temperature4, relhumidity4;
    
        subscriber1.recv(&update1);
        subscriber2.recv(&update2);
        subscriber3.recv(&update3);
        subscriber4.recv(&update4);
    
        std::istringstream iss1(static_cast<char*>(update1.data()));
        std::istringstream iss2(static_cast<char*>(update2.data()));
        std::istringstream iss3(static_cast<char*>(update3.data()));
        std::istringstream iss4(static_cast<char*>(update4.data()));
    
        iss1 >> zipcode1 >> temperature1 >> relhumidity1;
        iss2 >> zipcode2 >> temperature2 >> relhumidity2;
        iss3 >> zipcode3 >> temperature3 >> relhumidity3;
        iss4 >> zipcode4 >> temperature4 >> relhumidity4;
    
        total_temp1 += temperature1;
        total_temp2 += temperature2;
        total_temp3 += temperature3;
        total_temp4 += temperature4;
    }
    
    std::cout << "Average temperature for zipcode '" << filter << "' was "
              << (int)(total_temp1 / update_nbr) << "F" << std::endl;
    std::cout << "Average temperature for zipcode '" << filter << "' was "
              << (int)(total_temp2 / update_nbr) << "F" << std::endl;
    std::cout << "Average temperature for zipcode '" << filter << "' was "
              << (int)(total_temp3 / update_nbr) << "F" << std::endl;
    std::cout << "Average temperature for zipcode '" << filter << "' was "
              << (int)(total_temp4 / update_nbr) << "F" << std::endl;
    

    请注意,以上代码是获取建议/建议的示例代码。

    我想知道它是否如上面的示例代码所示是一个好的选择?

    1 回复  |  直到 8 年前
        1
  •  2
  •   user3666197    8 年前

    仍在等待任何定量事实,但让我们开始:

    ZeroMQ是一种使用智能启用工具的概念,而底层系统编程被ZeroMQ核心元素隐藏 Context -发动机。

    也就是说,作为可扩展的正式通信模式原型的高级工具,提供了某种人类模仿行为-- PUB 出版商确实“发布”, SUB 用户可以“订阅”, REQ 请求者可以“请求”, REP 回复者确实可以“回复”,等等。

    这些具有行为的接入点可以 .bind()/.connect() 进入某种分布式行为基础设施,一旦被授予一些基本规则。其中一条规则是不必理会实际的运输类别,所有这些都是功能丰富的技术,目前跨越了 { inproc:// | ipc:// | tcp:// | pgm:// | epgm:// | vmci:// } ,低级细节 Context() -实例将处理所有这些 对你的高层行为透明。忘了这件事吧。另一条规则是,您可以确保发送的每封邮件要么没有错误,要么根本没有错误——在这方面没有任何妥协,没有任何令人痛苦的垃圾被交付给愚弄或破坏收件人的访问点后处理。

    如果不理解这一点,ZeroMQ就无法最大限度地为我们提供这种豪华工具所带来的舒适感和强大功能。

    回到你的困境:

    在说了上面的几句话之后,您的主要架构还不清楚,在这里仍然可以帮助您。

    ZeroMQ抽象的分布式行为套接字工具主要是一个- [SERIAL] { .bind() | .connect() } -与套接字关联的消息可以任意重新排序纯序列消息流。

    这意味着,在任何情况下- [CONCURRENT] 进程调度或在极端情况下- [PARALLEL] 流程调度在技术上是精心安排的,一个“纯”- [序列号] 传送通道不允许 { [CONCURRENT] | [PARALLEL] } -系统将继续提供此类流程调度模式,并将事件/处理流程分割为“纯”- [序列号] 消息序列。

    A) 确实如此 可能既是一个理由也是一个必须 用于引入多个独立操作的ZeroMQ分布式行为套接字实例。

    B) 另一方面,对全球分布式系统行为一无所知, 没有人能确定 ,进入多个独立操作的套接字实例是否不仅仅是浪费时间和资源,交付的数据不合理地低于平均水平,或者由于极端错误或完全缺少初始工程决策而导致的端到端系统行为性能不可接受。


    表演

    不要在这个领域猜测,永远不要。而是从第一个定量声明的需求开始,在此基础上,技术合理的设计将能够继续并定义资源映射和性能调整到平台极限所需的所有步骤。

    在过去的二十年里,ZeroMQ凭借其终极的性能特点和设计与性能,为实现这一目标做出了卓越的准备;工程团队在完善可扩展性和性能包络方面做了大量工作,同时将延迟保持在一个特定程序员难以实现的水平上。实际上,这是一个隐藏在ZeroMQ基础中的伟大的系统编程。

    " “--好的, 定义 大小 --交付 1E+9 消息具有 1 [B] 除了传递1E+3消息外,in size还有其他性能调整 1.000.000 [B] 在尺寸上。

    " 尽可能快 “--好的, 定义 快速的 对于 给定大小 信息的预期节奏为1/s~ 1 [Hz] ,10/s~ 10 [Hz] ,1000/s~ 1 [kHz]

    当然,在某些特定情况下,这种需求组合可能会偏离当代计算设备能力的范围。在任何编程开始之前,都必须对其进行最好的审查,因为否则你只会破坏一件事情上的一些编程工作,而这件事情永远不会成功,因此最好在可接受的资源和成本范围内,有一个解决方案架构可行和可行的积极证明。

    所以,如果你的项目需要什么, 第一 定义并定量说明实际情况, 下一个 解决方案架构可以开始对其进行分类,并提供决策、哪些工具和哪些工具配置可以与定义的功能和性能目标的目标级别相匹配。

    从盖屋顶开始盖房子,永远无法回答如何布置地下室墙的问题。铁混凝土铠装的厚度足够,但不会超过设计厚度,这将承载未知数量的高耸建筑地板。有一个已经建成的屋顶很容易出现,但与系统和严格的设计无关;工程实践。