【问题标题】:ZeroMQ dealer--to-dealer high latency compared to winsock与 winsock 相比,ZeroMQ 经销商 - 经销商的高延迟
【发布时间】:2014-10-07 12:01:37
【问题描述】:

我的公司正在研究使用 ZeroMQ 作为传输机制。首先,我对性能进行了基准测试,以了解我正在玩什么。

因此,我创建了一个应用程序,将 zmq 经销商对经销商设置与 winsock 进行比较。我测量了从客户端发送同步消息到服务器的往返时间,然后计算平均值。

这里运行winsock的服务器:

DWORD RunServerWINSOCKTest(DWORD dwPort)
{
    WSADATA wsaData;
    int iRet = WSAStartup(MAKEWORD(2, 2), &wsaData);
    if (iRet != NO_ERROR)
    {
        printf("WSAStartup failed with error: %d\n", iRet);
        return iRet;
    }

    struct addrinfo hints;
    ZeroMemory(&hints, sizeof(hints));
    hints.ai_family = AF_INET;
    hints.ai_socktype = SOCK_STREAM;
    hints.ai_protocol = IPPROTO_TCP;
    hints.ai_flags = AI_PASSIVE;

    struct addrinfo *result = NULL;
    iRet = getaddrinfo(NULL, std::to_string(dwPort).c_str(), &hints, &result);
    if (iRet != 0)
    {
        WSACleanup();
        return iRet;
    }

    SOCKET ListenSocket = socket(result->ai_family, result->ai_socktype, result->ai_protocol);
    if (ListenSocket == INVALID_SOCKET)
    {
        freeaddrinfo(result);
        WSACleanup();
        return WSAGetLastError();
    }

    iRet = bind(ListenSocket, result->ai_addr, (int)result->ai_addrlen);
    if (iRet == SOCKET_ERROR)
    {
        freeaddrinfo(result);
        closesocket(ListenSocket);
        WSACleanup();
        return WSAGetLastError();
    }

    freeaddrinfo(result);
    iRet = listen(ListenSocket, SOMAXCONN);
    if (iRet == SOCKET_ERROR)
    {
        closesocket(ListenSocket);
        WSACleanup();
        return WSAGetLastError();
    }

    while (true)
    {
        SOCKET ClientSocket = accept(ListenSocket, NULL, NULL);
        if (ClientSocket == INVALID_SOCKET)
        {
            closesocket(ListenSocket);
            WSACleanup();
            return WSAGetLastError();
        }
        char value = 0;
        setsockopt(ClientSocket, IPPROTO_TCP, TCP_NODELAY, &value, sizeof(value));

        char recvbuf[DEFAULT_BUFLEN];
        int recvbuflen = DEFAULT_BUFLEN;
        do {

            iRet = recv(ClientSocket, recvbuf, recvbuflen, 0);
            if (iRet > 0) {
            // Echo the buffer back to the sender
                int iSendResult = send(ClientSocket, recvbuf, iRet, 0);
                if (iSendResult == SOCKET_ERROR)
                {
                    closesocket(ClientSocket);
                    WSACleanup();
                    return WSAGetLastError();
                }
            }
            else if (iRet == 0)
                printf("Connection closing...\n");
            else  {
                closesocket(ClientSocket);
                WSACleanup();
                return 1;
            }

        } while (iRet > 0);

        iRet = shutdown(ClientSocket, SD_SEND);
        if (iRet == SOCKET_ERROR)
        {
            closesocket(ClientSocket);
            WSACleanup();
            return WSAGetLastError();
        }
        closesocket(ClientSocket);
    }
    closesocket(ListenSocket);

    return WSACleanup();
}

这是运行winsock的客户端:

