【问题标题】:Strange behavior of the ZeroMQ PUB/SUB pattern with TCP as transport layer使用 TCP 作为传输层的 ZeroMQ PUB/SUB 模式的奇怪行为
【发布时间】:2018-04-13 11:04:50
【问题描述】:

为了设计我们的 API/消息,我对我们的数据进行了一些初步测试:

Protobuf V3 消息:

message TcpGraphes {
    uint32 flowId                   = 1;
    repeated uint64 curTcpWinSizeUl = 2; // max 3600 elements
    repeated uint64 curTcpWinSizeDl = 3; // max 3600 elements
    repeated uint64 retransUl       = 4; // max 3600 elements
    repeated uint64 retransDl       = 5; // max 3600 elements
    repeated uint32 rtt             = 6; // max 3600 elements
}

消息构建为多部分消息,以便为客户端添加过滤器功能

使用 10 个 python 客户端进行测试:5 个在同一台 PC (localhost) 上运行,5 个在外部 PC 上运行。 使用的协议是 TCP。每秒发送大约 200 条消息

结果:

  1. 本地客户端正在工作:他们收到每条消息
  2. 远程客户端丢失了一些消息(吞吐量似乎被服务器限制为每个客户端 1Mbit/s)

服务器代码(C++):

// zeroMQ init
zmq_ctx = zmq_ctx_new();
zmq_pub_sock = zmq_socket(zmq_ctx, ZMQ_PUB);
zmq_bind(zmq_pub_sock, "tcp://*:5559");

每秒大约有 200 条消息循环发送:

std::string serStrg;
tcpG.SerializeToString(&serStrg);
// first part identifier: [flowId]tcpAnalysis.TcpGraphes
std::stringstream id;
id << It->second->first << tcpG.GetTypeName();
zmq_send(zmq_pub_sock, id.str().c_str(), id.str().length(), ZMQ_SNDMORE);
zmq_send(zmq_pub_sock, serStrg.c_str(), serStrg.length(), 0);

客户端代码(python):

ctx = zmq.Context()
sub = ctx.socket(zmq.SUB)
sub.setsockopt(zmq.SUBSCRIBE, '')
sub.connect('tcp://x.x.x.x:5559')
print ("Waiting for data...")
while True:
    message = sub.recv() # first part (filter part, eg:"134tcpAnalysis.TcpGraphes")
    print ("Got some data:",message)
    message = sub.recv() # second part (protobuf bin)

我们查看了 PCAP 并且服务器没有使用可用的全部带宽,我可以添加一些新订阅者,删除一些现有订阅者,每个远程订阅者“仅”获得 1Mbit/s。

我已经测试了两台 PC 之间的 Iperf3 TCP 连接,我达到了 60Mbit/s。

运行 python 客户端的 PC 最后拥有大约 30% 的 CPU。 为了避免打印输出,我已经最小化了客户端运行的控制台,但它没有任何效果。

这是 TCP 传输层(PUB/SUB 模式)的正常行为吗?这是否意味着我应该使用 EPGM 协议?

配置:

  • 用于服务器的 windows xp
  • windows 7 用于 python 远程客户端
  • 使用zmq 4.0.4版

【问题讨论】:

  • 单击 [+1] 以提供版本详细信息,(值得重新编辑帖子以包含完整的 MCVE-代码, 在服务器端和客户端正确的多部分消息处理。这个原始代码 sn-ps 违反了 Minimum-Complete-V可验证-E代码示例,可运行并自我演示所呈现的不良行为,以便与社区成员协商。
  • 刚刚添加了发送消息的代码(循环)

标签: windows zeromq


【解决方案1】:

表现出兴趣?

好的,让我们先更充分地利用资源:

// //////////////////////////////////////////////////////
// zeroMQ init
// //////////////////////////////////////////////////////

zmq_ctx = zmq_ctx_new();

int aRetCODE = zmq_ctx_set( zmq_ctx, ZMQ_IO_THREADS, 10 );

assert( 0 == aRetCODE );

zmq_pub_sock = zmq_socket(  zmq_ctx, ZMQ_PUB );

    aRetCODE = zmq_setsockopt( zmq_pub_sock, ZMQ_AFFINITY, 1023 );
    //                                                     ^^^^
    //                                                     ||||
    //                                 (:::::::::::)-------++++
    // >>> print ( "[{0: >16b}]".format( 2**10 - 1 ) ).replace( " ", "." )
    // [......1111111111]
    //        ||||||||||
    //        |||||||||+---- IO-thread 0
    //        ||||||||+----- IO-thread 1
    //        |......+------ IO-thread 2
    //        ::             :         :
    //        |+------------ IO-thread 8
    //        +------------- IO-thread 9
    //
    // API-defined AFFINITY-mapping

具有更新 API 的非 Windows 平台还可以触及调度程序细节并更好地调整 O/S 端的优先级。


网络?

好的,让我们先更充分地利用资源:

    aRetCODE = zmq_setsockopt( zmq_pub_sock, ZMQ_TOS, <_a_HIGH_PRIORITY_ToS#_> );

将整个基础架构转换为 epgm:// ?

好吧,如果有人希望进行实验并获得保证 E2E 的资源。

【讨论】:

    猜你喜欢
    • 2013-06-20
    • 1970-01-01
    • 2015-08-21
    • 1970-01-01
    • 2016-06-02
    • 1970-01-01
    • 1970-01-01
    • 2013-07-22
    • 2011-09-08
    相关资源
    最近更新 更多