Linux高性能编程_无锁队列

Linux高性能编程_无锁队列

目录

大家好,这里是物联网心球。

今天我们讨论的主题是无锁队列,有过代码性能优化经验的小伙伴们,应该都听过无锁队列,无锁队列在很多高性能开源项目都有运用,比如:XDP,DPDK等。

1.无锁队列简介

无锁队列是一种数据结构,它在并发编程中不依赖于传统的锁定机制(如互斥锁或信号量)来同步访问,而是利用原子操作、内存屏障或者其他非阻塞的同步原语。无锁算法通常用于提高并发性能和降低死锁风险。

优点:

  1. 高并发性:无锁队列避免了锁竞争,使得多个线程能够并行地添加和移除元素,提高了系统的吞吐量。

  2. 减少阻塞:由于没有显式锁,读写操作不会被其他线程阻塞,提高了响应时间和系统利用率。

  3. 减少上下文切换:无锁操作通常导致更少的线程上下文切换,减少了CPU开销。

缺点:

  1. 复杂性增加:实现无锁队列往往比使用锁更为复杂,需要对底层硬件和操作系统有更好的理解。

  2. 代码可读性和调试困难:由于缺乏锁的直观性,代码可能变得难以理解和维护,出错排查难度大。

  3. 适应特定环境:并非所有情况都适合无锁设计,比如处理器不支持原子操作或者内存模型不保证顺序一致性时,效果可能会大打折扣。

2.无锁队列实现原理

无锁队列的实现相对来说比较复杂,要想实现无锁队列,我们得了解一些背景知识:原子操作,内存屏障,内存顺序。

2.1 原子操作

原子操作是无锁化编程的基础,原子操作前面已经介绍过,请参考文章:

Linux高性能编程_原子操作

2.2 内存屏障

内存屏障其实就是防止指令重排的一种技术,编译器和CPU会背着程序员偷偷修改指令执行的顺序。程序没有按照程序设计的逻辑执行,导致程序出错。

如下图,程序设计的逻辑是顺序执行指令1,指令2,指令3,指令4。

Linux高性能编程_无锁队列 图1

编译器或CPU重排了指令1和指令3,先执行指令3再执行指令1,程序运行结果出错。

Linux高性能编程_无锁队列 图2

为了防止指令被重排,我们在指令2和指令3之间增加一个内存屏障,内存屏障之前的代码和内存屏障之后的代码将禁止重排,禁止重排后可以保证指令1和指令2先于指令3和指令4执行。

当然指令1和指令2之间以及指令3和指令4之间可能存在指令重排,我们可以继续增加内存屏障解决。

Linux高性能编程_无锁队列 图3

内存屏障是一个很复杂的概念,本文只是简单的介绍了内存屏障的工作原理,方便理解后续的概念。

2.3 内存顺序

内存顺序的定义如下:

typedef enum memory_order {
  memory_order_relaxed = __ATOMIC_RELAXED,
  memory_order_consume = __ATOMIC_CONSUME,
  memory_order_acquire = __ATOMIC_ACQUIRE,
  memory_order_release = __ATOMIC_RELEASE,
  memory_order_acq_rel = __ATOMIC_ACQ_REL,
  memory_order_seq_cst = __ATOMIC_SEQ_CST
} memory_order;

内存顺序通常是和原子操作共同使用,这样原子操作会增加额外的功能(如:内存屏障功能),无锁队列将会用到三种内存顺序:

1)memory_order_relaxed

  • 只保证原子操作,不保证指令顺序。

2)memory_order_acquire

  • 用于获取(load)原子变量值相关的操作,比如:atomic_load_explicit操作。

  • 对于使用memory_order_acquire内存顺序的指令,该指令后面的所有读写操作不能重排在该指令之前。

  • 当前线程执行的memory_order_acquire指令能够保证读到其他线程 memory_order_release指令之前的所有内存写入操作。

3)memory_order_release

  • 用于设置(store)原子变量值相关的操作,比如:atomic_store_explicit操作。

  • 对于使用memory_order_release内存顺序的指令,该指令之前的所有读写操作不能重排在该指令之后。

  • 当前线程memory_order_release指令之前的所有内存写操作对于其他线程 memory_order_acquire指令之后可见。

注意:能否正确理解这三种内存顺序是学好无锁队列的关键!!!

三种内存顺序除了红色字体部分难理解,其他的部分相对来说比较简单,所以我们重点讲解红色字体部分。

memory_order_acquire和memory_order_release内存顺序通常来说是成对使用,实现多线程互斥访问的功能。