DWORD RunClientWINSOCKTest(std::string strAddress, DWORD dwPort, DWORD dwMessageSize)
{
    WSADATA wsaData;
    int iRet = WSAStartup(MAKEWORD(2, 2), &wsaData);
    if (iRet != NO_ERROR)
    {
        return iRet;
    }

    SOCKET ConnectSocket = INVALID_SOCKET;
    struct addrinfo *result = NULL,  *ptr = NULL, hints;


    ZeroMemory(&hints, sizeof(hints));
    hints.ai_family = AF_UNSPEC;
    hints.ai_socktype = SOCK_STREAM;
    hints.ai_protocol = IPPROTO_TCP;

    int iResult = getaddrinfo(strAddress.c_str(), std::to_string(dwPort).c_str(), &hints, &result);
    if (iResult != 0) {
        WSACleanup();
        return 1;
    }

    for (ptr = result; ptr != NULL; ptr = ptr->ai_next) {
        ConnectSocket = socket(ptr->ai_family, ptr->ai_socktype, ptr->ai_protocol);
        if (ConnectSocket == INVALID_SOCKET) {
            WSACleanup();
            return 1;
        }

        iResult = connect(ConnectSocket, ptr->ai_addr, (int)ptr->ai_addrlen);
        if (iResult == SOCKET_ERROR) {
            closesocket(ConnectSocket);
            ConnectSocket = INVALID_SOCKET;
            continue;
        }
        break;
    }

    freeaddrinfo(result);

    if (ConnectSocket == INVALID_SOCKET) {
        WSACleanup();
        return 1;
    }


    // Statistics
    UINT64 uint64BytesTransmitted = 0;
    UINT64 uint64StartTime = s_TimeStampGenerator.GetHighResolutionTimeStamp();
    UINT64 uint64WaitForResponse = 0;

    DWORD dwMessageCount = 1000000;

    CHAR cRecvMsg[DEFAULT_BUFLEN];
    SecureZeroMemory(&cRecvMsg, DEFAULT_BUFLEN);

    std::string strSendMsg(dwMessageSize, 'X');

    for (DWORD dwI = 0; dwI < dwMessageCount; dwI++)
    {
        int iRet = send(ConnectSocket, strSendMsg.data(), strSendMsg.size(), 0);
        if (iRet == SOCKET_ERROR) {
            closesocket(ConnectSocket);
            WSACleanup();
            return 1;
        }
        uint64BytesTransmitted += strSendMsg.size();

        UINT64 uint64BeforeRespone = s_TimeStampGenerator.GetHighResolutionTimeStamp();
        iRet = recv(ConnectSocket, cRecvMsg, DEFAULT_BUFLEN, 0);
        if (iRet < 1)
        {
            closesocket(ConnectSocket);
            WSACleanup();
            return 1;
        }
        std::string strMessage(cRecvMsg);

        if (strMessage.compare(strSendMsg) == 0)
        {
            uint64WaitForResponse += (s_TimeStampGenerator.GetHighResolutionTimeStamp() - uint64BeforeRespone);
        }
        else
        {
            return NO_ERROR;
        }
}

    UINT64 uint64ElapsedTime = s_TimeStampGenerator.GetHighResolutionTimeStamp() - uint64StartTime;
    PrintResult(uint64ElapsedTime, uint64WaitForResponse, dwMessageCount, uint64BytesTransmitted, dwMessageSize);

    iResult = shutdown(ConnectSocket, SD_SEND);
    if (iResult == SOCKET_ERROR) {
        closesocket(ConnectSocket);
        WSACleanup();
        return 1;
    }
    closesocket(ConnectSocket);
    return WSACleanup();
}

这里是运行ZMQ的服务器(经销商)

DWORD RunServerZMQTest(DWORD dwPort)
{
    try
    {
        zmq::context_t context(1);
        zmq::socket_t server(context, ZMQ_DEALER);

        // Set options here
        std::string strIdentity = s_set_id(server);
        printf("Created server connection with ID: %s\n", strIdentity.c_str());

        std::string strConnect = "tcp://*:" + std::to_string(dwPort);
        server.bind(strConnect.c_str());

        bool bRunning = true;
        while (bRunning)
        {
            std::string strMessage = s_recv(server);

            if (!s_send(server, strMessage))
            {
                return NO_ERROR;
            }
        }
    }
    catch (zmq::error_t& e)
    {
        return (DWORD)e.num();
    }

return NO_ERROR;

}

这里是运行ZMQ的客户端(dealer)

