【问题标题】:What causes cases with high ZeroMQ latency and how to avoid them?是什么导致 ZeroMQ 延迟高的情况以及如何避免它们?
【发布时间】:2020-09-23 05:49:10
【问题描述】:

我尝试使用 ZeroMQ 进行快速消息传递。消息需要在 1 [ms] 内送达。我做了一些测试(inproc,Linux 上的单进程,没有 TCP),发现通常没有问题。延迟约为10 - 100 [us],具体取决于发送消息的频率(为什么?)。但是,有时会在6 [ms] 之后收到消息,这是不可接受的。

某些消息延迟的原因可能是什么?

也许进程被抢占了?

还是因为使用了轮询 (zmq_poll())?

我的测试结果示例:

avg lag =    28    [us]
max lag =  5221    [us]
std dev =    25.85 [us]
big lag =   180    x above 200 [us]

“大滞后”表示延迟超过200 [us] 的情况数。在我的测试中,发送了 500 000 条消息,因此值 180 意味着 200 [us] 上的延迟记录在 180 / 500000 = 0,036% 中。这是一个很低的数字,但我希望它为零。即使以平均延迟为代价。

测试源码如下:

#include <stdlib.h>
#include <math.h>
#include <zmq.h>
#include <pthread.h>

#define SOCKETS_NUM 5
#define RUNS 100000

void *context;
int numbers[SOCKETS_NUM];
struct {
    struct timespec send_time;
    struct timespec receive_time;
} times[SOCKETS_NUM * RUNS], *ptimes;

static void * worker_thread(void * dummy) {
    int * number = dummy;
    char endpoint[] = "inproc://endpointX";
    endpoint[17] = (char)('0' + *number);
    void * socket = zmq_socket(context, ZMQ_PUSH);
    zmq_connect(socket, endpoint);
    struct timespec sleeptime, remtime;
    int rnd = rand() / 3000;
    sleeptime.tv_sec = 0;
    sleeptime.tv_nsec = rnd;
    nanosleep(&sleeptime, &remtime);
    clock_gettime(CLOCK_REALTIME, &(ptimes[*number].send_time));
    zmq_send(socket, "Hello", 5, 0);
    zmq_close(socket);
    return NULL;
}

static void run_test(zmq_pollitem_t items[]) {
    pthread_t threads[SOCKETS_NUM];
    for (int i = 0; i < SOCKETS_NUM; i++) {
        pthread_create(&threads[i], NULL, worker_thread, &numbers[i]);
    }

    char buffer[10];
    int to_receive = SOCKETS_NUM;
    for (int i = 0; i < SOCKETS_NUM; i++) {
        int rc = zmq_poll(items, SOCKETS_NUM, -1);
        for (int j = 0; j < SOCKETS_NUM; j++) {
            if (items[j].revents & ZMQ_POLLIN) {
                clock_gettime(CLOCK_REALTIME, &(ptimes[j].receive_time));
                zmq_recv(items[j].socket, buffer, 10, 0);
            }
        }
        to_receive -= rc;
        if (to_receive == 0) break;
    }

    for (int i = 0; i < SOCKETS_NUM; i++) {
        pthread_join(threads[i], NULL);
    }
}