Linux高性能编程_无锁队列 图4

3.环形无锁队列

3.1 环形无锁队列整体介绍

一个完整的环形无锁队列需要具备很多要素,其中最重要的几个要素如下:

  • 固定大小的数组,用于存放队列元素。

  • 数组长度,用于计算已插入元素(entries)和空闲元素(free_entries)。

  • 生产者索引,生产者索引分为头部(prod_head)和尾部(prod_tail),且都是原子整型变量。

  • 消费者索引,消费者索引分为头部(cons_head)和尾部(cons_tail),且都是原子整型变量。

Linux高性能编程_无锁队列 图5

为了便于讨论无锁队列,我们把环形无锁队列看成是一个数组。通过数组我们很容易理解无锁队列入队和出队操作。

Linux高性能编程_无锁队列 图6

3.2 环形无锁队列入队

多个线程同时入队时,为了保证入队的正确性,需要将入队操作分为三个步骤:移动生产者头部,插入队列元素,移动生产者尾部。我们需要重点关注每一个步骤如何保证多线程安全。1)移动生产者头部

Linux高性能编程_无锁队列 图7

无锁队列将生产者索引分为头部和尾部,目的是为了保证多线程入队安全。如上图,3个线程同时入队,入队先后顺序为:线程1>线程2>线程3。

  • 线程1:记录当前生产者头部位置(old_head),根据入队元素个数计算出新的生产者头部(new_head),使用CAS原子操作将prod_head设置成new_head的值。

  • 线程2:记录当前生产者头部位置(old_head),此时的old_head为线程1更新后的prod_head值,根据入队元素个数计算出新的生产者头部(new_head),将prod_head设置成new_head的值。

  • 线程3:同线程2。示例代码:

int mini_ring_move_prod_head(struct mini_ring *ring, uint32_t n, unsigned int *cur_head, unsigned int *next_head) {
    uint32_t cons_tail;
    int free_entries;
    bool suc;
    *cur_head = atomic_load_explicit(&ring->prod.head, memory_order_relaxed);
    do {
        cons_tail = atomic_load_explicit(&ring->cons.tail, memory_order_acquire);
        free_entries = ring->size - (*cur_head - cons_tail);
        if (n > free_entries) {
            return -1;
        }
        *next_head = *cur_head + n;
        suc = atomic_compare_exchange_strong_explicit(&ring->prod.head,
                cur_head,
                *next_head,
                memory_order_relaxed,
                memory_order_relaxed);
    } while (suc == false);
    return 0;
}

2)插入队列元素

Linux高性能编程_无锁队列 图8

old_head和new_head为队列元素插入的起始位置和结束位置。移动头部完成后,线程1,线程2,线程3可以同时插入队列元素。为什么3个线程同时插入队列元素不会相互干扰,大家可以思考一下。 原因在于3个线程的old_head和new_head没有重叠,所以没有藕合关系,就能同时插入队列元素。

3)移动生产者尾部 Linux高性能编程_无锁队列 图9

移动尾部不是简单把prod_tail设置成new_head的值,如果3个线程同时设置prod_tail的值,那么最终prod_tail的值将变得不确定,如果prod_tail的值不确定会导致无法正确记录entries(已插入元素)的起始位置。移动尾部需要具备一定的条件才能移动,条件就是线程prod_tail和old_head值相等才能移动,该条件可以限制未满足条件的线程修改prod_tail的值,并且一个线程移动完尾部,又会有其他线程满足该条件,继续移动尾部,确保多项成安全,此时3个线程只有线程1满足情况。

  • 线程1:设置prod_tail为new_head,移动尾部完成。
  • 线程2:线程1更新完prod_tail后,线程2检测到prod_tail等于old_head值,设置prod_tail为new_head,移动尾部完成。
  • 线程3:同线程2。

示例代码:

int mini_ring_update_prod_tail(struct mini_ring *ring, uint32_t *cur_head, uint32_t *next_head) {
    while(atomic_load_explicit(&ring->prod.tail, memory_order_relaxed) != *cur_head) {}

    atomic_store_explicit(&ring->prod.tail, *next_head, memory_order_release);
}

3.3 环形无锁队列出队

环形无锁队列出队和入队非常相似,同样分为三个步骤:移动消费者头部,移出队列元素,移动消费者尾部。这里我只贴出过程图,小伙伴们自行分析实现原理。

1)移动消费者头部 Linux高性能编程_无锁队列 图10 2)移出队列元素

Linux高性能编程_无锁队列 图11 3)移动消费者尾部

Linux高性能编程_无锁队列 图12