DWORD RunClientZMQTest(std::string strAddress, DWORD dwPort, DWORD dwMessageSize)
{
    try
    {
        zmq::context_t ctx(1);
        zmq::socket_t client(ctx, ZMQ_DEALER); // ZMQ_REQ

        // Set options here
        std::string strIdentity = s_set_id(client);

        std::string strConnect = "tcp://" + strAddress + ":" + std::to_string(dwPort);
        client.connect(strConnect.c_str());

        if(s_send(client, "INIT"))
        {
            std::string strMessage = s_recv(client);
            if (strMessage.compare("INIT") == 0)
            {
                printf("Client[%s] connected to: %s\n", strIdentity.c_str(), strConnect.c_str());
            }
            else
            {
                return NO_ERROR;
            }
        }
        else
        {
            return NO_ERROR;
        }


        // Statistics
        UINT64 uint64BytesTransmitted   = 0;
        UINT64 uint64StartTime          = s_TimeStampGenerator.GetHighResolutionTimeStamp();
        UINT64 uint64WaitForResponse    = 0;

        DWORD dwMessageCount = 10000000;


        std::string strSendMsg(dwMessageSize, 'X');
        for (DWORD dwI = 0; dwI < dwMessageCount; dwI++)
        {
            if (s_send(client, strSendMsg))
            {
                uint64BytesTransmitted += strSendMsg.size();

                UINT64 uint64BeforeRespone = s_TimeStampGenerator.GetHighResolutionTimeStamp();
                std::string strRecvMsg = s_recv(client);
                if (strRecvMsg.compare(strSendMsg) == 0)
                {
                    uint64WaitForResponse += (s_TimeStampGenerator.GetHighResolutionTimeStamp() - uint64BeforeRespone);
                }
                else
                {
                    return NO_ERROR;
                }
            }
            else
            {
                return NO_ERROR;
            }
        }
        UINT64 uint64ElapsedTime = s_TimeStampGenerator.GetHighResolutionTimeStamp() - uint64StartTime;
        PrintResult(uint64ElapsedTime, uint64WaitForResponse, dwMessageCount, uint64BytesTransmitted, dwMessageSize);
    }
    catch (zmq::error_t& e)
    {
        return (DWORD)e.num();
    }

    return NO_ERROR;
    }

我在本地运行基准测试,消息大小为 5 个字节,我得到以下结果:

WINSOCK

Messages sent:                 1 000 000
Time elapsed (us):            48 019 415
Time elapsed (s):                     48.019 415
Message size (bytes):                  5
Msg/s:                            20 825
Bytes/s:                         104 125
Mb/s:                                  0.099
Total   response time (us):   24 537 376
Average repsonse time (us):           24.0

和

ZeroMQ

Messages sent:                 1 000 000
Time elapsed (us):           158 290 708
Time elapsed (s):                    158.290 708    
Message size (bytes):                  5
Msg/s:                             6 317
Bytes/s:                          31 587
Mb/s:                                  0.030
Total   response time (us):  125 524 178    
Average response time (us):          125.0

谁能解释为什么使用 ZMQ 时平均响应时间要长得多?

我们的目标是找到一种设置,让我可以异步发送和接收消息而无需回复。如果这可以通过与经销商-经销商不同的设置来实现,请告诉我!

【问题讨论】:

  • 我快速浏览了 0MQ 网站,发现主要是商业支持和错误跟踪,这个问题似乎不适合(特别是因为您仍在调查)。这是一个棘手的问题,因为在套接字顶部添加一些堆栈预计会增加延迟,但这差不多。测试的主要区别是使用 std::string 可能会有所不同(但不是那么大)。此外,请确保您在发布模式下启用了编译器优化标志,以确保测试中不涉及额外检查。
  • 我对两种实现都使用了相同的编译器优化标志,所以这应该是个问题。 std::string 只创建一次,然后数据被引用多次,所以不能这样。感谢您的输入
  • 有一些测试将动态数组与向量进行比较,并显示使用相同优化标志的向量要慢得多,但事实证明,使用这些标志对未使用的向量进行了额外检查开启优化的发布模式。服务器和客户端每次都会收到一个字符串,然后由服务器发回。即使字符串在堆栈上,它的实际数据也可以是动态的。
  • 请注意,我不希望 std::string 与堆栈处理和网络流量相比有很大的不同。旗帜可以。顺便说一句,您是否已经使用wireshark 之类的网络监视器查看了流量?
  • 您还提到了使用wireshark,这有什么说明吗?排队消息是可能的。您需要具有全面 ZeroMQ 知识的人来验证这一点。底线是,当然对于简单的情况,添加堆栈和库通常会增加延迟,即使您测试的内容似乎超出预期。主要问题是 1. 额外的延迟是否可以接受? 2. 更轻松的交通处理带来的额外好处值得额外的延迟吗? 3. 现实生活中的延迟差异更好吗?

