【问题标题】:How can I both Send and Receive from a Router socket in ZeroMQ or NetMQ?如何在 ZeroMQ 或 NetMQ 中从路由器套接字发送和接收?
【发布时间】:2016-10-12 03:58:13
【问题描述】:

我在 NetMQ v4 中有一个经销商 路由器设置,我可以在任何方向异步发送和接收消息,没有任何问题。

我现在想将其形式化为一个抽象,其中服务器(Router)侦听任何传入消息,但它还需要按需将消息广播到任何连接的客户端(Dealers )。

我试图避免使用 Pub Sub 套接字,因为我需要订阅者也向服务器发送消息。与我尝试实现的最接近的模式是 WebSocket 客户端-服务器通信。

监听客户端消息的第一部分是这样完成的:

using (var server = new RouterSocket("@tcp://*:80"))
{
    var addresses = new HashSet<string>();
    while (true)
    {
        var msg = server.ReceiveMultipartMessage();

        var address = Encoding.UTF8.GetString(msg[0].Buffer);
        var payload = Encoding.UTF8.GetString(msg[2].Buffer);
        Console.WriteLine("[Server] - Client: {0} Says: {1}", address, payload);

        var contains = addresses.Contains(address);
        if (!contains) { addresses.Add(address); }            

        msg.Clear();
        msg.Append(address);
        msg.AppendEmptyFrame();
        msg.Append("Reply for: " + address);
        server.SendMultipartMessage(msg);
    }
}

现在考虑到套接字不是线程安全的,我一直在寻找一种方法来向所有客户端广播消息(根据需要来自不同的线程)。

我可能可以在循环中使用TryReceiveMultipartMessage 方法,而不是设置超时,之后我可以检查队列中是否有任何广播消息,然后遍历每个发送此类消息的客户端。比如:

using (var server = new RouterSocket("@tcp://*:80"))
{
    var addresses = new HashSet<string>();

    var msg = new NetMQMessage();
    while (true)
    {
        var clientHasMsg = server.TryReceiveMultipartMessage(TimeSpan.FromSeconds(1), ref msg);
        if (!clientHasMsg)
        {
            // Check any incoming broacast then loop through all the clients
            // sending each the brodcast msg
            var broadMsg = new NetMQMessage();
            foreach (var item in addresses)
            {
                broadMsg.Append(item);
                broadMsg.AppendEmptyFrame();
                broadMsg.Append("This is a broadcast");
                server.SendMultipartMessage(broadMsg);
                broadMsg.Clear();
            }

            // Go back into the loop waiting for client messages
            continue;
        }

        var address = Encoding.UTF8.GetString(msg[0].Buffer);
        var payload = Encoding.UTF8.GetString(msg[2].Buffer);
        Console.WriteLine("[Server] - Client: {0} Says: {1}", address, payload);

        var contains = addresses.Contains(address);
        if (!contains) { addresses.Add(address); }

        msg.Clear();
        msg.Append(address);
        msg.AppendEmptyFrame();
        msg.Append("Reply for: " + address);
        server.SendMultipartMessage(msg);
    }
}

这在某种程度上感觉不对,主要是因为:

  • 什么值的超时是一个好的值? 1 秒、100 毫秒等;
  • 这是最高效/性能最好的解决方案吗,因为该程序将用于连接 10 万多个客户端,每个客户端每秒发送数千条消息。

非常感谢任何关于什么是最好的方法的指针。

【问题讨论】:

    标签: c# .net zeromq netmq


    【解决方案1】:

    你可以使用netmqqueue,它是多生产者单消费者队列。您可以将其添加到 NetMQPoller 并从多个线程排队而无需锁定。

    【讨论】:

    • 我刚刚在Device 上阅读您的博客,并认为这将是一个不错的选择,直到您提到Queue :-) 您是否在某处有任何示例(除了@ 上的几行987654321@)?
    • 好的,我想我可以在没有其他示例的情况下完成这个工作,只需再问一个问题,NetMQQueue 和 NetMQSscheduler 有什么区别? v4 中的调度程序是否已过时?
    • NetMQScheduler 已过时(现在是 NetMQPoller 的一部分),无论如何 NetMQScheduler 是任务队列,NetMQQueue 是任何类型的队列。
    • 链接的示例显示每个接收就绪的单个出队,但套接字通知(至少)是边缘触发的。这是个问题吗?
    • OPs Q 直接指向 0MQ 的最小实际使用量,但我找不到使用队列或第二个套接字将传出消息传递给允许触摸原始套接字的单个线程的示例并且任何阻塞读取(或命中发送 HWM)意味着您无法关闭,因为您无法从另一个线程关闭套接字(或者如果您这样做,它不会解除阻塞)。那么...您能否提供一个使用 NetMQPoller、NetMQQueue、一些“有目的的”套接字以及将其全部关闭的方法的更完整示例?
    【解决方案2】:

    我认为 PUB/SUB 是满足 100k+ 客户需求的合适方法。然而,这并不意味着您不能与服务器通信:使用 DEALER/ROUTER。为什么您认为这种解决方案不可接受?

    【讨论】:

    • 我不确定我是否理解正确。您是说根据我的上述解决方案使用 Dealer/Router 进行 Pub/Sub 对您来说是合理的吗?
    猜你喜欢
    • 2010-11-14
    • 1970-01-01
    • 2011-07-11
    • 1970-01-01
    • 1970-01-01
    • 2016-02-21
    • 2013-04-01
    • 1970-01-01
    • 2018-09-24
    相关资源
    最近更新 更多