【问题标题】:Using zmq::proxy with REQ/REP pattern使用带有 REQ/REP 模式的 zmq::proxy
【发布时间】:2019-09-10 20:58:38
【问题描述】:

我试图了解 zmq::proxy 是如何工作的,但我遇到了问题:我希望将消息路由到正确的工作人员,但似乎身份和 evelopes 被忽略了:在例如,我想将消息从 client1 路由到 worker2,并将消息从 client2 路由到 worker1,但似乎消息是根据基于“第一个可用工作者”的规则提供的。 是我做错了什么,还是我误解了身份的工作原理?

#include <atomic>
#include <cassert>
#include <chrono>
#include <iostream>
#include <thread>
#include <mutex>

#include <zmq.hpp>
#include <zmq_addon.hpp>

using namespace zmq;
std::atomic_bool running;
context_t context(4);
std::mutex mtx;

void client_func(std::string name, std::string target, std::string message)
{
    std::this_thread::sleep_for(std::chrono::seconds(1));

    socket_t request_socket(context, socket_type::req);
    request_socket.connect("inproc://router");
    request_socket.setsockopt( ZMQ_IDENTITY, name.c_str(), name.size());

    while(running)
    {   
        multipart_t msg;
        msg.addstr(target);
        msg.addstr("");
        msg.addstr(message);

        std::cout << name << "sent a message: " << message << std::endl;
        msg.send(request_socket);
        multipart_t reply;
        if(reply.recv(request_socket))
        {
            std::unique_lock<std::mutex>(mtx);
            std::cout << name << " received a reply!" << std::endl;

            for(size_t i = 0 ; i < reply.size() ; i++)
            {
                std::string theData(static_cast<char*>(reply[i].data()),reply[i].size());
                std::cout << "Part " << i << ": " << theData << std::endl;
            }

        }

        std::this_thread::sleep_for(std::chrono::seconds(1));
    }

    request_socket.close();
}


void worker_func(std::string name, std::string answer)
{
    std::this_thread::sleep_for(std::chrono::seconds(1));

    socket_t response_socket(context, socket_type::rep);
    response_socket.connect("inproc://dealer");
    response_socket.setsockopt( ZMQ_IDENTITY, "resp", 4);

    while(running)
    {
        multipart_t request;

        if(request.recv(response_socket))
        {
            std::unique_lock<std::mutex>(mtx);

            std::cout << name << " received a request:" << std::endl;
            for(size_t i = 0 ; i < request.size() ; i++)
            {
                std::string theData(static_cast<char*>(request[i].data()),request[i].size());
                std::cout << "Part " << i << ": " << theData << std::endl;
            }

            std::string questioner(static_cast<char*>(request[0].data()),request[0].size());

            multipart_t msg;
            msg.addstr(questioner);
            msg.addstr("");
            msg.addstr(answer);

            msg.send(response_socket);
        }
    }

    response_socket.close();
}


int main(int argc, char* argv[])
{
    running = true;

    zmq::socket_t dealer(context, zmq::socket_type::dealer);
    zmq::socket_t router(context, zmq::socket_type::router);
    dealer.bind("inproc://dealer");
    router.bind("inproc://router");

    std::thread client1(client_func, "Client1", "Worker2", "Ciao");
    std::thread client2(client_func, "Client2", "Worker1", "Hello");
    std::thread worker1(worker_func, "Worker1","World");
    std::thread worker2(worker_func, "Worker2","Mondo");

    zmq::proxy(dealer,router);

    dealer.close();
    router.close();

    if(client1.joinable())
        client1.join();

    if(client2.joinable())
        client2.join();

    if(worker1.joinable())
        worker1.join();

    if(worker2.joinable())
        worker2.join();

    return 0;
}

【问题讨论】:

    标签: c++ c++11 zeromq


    【解决方案1】:

    来自docs

    当前端是 ZMQ_ROUTER 套接字,后端是 ZMQ_DEALER 套接字时,代理应充当共享队列,收集来自一组客户端的请求,并将这些请求公平地分配给一组服务。请求应从前端连接公平排队,并均匀分布在后端连接上。回复将自动返回给发出原始请求的客户端。

    代理处理多个客户端并使用多个工作人员来处理请求。身份用于将响应发送到正确的客户端。您不能使用标识来“选择”特定的工作人员。

    【讨论】:

      猜你喜欢
      • 2014-01-16
      • 1970-01-01
      • 1970-01-01
      • 2017-02-13
      • 1970-01-01
      • 1970-01-01
      • 2019-12-01
      • 1970-01-01
      • 2016-11-17
      相关资源
      最近更新 更多