标签: c++ visual-c++ zeromq winsock low-latency


【解决方案1】:

这只是对您问题的一小部分的回答,但这里是 -

为什么需要经销商/经销商?我假设是因为通信可以从任一点开始?您没有束缚到经销商/经销商,特别是它限制您只能使用两个端点,如果您曾经在通信的任一侧添加另一个端点,例如第二个客户端,那么每个客户端都会只收到一半的消息,因为庄家是严格循环的。

异步通信所需的是经销商和/或路由器套接字的某种组合。两者都不需要响应,主要区别在于它们如何选择向哪个连接的对等方发送消息:

  • 如前所述,经销商是严格循环的,它将发送给串联的每个连接的对等方
  • 路由器严格来说是一个寻址消息,您必须知道要发送到的对等方的“名称”才能从那里获取消息。

这两种套接字类型可以一起工作,因为经销商套接字(和请求套接字,经销商是“请求类型”套接字)将它们的“名称”作为消息的一部分发送,路由器套接字可以使用该消息将数据发回。这是一个请求/回复范式,您会在the guide 中的所有示例中看到这种范式,但您可以将该范式转变为您正在寻找的内容,特别是经销商和路由器都不需要回复。

在不了解您的全部要求的情况下,我无法告诉您我会选择哪种 ZMQ 架构,但总的来说,我更喜欢路由器套接字的可扩展性,处理适当的寻址比将所有内容硬塞到单个对等体中更容易...你会看到不要做路由器/路由器的警告,我同意他们的观点,你应该在尝试之前了解你在做什么,但是了解你在做什么,实现不是那样很难。


如果它符合您的要求,您还可以选择为每一端设置一个 pub 套接字,如果实际上没有回复 永远,则每个端都有一个子套接字。如果它严格来说是从源到目标的数据馈送,并且对等点都不需要对其发送的内容有任何反馈,那么这可能是最佳选择,即使这意味着您每端处理两个套接字而不是一个。


这些都不能直接解决性能问题,但要了解的重要一点是 zmq 套接字针对特定用例进行了优化,正如 John Jefferies 的回答中所指出的那样,您正在打破经销商套接字的用例测试中的消息传递严格同步。首先要确定您的 ZMQ 架构,然后模拟实际的消息流,特别是不添加任意等待和同步性,这将必然改变方式吞吐量看起来就像您正在测试它,几乎按照定义。

【讨论】:

  • 感谢您的全面回复!我正在进行的设计必须能够处理双向的异步消息传递。我有一个可行的解决方案,但性能比我的 winsock 解决方案 (IOCP) 差大约 50%。这促使我深入实施,以找到导致我进行此单元测试的潜在瓶颈。我问这个问题的原因是为了得到一些关于做什么和不做什么的指导方针。我会听取您的建议并研究路由器路由器设计,希望它会产生更好的性能!再次感谢
  • 我不一定认为路由器/路由器会产生更好的性能,特别是如果您的通信保证是一对一的,那么经销商/经销商将很多 更简单,如果保证多对一,那么经销商/路由器也将更更简单。如果它可能是多对多的,那么路由器/路由器可能有意义,或者多插槽设计可能更有意义,具体取决于您的具体情况。无论哪种情况,关于经销商和路由器,您需要记住的重要一点是它们针对异步通信进行了优化,尽量不要破坏这一点
  • 我明白你的意思,不过还有一件事。我也尝试了相同的应用程序,但使用了 REQ-REP(更适合同步传输,对吗?),但它也没有产生任何改进。您对此有何看法,或者用于同步传输的合适套接字类型是什么?再次感谢您的意见
  • 我一点也不惊讶 REQ-REP 在该测试中提供相同的性能。 REQ/REP 套接字类型基于 DEALER/ROUTER,但本质区别在于它们强制同步消息传递。由于您的测试代码已经强制执行同步消息传递,因此性能将是相同的。通常,应用程序会希望为某些功能或其他功能进行异步消息传递;纯 REQ/REP 用于培训/初学者:-)。
  • 是的,正如@JohnJefferies 所说,性能问题源于消息传递模式,而不是套接字类型。通过使用异步消息传递,您只会看到 ZMQ 性能的真正限制是什么。我谈论套接字类型的原因是因为某些套接字类型不会妨碍它。根据定义,同步消息更关注对话的结构(request-reply-request-reply),而不是它的速度。使用不那么自以为是的协议(如 winsock),您将保持更高水平的速度,但它仍然无法达到异步的可能。
