Murphy

Murphy

@Murphy

  • 无锁循环队列

    内存屏障/原子操作的一种常见应用

    单生产单消费

    KFIFO

    Linux 内核队列改写的 c++ 版本

    不需要使用任何原子变量,巧妙应用内存屏障,在读/写结束后才更新头尾指针,保证了当前线程正在读/写的部分对于对家线程不可见

    template <class ValueType>
    class kfifo {
    private:
    	ValueType* buffer;	/* the buffer holding the data */
    	unsigned int size;	/* the size of the allocated buffer */
    	unsigned int in;	/* data is added at offset (in % size) */
    	unsigned int out;	/* data is extracted from off. (out % size) */
    
    public:
        kfifo(unsigned int sz) : size((1<<sz)), in(0), out(0) {
            buffer = new ValueType[size];
        }
    
        kfifo(const kfifo& ) = delete;
        kfifo(const kfifo&& ) = delete;
        void operator = (const kfifo&& ) = delete;
    
        ~kfifo() {
            delete[] buffer;
        }
        
        bool push(const char *item) {
    
    /*******************两种自旋方式都是正确的************************/
            
    // MODE1:        
            std::atomic_thread_fence(std::memory_order_acquire);
            if (size - in + out == 0) {
                return false;
            } // 函数外 while 自旋
    // MODE2:        
            // while (size - in + out == 0){
            //     std::atomic_thread_fence(std::memory_order_acquire);
    
            // } // 函数内 while 自旋
    
            std::atomic_thread_fence(std::memory_order_acq_rel);
    
            std::memcpy(buffer+(in & (size - 1)), item, sizeof(ValueType));
    
            std::atomic_thread_fence(std::memory_order_release);
    
            ++ in;
            return true;
        }
    
        bool pop(char *item) {
    // MODE1:  
            std::atomic_thread_fence(std::memory_order_acquire);
            if (in == out) {
                return false;
            } // 函数外 while 自旋
    // MODE2:  
            // while (in == out){
            //     std::atomic_thread_fence(std::memory_order_acquire);
        
            // } // 函数内 while 自旋
    
            std::atomic_thread_fence(std::memory_order_acquire);
    
            std::memcpy(item, buffer+(out & (size - 1)), sizeof(ValueType));
    
            std::atomic_thread_fence(std::memory_order_acq_rel);
    
            ++ out;
            return true;
        }
    };
    

    如果一次性拷贝的长度较大或是长度不一可以写成以下形式

        unsigned int push(const char *buf, unsigned int len)
        {
            unsigned int l;
    
            len = min(len, size - in + out);
    
            std::atomic_thread_fence(std::memory_order_acq_rel);
    
            l = min(len, size - (in & (size - 1)));
            memcpy(buffer + (in & (size - 1)), buf, l);
    
            memcpy(buffer, buf + l, len - l);
    
            std::atomic_thread_fence(std::memory_order_release);
    
            in += len;
    
            return len;
        }
    
        unsigned int pop(char *buf, unsigned int len)
        {
            unsigned int l;
    
            len = min(len, in - out);
    
            std::atomic_thread_fence(std::memory_order_acquire);
    
            l = min(len, size - (out & (size - 1)));
            memcpy(buf, buffer + (out & (size - 1)), l);
    
            memcpy(buf + l, buffer, len - l);
    
            std::atomic_thread_fence(std::memory_order_acq_rel);
    
            out += len;
    
            return len;
        }
    

    SPSC

    使用CAS (compare and swap) 原子操作实现

    template <typename T>
    class spsc {
        T *data;
        std::atomic<size_t> head{0}, tail{0};
        size_t Cap;
    public:
        spsc(size_t siz): Cap(1<<siz) {
            data = new T[Cap];
        }
        spsc(const spsc&) = delete;
        spsc &operator=(const spsc&) = delete;
        spsc &operator=(const spsc&) volatile = delete;
    
        bool push(const T &val) { 
            
    /*************************注释部分为错误版本***********************************/
            // size_t t;
            // do {
            //     t = tail.load(std::memory_order_acquire);
            // } while ((t + 1) & (Cap-1) == head.load(std::memory_order_acquire));
    
            size_t t = tail.load(std::memory_order_relaxed);
            if ((t + 1) % Cap == head.load(std::memory_order_acquire)) return false;
            std::memcpy(data + t, &val, sizeof(T));
            tail.store((t + 1) & (Cap-1), std::memory_order_release);
            return true;
        }
    
        bool pop(T &val) { 
            // size_t h;
            // do {
            //     h = head.load(std::memory_order_acquire);
            // } while (h == tail.load(std::memory_order_acquire));
            size_t h = head.load(std::memory_order_relaxed);
            if (h == tail.load(std::memory_order_acquire)) return false;
            std::memcpy(&val, data + h, sizeof(T));
            head.store((h + 1) & (Cap-1), std::memory_order_release);
            return true;
        }
    };
    

    多生产多消费

    Disruptor

    实现原理

    4 个位置变量

    • lastRead 最后一个已读内容位置
    • lastWrote 最后一个已写内容位置
    • lastDispatch 最后一个派发给消费者的槽位序号
    • writableSeq 当前可写的槽位序号

    记录 readableSeq=lastDispatch+1 表示当前可读槽位序号

    他们从小到大依次为:lastRead,readableSeq,lastWrote,writableSeq

    实现逻辑为

    • 对于生产者而言,先申请 writableSeq 槽位,更新 writableSeq;写完之后,等待在它之前的槽位都被写完才更新 lastWrote

    • 对于消费者而言,先申请 readableSeq 槽位,更新 readableSeq;读完之后,等待在它之前的槽位都被读完才更新 lastRead

    由以上两点可知

    • lastRead~readableSeq 为正在被消费的部分
    • readableSeq~lastWrote 为已经生产完的部分,对于消费者可见
    • lastWrote~writableSeq 为正在被生产的部分
    • writableSeq~lastRead 为已经消费完的部分,对于生产者可见

    因此只有被更新完毕的部分才对对家线程可见

    一些优化

    • 首先是将常用的变量用 alignas(64) 进行内存对齐,以防交替冲突

    • 加速取余,将队列大小 N 设置为 2 的次幂的编译时常量,将 x%N 优化为x & (N-1)

    • 避免取余,将 4 个位置变量设置为unsigned int,自动处理溢出

    template<class ValueType , size_t N = DefaultRingBufferSize>
    class Disruptor
    {
    public:
        Disruptor() : _lastRead(-1L) , _lastWrote(-1L), _lastDispatch(-1L), _writableSeq(0L) , _stopWorking(false){};
        ~Disruptor() {};
    
        Disruptor(const Disruptor&) = delete;
        Disruptor(const Disruptor&&) = delete;
        void operator=(const Disruptor&) = delete;
    
        static_assert(((N > 0) && ((N& (~N + 1)) == N)),
            "RingBuffer's size must be a positive power of 2");
    
        template<typename T>
        bool push(ValueType& val)
        {
            const uint64_t writableSeq = _writableSeq.fetch_add(1);
            while (writableSeq - _lastRead > N)
            {
                if (_stopWorking.load()) {
                    // throw std::runtime_error("writting when stopped disruptor");
                    return false;
                }
                std::this_thread::yield();
            }
            std::memcpy(_ringBuf + (writableSeq & (N - 1)), &val, sizeof(ValueType));
    
            while (writableSeq - 1 != _lastWrote);
            _lastWrote = writableSeq;
            return true;
        };
    
        bool pop(ValueType& val) {
            const uint64_t readableSeq = _lastDispatch.fetch_add(1) + 1;
            while (readableSeq > _lastWrote)
            {
                if (_stopWorking.load())
                {
                    // throw std::runtime_error("reading when stopped disruptor");
                    return false; 
                }
                std::this_thread::yield();
            }
            std::memcpy(&val, _ringBuf + (readableSeq & (N-1)), sizeof(ValueType));
            while (readableSeq - 1 != _lastRead);
            _lastRead = readableSeq;
            return true;
        }
    
        bool empty() {return _writableSeq - _lastRead == 1;}
    
        //通知 disruptor 停止工作,调用该函数后,若 buffer 已经全部处理完,那么获取可读下标时只会获取到 -1L
        void stop() {_stopWorking.store(true);}
    
    private:
        
        alignas(64) std::atomic_bool _stopWorking;
        
        alignas(64) std::atomic_uint64_t _lastRead;
        
        alignas(64) std::atomic_uint64_t _lastWrote;
        
        alignas(64) std::atomic_uint64_t _lastDispatch;
        
        alignas(64) std::atomic_uint64_t _writableSeq;
        
        ValueType _ringBuf[N];
    };
    

    RingBuffer

    使用CAS (compare and swap) 操作实现

    实现原理

    将非原子过程用do while循环包起来,最后的判断设计为一个 CAS 操作。如果 CAS 操作失败,意味着该过程被其他线程抢占,那么重新执行该过程,不断重试直到强占到执行权。

    注意,必须先读取再更新指针才能保证并发正确,但考虑到

    • push 时数据拷贝必须放在循环体外,否则强占到执行权的线程写入的数据有可能被其他未强占到执行权的线程修改
    • pop 时在循环体内拷贝数据不会导致错误,但没有强占到执行权的线程会反复的拷贝数据,对 CPU 资源消耗很大

    因此引入两个新指针

    • write,表示 push 操作写完的位置
    • read,表示 pop 操作读完的位置

    由此可知

    • read~head 部分正在被读
    • head~write 部分可读
    • write~tail 部分正在被写
    • tail~read 部分可写

    (:其实和 Disruptor 是一样的

    template <typename T>
    class RingBuffer {
    public:
        RingBuffer(size_t siz): Cap(1<<siz){ data = new T[Cap];}
        RingBuffer(const RingBuffer&) = delete;
        RingBuffer &operator=(const RingBuffer&) = delete;
        RingBuffer &operator=(const RingBuffer&) volatile = delete;
    
        bool push(const T &val) {
            size_t t, w;
            do {
                t = tail.load();
                if ((t + 1) % Cap == head.load())
                    return false;
            } while (!tail.compare_exchange_weak(t, (t + 1) % Cap));
            std::memcpy(data+t, &val, sizeof(T));
            do {
                w = t;
            } while (!write.compare_exchange_weak(w, (w + 1) % Cap));
            return true;
        }
    
        bool pop(T &val) {
            size_t h;
            do {
                h = head.load();
                if (h == write.load()) 
                    return false;
                std::memcpy(&val, data+h, sizeof(T));
            } while (!head.compare_exchange_strong(h, (h + 1) % Cap));
            return true;
        }
    private:
        T* data;
        std::atomic<size_t> head{0}, tail{0}, write{0};
        uint64_t Cap;
    };
    

    测试

    using ValueType = char[64],以 512bit 的字符数组为一单位元素进行测试

    性能测试

    线程均绑定隔离核

    节选结果如下:

    *:~/fastQ/bin$ ./test_sng --op 3e9
    Sng_Queue Performance Test for Queue size 1<<10, 1 producers, 1 consumers, 3000000000 operations per producer, 0 nanoseconds per production, 0 nanoseconds per consumption
    Mutex CircularQueue Elapsed time: 350.593 seconds
    KFIFO Elapsed time: 34.319 seconds
    SPSC Elapsed time: 32.0383 seconds
    
    *:~/fastQ/bin$ ./test_muti --pd 3 --cs 3 --op 1e9
    Muti_Queue Performance Test for Queue size 1<<10, 3 producers, 3 consumers, 1000000000 operations per producer, 0 nanoseconds per production, 0 nanoseconds per consumption
    Muti-ptr Disruptor Elapsed time: 321.647 seconds
    Atomic RingBuffer Elapsed time: 710.012 seconds
    Mutex CircularQueue Elapsed time: 650.429 seconds
    
    *:~/fastQ/bin$ ./test_muti --pd 3 --cs 1 --op 1e9
    Muti_Queue Performance Test for Queue size 1<<10, 3 producers, 1 consumers, 1000000000 operations per producer, 0 nanoseconds per production, 0 nanoseconds per consumption
    Muti-ptr Disruptor Elapsed time: 234.385 seconds
    Atomic RingBuffer Elapsed time: 430.647 seconds
    Mutex CircularQueue Elapsed time: 732.851 seconds
    

    kfifospsc 速度比加锁快一个数量级,spsc 更优一些

    Disruptor 速度约为加锁的 2~3 倍

    RingBuffer 在消费者数量相对较多时表现不佳,怀疑是 CAS 操作导致竞争过于激烈浪费了 CPU 资源

    同等数据量的下,单生产单消费的 kfifo 的速度比多生产多消费的 Disruptor 快 7~8 倍

    同时,多生产多消费情况存在竞争,当竞争过于激烈时,无锁的性能就相对下降了

    正确性测试

    -O3 的情况下也都能通过,除了 spsc 注释中标注的错误方法,暂时还没有解决(:

    在性能测试时 spsc 的表现比 kfifo 更优秀一些,但开 -O3 后,如果将自旋的 while 改到 pushpop 函数内会出现死循环。有个现象是在自旋中sleep for 1 nanosecond就可以正确运行,不懂 ovo

    局限

    • 存储智能指针时它指向的对象并不能被及时的析构,出队后对象仍然在 data 数组里,并没有立即销毁
    • 为保证性能数组大小 N 应当在编译时确定,无法动态扩容。
    Post #23 ❤️ 2 likes
  • 内存屏障与编译器屏障(C++)

    内存屏障与内存顺序

    内存屏障用来阻止处理器重排内存操作,是硬件或编译器层面上的工具。

    内存顺序指在多线程程序中多个线程如何观察彼此的读写操作,是 C++ 提供的高级抽象。它通过指定不同的内存顺序模型隐式地插入内存屏障。程序员不需要直接操作内存屏障,而是通过原子操作与内存顺序模型来实现对内存访问顺序的控制。

    二者都是为了控制内存操作,防止指令重排、缓存一致性问题带来的意料之外的错误

    内存顺序

    讨论的是单线程内指令执行顺序对多线程影响的问题。不论携带怎样的内存顺序,单独的操作只会影响单个线程,但几个线程可能会经由变量的同步建立间接的顺序关系,从而被各自的内存顺序影响。

    内存顺序标记

    c++ 标准库定义了 6 种内存顺序,按照从宽松到严格的顺序:

    1. memory_order_relaxed:宽松操作,仅保证原子操作自身的原子性
    2. memory_order_consume:读屏障,只影响当前 load 原子变量及依赖变量的写入,其他与 memory_order_acquire 一样。C++20 中弃用
    3. memory_order_acquire:读屏障,当前线程所有的读写操作不得重排至当前操作之前,使得其他线程中相同原子变量在 release 之前的写入对当前线程时可见的。即对于相关内存位置的 load 施加“acquire”占有,一定读到最新值。
    4. memory_order_release:写屏障,当前线程所有的读写操作不得重排至当前操作之后,使得当前操作的写入对于其他线程中相同原子变量在 acquire 之后是可见的。即对于相关内存位置的 store 施加“release”释放,一定写入最新值。
    5. memory_order_acq_rel:读写屏障,当前线程中的读写操作不能重排至当前操作之后也不能重排至当前操作之前
    6. memory_order_seq_cst:读写屏障,对于所有线程看见的读写操作顺序是一致的

    其中

    • 所有原子操作默认的内存顺序标记是 std::memory_order_seq_cst
    • 两种读屏障:memory_order_consume 只影响当前变量和依赖它的变量;memory_order_acquire 影响当前操作前的所有变量
    • 两种写屏障:memory_order_acq_rel 只对当前线程有效;memory_order_seq_cst 对全局有效

    协作组合

    常见的几种协作组合:

    • 宽松顺序:memory_order_relaxed

      仅保证原子性,允许所有指令重排

    • 释放 - 获得序:memory_order_release+memory_order_acquire

      release 前的写操作一定到位,acquire 的读操作一定最新

    • 释放 - 消费序:memory_order_release+ memory_order_consume

      比释放获得序宽松一些,仅同步了原子操作本身的原子变量及与它产生依赖关系的变量

    • 序列一致顺序 memory_order_seq_cst

      拒绝一切重排,对所有线程可见

    举个例子:

    首先定义两个变量

    std::atomic<std::string *> ptr;
    int data;
    

    对于释放获得序:

    void producer() {
      auto *p = new std::string("Hello");
      data = 42;
      // 编译器不能将 data = 42 这条指令移动到 store 操作之后
      // memory_order_release 将修改后(写操作)的结果释放出来,其他线程使用 memory_order_acquire 可以观测到上述对内存写入的结果。此处保证 producer 的修改一定写入
      ptr.store(p, std::memory_order_release);
    }
    
    void consumer() {
      std::string *p2;
      // memory_order_acquire 读取到最新值,其他线程使用 memory_order_release 前的值一定可以读到。此处保证 consumer 一定读取到 producer 写入的值
      while (!(p2 = ptr.load(std::memory_order_acquire)))
        ;
      // 执行到此处说明 p2 是非空的,即 data = 42 的指令已经执行了(memory_order_release 保证)
      // 且此时 data 必定等于 42,p2 必定为“Hello”(memory_order_acquire 保证)
      assert(*p2 == "Hello");
      assert(data == 42); 
      // 不会出错
    }
    

    对于释放消费序:

    void producer() {
        std::string* p  = new std::string("Hello");
        data = 42;
        ptr.store(p, std::memory_order_release);
    }
     
    void consumer() {
        std::string* p2;
        while (!(p2 = ptr.load(std::memory_order_consume)))
            ;
        assert(*p2 == "Hello"); // 不会出错:*p2 从 ptr 携带依赖
        assert(data == 42); // 有可能出错:data 不从 ptr 携带依赖
    }
    

    内存屏障

    c++ 标准库中提供接口

    • std::atomic_thread_fence:显式的内存屏障。与内存顺序标记配合使用,将线程分为若干个阶段,每个阶段需要等待其他线程完成特定操作之后才能进入下一阶段。对于所有数据有效,不仅仅是原子变量。

    使用时可以将原子变量都设置为 memory_order_relaxed,仅依靠内存屏障来实现对内存的控制

    例子:

    std::atomic<bool> x,y;
    void write_x_then_y() {
        x.store(true, std::memory_order_relaxed); // 1
        std::atomic_thread_fence(std::memory_order_release); // 2.a 写屏障
        y.store(true, std::memory_order_relaxed); // 2.b
    }
    void read_y_then_x() {
        while(!y.load(std::memory_order_relaxed)); // 3.a
        std::atomic_thread_fence(std::memory_order_acquire); // 3.b 读屏障
        assert(!x.load(std::memory_order_relaxed));// 4
    }
    

    可以认为:

    • 2.a 写屏障向下和 2.b 写操作相结合,将 2.b 提升为 release-store
    • 3.b 读屏障向上和 3.a 读操作相结合,将 3.a 提升为 acquire-load
    • 同步关系与前面举的内存顺序的例子相似,保证了 1 一定先于 4

    内存屏障与编译器屏障

    编译器屏障仅影响编译器。防止编译器重排指令,确保最终生成的汇编代码保留了源代码中指定的顺序,但不会生成实际的硬件指令,因此不会影响 CPU 或内存的访问顺序。不能处理线程并发带来的缓存一致性、内存可见性问题,常用于单核处理器。

    内存屏障同时影响编译器与 CPU。不仅防止编译器重排操作,还会生成硬件指令影响 CPU 的内存访问顺序。确保线程之间的数据同步,常用于多核处理器。

    编译器屏障

    • std::atomic_signal_fence :编译器屏障,对所有数据有效。用于防止编译器在信号处理程序或与硬件交互时重排指令。
    • volatile 关键字
      • 对于某个 volatile 变量,与该变量相关的操作必须对内存进行读写,而不能仅仅读写寄存器中的缓存值。
      • 对于多个 volatile 变量,它们之间的访问顺序不会被编译器改变

    编译器屏障相当于告诉编译器,这些变量的值可能会在程序的控制之外发生变化,不可以假设其不变。常用于信号处理、硬件交互、嵌入式系统中断程序。

    举个例子:

    // thread1
    volatile bool tag = true;
    volatile int x = 0;
    while(tag);
    assert(!x);
    
    // thread2
    x = 10;
    tag = false;
    

    在单核处理器上“伪并发”这两个线程,volatile 在此处有两个作用:

    • 每次都从内存中读取 tag 最新值,避免 thread1 死循环
    • x 和 tag 的读写顺序和源代码中一致,即 x=10 一定发生在 tag=false 之前

    但在多核处理器上“真并发”这两个线程时,即使 volatile 强制编译器生成了从主存读写的代码,具体执行顺序仍然由 CPU 决定。此时就需要使用内存屏障来实现需求的逻辑,保证内存操作按照指定顺序执行。

    Post #20 ❤️ 4 likes