4.无锁队列性能测试

4.1 测试方法

测试代码包括2种队列:无锁队列和互斥锁队列。 将测试代码编译成可执行程序,通过time ./程序名 执行程序,对比互斥锁队列和无锁队列执行时间,不清楚time命令的使用方法,请参考文章:Linux高性能编程_原子操作 测试代码如下:注:由于篇幅有限,测试代码仅为部分代码,需要完整版代码请联系博主。

 #include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <stdbool.h>
#include <unistd.h>
#include <pthread.h>

#define PROD_THREAD_NUM (1) //生产者线程数量
#define CONS_THREAD_NUM (1) //消费者线程数量
#define TOTAL_PACKET_NUM (20000000) //测试入队数据包总数量

#define LOCK_TYPE (1)
#define FREE_LOCK (0)
#define MUTEX_LOCK (1)

struct mini_ring;
struct mini_ring *g_ring;
_Atomic(uint32_t) g_seq; //全局数据包序列号,每产生一个新的数据包加1
_Atomic(uint32_t) g_count; //全局出队数据包序列号,每出队一个数据包加1

int g_done;

pthread_mutex_t g_mutex;
pthread_cond_t g_cond;

struct mini_packet {
    uint32_t seq; //数据包序列号
};

struct mini_headtail {
    _Atomic(uint32_t) head;
    _Atomic(uint32_t) tail;
};

struct mini_ring {
    struct mini_headtail prod;
    struct mini_headtail cons;
    uint32_t head;
    uint32_t tail;
    uint32_t size;
    uint32_t mask;
};

int mini_ring_enqueue(struct mini_ring *ring, void *obj_table, uint32_t n) {
#if (LOCK_TYPE == FREE_LOCK)
    uint32_t cur_head;
    uint32_t next_head;

    int ret = mini_ring_move_prod_head(ring, n, &cur_head, &next_head);
    if (ret) {
        //printf("mini ring move prod head error\n");
        return -1;
    }
    mini_ring_enqueue_elems(ring, &cur_head, obj_table, n);
    mini_ring_update_prod_tail(ring, &cur_head, &next_head);
#elif (LOCK_TYPE == MUTEX_LOCK)
    uint32_t cur_head;
    uint32_t next_head;
    uint32_t free_entries;

    pthread_mutex_lock(&g_mutex);
    free_entries = ring->size - (ring->head - ring->tail);
    if (n > free_entries) {
        pthread_cond_signal(&g_cond);
        pthread_mutex_unlock(&g_mutex);
        return -1;
    }

    cur_head = ring->head;
    next_head = ring->head + n;
    mini_ring_enqueue_elems(ring, &cur_head, obj_table, n);
    ring->head = next_head;
    pthread_cond_signal(&g_cond);
    pthread_mutex_unlock(&g_mutex);
#else

#endif
    return 0;
}

int mini_ring_dequeue(struct mini_ring *ring, void *obj_table, uint32_t n) {
#if (LOCK_TYPE == FREE_LOCK)
    uint32_t cur_head;
    uint32_t next_head;

    int ret = mini_ring_move_cons_head(ring, n, &cur_head, &next_head);
    if (ret) {
        //printf("mini ring move cons head ret:%d error\n", ret);
        return -1;
    }

    mini_ring_dequeue_elems(ring, &cur_head, obj_table, n);
    mini_ring_update_cons_tail(ring, &cur_head, &next_head);
#elif (LOCK_TYPE == MUTEX_LOCK)
    uint32_t cur_tail;
    uint32_t next_tail;
    uint32_t entries;
    pthread_mutex_lock(&g_mutex);
    entries = ring->head - ring->tail;
#if 0
    if (n > entries) {
        pthread_mutex_unlock(&g_mutex);
        return -1;
    }
#else
    while((n > entries) && (g_done == 0)) {
        pthread_cond_wait(&g_cond, &g_mutex);
        entries = ring->head - ring->tail;
    }

    if (g_done != 0) {
        pthread_mutex_unlock(&g_mutex);
        return -1;
    }
#endif

    cur_tail = ring->tail;
    next_tail = ring->tail + n;
    mini_ring_dequeue_elems(ring, &cur_tail, obj_table, n);
    ring->tail = next_tail;
    pthread_mutex_unlock(&g_mutex);
#else
#endif
    return 0;
}