【解决方案2】:

你说你想异步发送和接收消息而不需要回复。然而到目前为止所做的测试都是完全同步的,本质上是请求-回复,但在经销商-经销商套接字上。那里没有计算出来的东西。为什么不运行更接近您的目标设计的测试?

ZeroMQ 通过将排队的消息聚合成一条消息,获得了相当多的“比 TCP 更快”的性能。显然,该机制不能在一次仅发送一条消息的纯同步设计中激活。

至于为什么这个非常小的消息完全同步发送和接收的特殊测试相对较慢,我不能说。你做过剖析吗?我要再说一遍,如果它看起来不像你的最终设计,那么运行这个测试并根据它做出决定是没有意义的。

ZeroMQ 代码中的 try/catch 块确实看起来很奇怪。这看起来不公平,因为winsock 测试不是这样写的。众所周知,在 try/catch 中有相当数量的开销。

【讨论】:

  • 我执行这个特定测试的原因是因为我必须使用经销商-经销商设置来进行异步消息传递(据我所知),但我想弄清楚开销。我们有一个“有效”的解决方案,我们异步发送消息,但我们在那里也得到了糟糕的性能,所以这就是我做低级测试的原因。我将尝试删除 try/catch,我目前正在使用wireshark 进行分析和查看流量。
  • 去掉了try/catch,没什么区别
【解决方案3】:

OP 问题是吞吐量问题,而不是延迟问题,并且可能与所提供示例中使用的模式有关。但是,您可能总是会发现 ZeroMQ 具有更高的延迟,我将对此进行解释,尽管在这种情况下它可能对 OP 没有用处。

ZeroMQ 通过缓冲消息工作。想象一下(只是作为一个基本说明)创建一个 std::string 并向其附加许多小字符串(数千个,每个都包括一个小标题以了解这些小段的大小),然后发送更大的字符串间隔为 100us、1000us、10ms 或其他。在接收端,接收大字符串,并根据与它一起发送的大小标头一次删除每个较小的消息。这使您可以分批发送数百万条消息(尽管std::string 显然是一个糟糕的选择),而无需一次发送数百万条非常小的消息的开销。因此,您可以充分利用网络资源并提高吞吐量,还可以创建基本的FIFO 行为。但是,您还创建了一个延迟以允许缓冲区填充,这意味着延迟增加。

想象一下(再次,仅作为基本说明):如果您花费半秒(包括字符串操作等)缓冲一百万条消息,这将导致几兆字节的较大字符串。现代网络可以在剩下的半秒内轻松发送这个更大的字符串。 1000000us(1 秒)/1000000 条消息将是 1us 每条消息,对吗?错误 - 所有消息都有半秒延迟以允许队列填满,导致所有消息的延迟最多增加半秒。 ZeroMQ 发送批次的速度比每个 500ms 快得多,但这说明延迟的增加仍然存在于 ZeroMQ 中,尽管它通常与 ms 类似。

【讨论】:

    猜你喜欢
    • 2016-08-23
    • 1970-01-01
    • 2012-02-26
    • 2010-09-25
    • 1970-01-01
    • 2011-03-21
    • 1970-01-01
    • 2011-12-02
    • 1970-01-01
    相关资源
    最近更新 更多