int main(void)
{
    context = zmq_ctx_new();
    zmq_ctx_set(context, ZMQ_THREAD_SCHED_POLICY, SCHED_FIFO);
    zmq_ctx_set(context, ZMQ_THREAD_PRIORITY, 99);
    void * responders[SOCKETS_NUM];
    char endpoint[] = "inproc://endpointX";
    for (int i = 0; i < SOCKETS_NUM; i++) {
        responders[i] = zmq_socket(context, ZMQ_PULL);
        endpoint[17] = (char)('0' + i);
        zmq_bind(responders[i], endpoint);
        numbers[i] = i;
    }

    time_t tt;
    time_t t = time(&tt);
    srand((unsigned int)t);

    zmq_pollitem_t poll_items[SOCKETS_NUM];
    for (int i = 0; i < SOCKETS_NUM; i++) {
        poll_items[i].socket = responders[i];
        poll_items[i].events = ZMQ_POLLIN;
    }

    ptimes = times;
    for (int i = 0; i < RUNS; i++) {
        run_test(poll_items);
        ptimes += SOCKETS_NUM;
    }

    long int lags[SOCKETS_NUM * RUNS];
    long int total_lag = 0;
    long int max_lag = 0;
    long int big_lag = 0;
    for (int i = 0; i < SOCKETS_NUM * RUNS; i++) {
        lags[i] = (times[i].receive_time.tv_nsec - times[i].send_time.tv_nsec + (times[i].receive_time.tv_sec - times[i].send_time.tv_sec) * 1000000000) / 1000;
        if (lags[i] > max_lag) max_lag = lags[i];
        total_lag += lags[i];
        if (lags[i] > 200) big_lag++;
    }
    long int avg_lag = total_lag / SOCKETS_NUM / RUNS;
    double SD = 0.0;
    for (int i = 0; i < SOCKETS_NUM * RUNS; ++i) {
        SD += pow((double)(lags[i] - avg_lag), 2);
    }
    double std_lag = sqrt(SD / SOCKETS_NUM / RUNS);
    printf("avg lag = %l5d    [us]\n", avg_lag);
    printf("max lag = %l5d    [us]\n", max_lag);
    printf("std dev = %8.2f [us]\n", std_lag);
    printf("big lag = %l5d    x above 200 [us]\n", big_lag);

    for (int i = 0; i < SOCKETS_NUM; i++) {
        zmq_close(responders[i]);
    }
    zmq_ctx_destroy(context);
    return 0;
}