void *prod_thread(void *arg) {
    int thread_num = (int)arg;
    uint64_t **obj_table = NULL;
    while(1) {
        //for (int i = 0; i < 10000; i++) ;
        struct mini_packet *pkt = mini_packet_alloc();
        if (!pkt) {
            printf("入队:%d个数据包,生产者线程:%d 退出!\n", TOTAL_PACKET_NUM, thread_num);
            break;
        }
        //printf("prod thread num:%d, new pkt seq:%u\n", thread_num, pkt->seq);
        obj_table = (uint64_t **)&pkt;
try_again:
        int ret = mini_ring_enqueue(g_ring, (void *)obj_table, 1);
        if (ret) {
            //printf("mini ring enqueue ret:%d error\n", ret);
            goto try_again;
        }
    }

    return NULL;
}

void *cons_thread(void *arg) {
    int thread_num = (int)arg;
    int ret = 0;
    uint64_t **obj_table = NULL;
    struct mini_packet *pkt = NULL;
    int failed_times = 0;
    uint32_t count = 0;
    while(1) {
        count = atomic_load(&g_count);
        if (count > TOTAL_PACKET_NUM) {
            printf("出队:%d个数据包,消费者线程:%d 退出!\n", count - 1, thread_num);
            break;
        }
        obj_table = (uint64_t **)&pkt;
        ret = mini_ring_dequeue(g_ring, (void *)obj_table, 1);
        if (ret < 0) {
            if (failed_times < 1000) {
                failed_times++;
            } else {
                //printf("dequeue failed times:%d\n", failed_times);
                usleep(10);
                failed_times = 0;
            }
            continue;
        }

#if 0
        if ((pkt->seq % 1000) == 0){
            printf("cons thread num:%d, dequeue pkt seq:%u\n", thread_num, pkt->seq);
        }
#endif
        mini_packet_free(pkt);
        atomic_fetch_add(&g_count, 1);
    }

    g_done++;

    return NULL;
}

int main(int argc, char *argv[]) {
    atomic_init(&g_seq, 0);
    atomic_init(&g_count, 0);
    pthread_mutex_init(&g_mutex, NULL);
    pthread_cond_init(&g_cond, NULL);
    g_done = 0;

    int size = 4096;
    g_ring = mini_ring_create(size);

    pthread_t prod_th[PROD_THREAD_NUM];
    pthread_t cons_th[CONS_THREAD_NUM];
    for (int i = 0; i < CONS_THREAD_NUM; i++) {
        pthread_create(&cons_th[i], NULL, cons_thread, (void *)i);
    }
    for (int i = 0; i < PROD_THREAD_NUM; i++) {
        pthread_create(&prod_th[i], NULL, prod_thread, (void *)i);
    }

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

    while(g_done < CONS_THREAD_NUM) {
        pthread_cond_signal(&g_cond);
        usleep(10);
        //printf("g_done:%d\n", g_done);
    }

    //printf("--g_done:%d\n", g_done);
    for (int i = 0; i < CONS_THREAD_NUM; i++) {
        pthread_join(cons_th[i], NULL);
    }

    pthread_mutex_destroy(&g_mutex);
    pthread_cond_destroy(&g_cond);

    return 0;
}

4.2 测试结果

测试硬件环境:树莓派4B,内存4GB。测试要求:2千万次入队和出队操作。1)1个生产者和1个消费者 无锁队列»>

Linux高性能编程_无锁队列 图13

结果分析:实际运行时间为8.3秒,系统运行时间为0.0秒,无系统调用。 互斥锁队列 »>

Linux高性能编程_无锁队列 图14

结果分析:互斥锁队列运行很不稳定,实际运行时间超过10分钟。

2)1个生产者和2个消费者 无锁队列»>

Linux高性能编程_无锁队列 图15

结果分析:实际运行时间为8.3秒,系统运行时间为0.1秒,无系统调用。 互斥锁队列»>

Linux高性能编程_无锁队列 图16

结果分析:实际运行时间为15.6秒,系统运行时间为19.9秒,存在大量系统调用。 3)2个生产者和1个消费者无锁队列»>

Linux高性能编程_无锁队列 图17

结果分析:实际运行时间为7.9秒,系统运行时间为0.0秒,无系统调用。 互斥锁队列»>

Linux高性能编程_无锁队列 图18

结果分析:实际运行时间为26.5秒,系统运行时间为29.7秒,存在大量系统调用。 4)2个生产者和2个消费者无锁队列»>

Linux高性能编程_无锁队列 图19

结果分析:实际运行时间为6.18秒,系统运行时间为0.0秒,无系统调用。 互斥锁队列»>

Linux高性能编程_无锁队列 图20

结果分析:实际运行时间为11.8秒,系统运行时间为22.1秒,存在大量系统调用。

← 返回文章列表