【问题讨论】:

    标签: performance real-time zeromq latency low-latency


    【解决方案1】:

    Q:“...我希望它为零。”

    说起来很酷,但很难做到。

    当您运行超快的内存映射inproc://传输类时,主要关注点将是Context()-处理的性能调整。在这里,您花费了如此多的设置开销和直接终止开销操作来发送 1E5-times 只是一个 5 [B],所以我想永远不会有与队列管理相关的问题,因为根本不会有任何“堆栈增长”。

    1)(假设我们让代码保持原样)这将是性能调整的一个自然步骤,至少设置 socket-CPU_core ZMQ_AFFINITY 的 ZeroMQ 映射(不跳跃或徘徊从核心到核心)。有趣的是,如果 PUSH 端有那么多 ~ 5E5 套接字设置/终止,每个都不会发送超过一次的5 [B] 在内存映射行上,可以通过使用 SOCKETS_NUM I/O 线程配置 context-实例获得一些帮助(对于那些大的开销和维护),使用 @ 987654330@ 设置(争取“实时”-ness,使用SCHED_FIFO,只有一个 I/O 线程并没有多大帮助,不是吗?)

    2) 下一级实验是重新平衡 ZMQ_THREAD_AFFINITY_CPU_ADD 映射(全局 context 的 I/O 线程到 CPU 核心)和每个插槽设置的ZMQ_AFFINITY 映射到context 的I/O 线程。拥有足够数量的 CPU 内核,让多个 I/O 线程组为一个 socket-instance 服务于同一个 CPU 内核上保持“在一起”,可能会带来一些性能/超低延迟优势,然而,在这里,我们进入了一个领域,在没有任何体内测试和验证的情况下,实际硬件和真实系统的后台工作负载以及用于这种“实时”野心的实验的“备用”资源开始变得难以预测。

    3) 调整每个套接字 zmq_setsockopt() 参数可能会有所帮助,但除非纳米级套接字寿命(而不是昂贵的一次性使用的“一次性消耗品”),否则不要指望从这里取得任何突破。

    4) 尝试以纳秒分辨率进行测量,如果用于某事物的“持续时间”则越多,CLOCK_MONOTONIC_RAW 应该使用它,从而避免ntp-注入调整,天文学-校正闰秒注入等。

    5)zmq_poll()-strategy:我不会走这条路的。使用timeout == -1 会阻止整个马戏团。在任何分布式计算系统中,我都强烈反对这一点,尤其是在一个具有“实时”野心的系统中。将PULL-side 旋转到最大性能可以通过在任一侧具有 1:1 PUSH/PULL 线程,或者如果试图挑战修饰,拥有 5-PUSH-er 线程,就像你拥有的那样,并在一个单一的零拷贝上收集所有入口消息PULL-er(更容易轮询,可以使用基于有效负载的索引助手,发送端时间戳将接收端时间戳放入其中),无论如何,阻塞式轮询器几乎是挑战任何低延迟软实时玩具的反模式。

    无论如何,不​​要犹豫重构代码并使用分析工具更好地查看您“获取” big_lag-s 的位置(我的猜测在上面)

    #include <stdlib.h>
    #include <math.h>
    #include <zmq.h>
    #include <pthread.h>
    
    #define SOCKETS_NUM      5
    #define        RUNS 100000
    
    void *context;
    int   numbers[SOCKETS_NUM];
    struct {
        struct timespec send_time;
        struct timespec recv_time;
    } times[SOCKETS_NUM * RUNS],
     *ptimes;
    
    static void *worker_thread( void *dummy ) { //-------------------------- an ovehead expensive one-shot PUSH-based "Hello"-sender & .close()
        
        int   *number       = dummy;
        char   endpoint[]   = "inproc://endpointX";
               endpoint[17] = (char)( '0' + *number );
        int    rnd          = rand() / 3000;
        void  *socket       = zmq_socket( context, ZMQ_PUSH );
                
        struct timespec   remtime,
                        sleeptime;
                        sleeptime.tv_sec  = 0;
                        sleeptime.tv_nsec = rnd;
                        
        zmq_connect( socket, endpoint );
        
        nanosleep( &sleeptime, &remtime ); // anything betweed < 0 : RAND_MAX/3000 > [ns] ... easily >> 32, as #define RAND_MAX    2147483647 ~ 715 827 [ns]
        
        clock_gettime( CLOCK_REALTIME, &( ptimes[*number].send_time) ); //............................................................................ CLK_set_NEAR_SEND
                                                                        // any CLOCK re-adjustments may and will skew any non-MONOTONIC_CLOCK
        
        zmq_send(  socket, "Hello", 5, 0 );
        zmq_close( socket );
        
        return NULL;
    }
    
    static void run_test( zmq_pollitem_t items[] ) { //--------------------- zmq_poll()-blocked zmq_recv()-orchestrator ( called ~ 1E5 x !!! resources' nano-use & setup + termination overheads matter )
        
        char      buffer[10];
        int       to_receive = SOCKETS_NUM;
        pthread_t threads[SOCKETS_NUM];
        
        for ( int i = 0; i < SOCKETS_NUM; i++ ) { //------------------------ thread-maker ( a per-socket PUSH-er[]-s )
            pthread_create( &threads[i], NULL, worker_thread, &numbers[i] );
        }
        
        for ( int i = 0; i < SOCKETS_NUM; i++ ) { //------------------------ [SERIAL]-------- [i]-stepping
            
            int rc = zmq_poll( items, SOCKETS_NUM, -1 ); //----------------- INFINITE ??? --- blocks /\/\/\/\/\/\/\/\/\/\/\ --- several may flag ZMQ_POLLIN
            
            for ( int j = 0; j < SOCKETS_NUM; j++ ) { //-------------------- ALL-CHECKED in a loop for an items[j].revents
                
                if ( items[j].revents & ZMQ_POLLIN ) { //------------------- FIND IF IT WAS THIS ONE
                    
                    clock_gettime( CLOCK_REALTIME, &( ptimes[j].recv_time ) );//...................................................................... CLK_set_NEAR_poll()_POSACK'd R2recv
                    
                    zmq_recv( items[j].socket, buffer, 10, 0 ); //---------- READ-IN from any POSACK'd by zmq_poll()-er flag(s)
                }
            }
            to_receive -= rc; // ---------------------------------------------------------------------------------------------- SUB rc
            if (to_receive == 0) break;
        }
    
        for ( int i = 0; i < SOCKETS_NUM; i++ ) { //------------------------ thread-killer
            
            pthread_join( threads[i], NULL );
        }
    }
    
    int main( void ) {
        
                     context = zmq_ctx_new();
        zmq_ctx_set( context, ZMQ_THREAD_SCHED_POLICY, SCHED_FIFO );
        zmq_ctx_set( context, ZMQ_THREAD_PRIORITY, 99 );
        
        void *responders[SOCKETS_NUM];
        char  endpoint[] = "inproc://endpointX";
        
        for ( int i = 0; i < SOCKETS_NUM; i++ ) {
            
            responders[i] = zmq_socket( context, ZMQ_PULL ); // ------------ PULL instances into []
            endpoint[17] = (char)( '0' + i );
            zmq_bind( responders[i], endpoint ); //------------------------- .bind()
            numbers[i] = i;
        }
    
        time_t tt;
        time_t t = time(&tt);
        srand( (unsigned int)t );
    
        zmq_pollitem_t poll_items[SOCKETS_NUM];
        
        for ( int i = 0; i < SOCKETS_NUM; i++ ) { //------------------------ zmq_politem_t array[] ---pre-fill---
            poll_items[i].socket = responders[i];
            poll_items[i].events = ZMQ_POLLIN;
        }
    
        ptimes = times;
        
        for ( int i = 0; i < RUNS; i++ ) { //------------------------------- 1E5 RUNs
            run_test( poll_items ); // -------------------------------------     RUN TEST
            ptimes += SOCKETS_NUM;
        }
    
        long int lags[SOCKETS_NUM * RUNS];
        long int total_lag = 0;
        long int   max_lag = 0;
        long int   big_lag = 0;
        
        for ( int i = 0; i < SOCKETS_NUM * RUNS; i++ ) {
            lags[i] = (   times[i].recv_time.tv_nsec
                      -   times[i].send_time.tv_nsec
                      + ( times[i].recv_time.tv_sec
                        - times[i].send_time.tv_sec
                          ) * 1000000000
                        ) / 1000; // --------------------------------------- [us]
            if ( lags[i] > max_lag ) max_lag = lags[i];
            total_lag += lags[i];
            if ( lags[i] > 200 )     big_lag++;
        }
        
        long int avg_lag = total_lag / SOCKETS_NUM / RUNS;
        double        SD = 0.0;
        
        for ( int i = 0; i < SOCKETS_NUM * RUNS; ++i ) {
            SD += pow( (double)( lags[i] - avg_lag ), 2 );
        }
        
        double std_lag = sqrt( SD / SOCKETS_NUM / RUNS );
        
        printf("avg lag = %l5d    [us]\n", avg_lag);
        printf("max lag = %l5d    [us]\n", max_lag);
        printf("std dev = %8.2f [us]\n", std_lag);
        printf("big lag = %l5d    x above 200 [us]\n", big_lag);
    
        for ( int i = 0; i < SOCKETS_NUM; i++ ) {
            zmq_close( responders[i] );
        }
        zmq_ctx_destroy( context );
        
        return 0;
    }
    

    使用nanosleep 进行随机(不是基数,安全地在任何控制循环活动之外)睡眠是相当冒险的奢侈,因为在早期的内核中会导致问题:

    为了支持需要更精确暂停的应用程序(例如,为了控制一些对时间要求严格的硬件),nanosleep() 将通过忙等待处理长达 2 毫秒的暂停 从实时策略(如SCHED_FIFO 或SCHED_RR)下调度的线程调用时,精度为微秒。这个特殊的扩展在内核 2.5.39 中被删除,因此在当前的 2.4 内核中仍然存在,但在 2.6 内核中没有。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-04-11
      • 2011-12-10
      • 1970-01-01
      • 1970-01-01
      • 2011-05-17
      • 2023-03-12
      • 2020-06-24
      • 2010-09-09
      相关资源
      最近更新 更多