💡 SPSC란?

Single-Producer / Single-Consumer 의 약자로 생산자와 소비자가 1:1인 관계를 뜻합니다.

리눅스에서 소켓을 다루다가 스레드간 데이터 전달이 필요해서 알아보니 SPSC Queue가 있다는것을 알게 되었습니다. 제 경우는 클라이언트의 연결을 담당하는는스레드에서 데이터 수신을 담당하는 스레드로 소켓(fd)을 전달해야 했는데 공용 컨테이너와 뮤텍스를 사용하면 스레드가 늘어나면 늘어날수록 오버헤드가 늘 것 같았습니다.

그래서 연결을 담당하는 스레드와 데이터를 수신하는 스레드의 생산-소비가 1:1인 관계에서 뮤텍스 없이 Lock-free 하게 구현해 보려고 합니다. 👍

🚀 구현

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
// spsc_queue.hpp

#ifndef __SPSC_QUEUE_HPP__
#define __SPSC_QUEUE_HPP__

#include <atomic>
#include <cstddef>

// 타입 일반화를 위한 template
template <typename T, std::size_t size>
class SPSCQueue {
public:
    // 해당 객체의 복사/이동 금지
    SPSCQueue(const SPSCQueue &) = delete;
    SPSCQueue& operator=(const SPSCQueue&) = delete;
    SPSCQueue(SPSCQueue&&) = delete;
    SPSCQueue& operator=(SPSCQueue&&) = delete;

    // 기본생성자
    SPSCQueue()
        : _capacity(pow2RoundUp(size))
        , _top(0), _bottom(0) {
        _chunk = new T[_capacity];
    }
    virtual ~SPSCQueue() {
        delete[] _chunk;
    }

    // Forwarding reference
    template <typename U>
    bool push(U&& data) {
        std::size_t top = _top.load(std::memory_order_relaxed);
        std::size_t bottom = _bottom.load(std::memory_order_acquire);
        if(top - bottom >= _capacity) {
            return false;
        }
        
        _chunk[top & (_capacity - 1)] = std::forward<U>(data);
        _top.store(top + 1, std::memory_order_release);

        return true;
    }

    bool pop(T& data) {
        std::size_t top = _top.load(std::memory_order_acquire);
        std::size_t bottom = _bottom.load(std::memory_order_relaxed);
        if(top - bottom == 0) {
            return false;
        }

        data = std::move(_chunk[bottom & (_capacity - 1)]);
        _bottom.store(bottom + 1, std::memory_order_release);

        return true;
    }

    std::size_t getCapacity() { return _capacity; }
    std::size_t getSize() { 
        return _bottom.load(std::memory_order_relaxed) - 
        _top.load(std::memory_order_relaxed);
    }

private:
    T* _chunk;
    std::size_t _capacity;
    std::atomic<std::size_t> _top, _bottom;
줌
    // 2^x 형태로 만들어줌
    std::size_t pow2RoundUp(std::size_t number) {
        number--;
        number |= number >> 1;
        number |= number >> 2;
        number |= number >> 4;
        number |= number >> 8;
        number |= number >> 16;
        number |= number >> 32;

        return ++number;
    }
};

#endif //__SPSC_QUEUE_HPP__

💡 push() pop() 함수가 핵심!

생선과 소비의 관계가 1:1 이므로 생산자는 push 할 때 _top 인덱스만 사용하고 소비자는 pop 할 때 _bottom 인덱스만 사용합니다. 물론 큐가 찼는지, 비었는지 확인하기 위해 _top _bottom 모두 확인은 합니다.

push() 생산자 입장에서 생각해보기

생산자가 데이터를 넣고나서 변경(write)해야 할 인덱스(변수)는 _top입니다. 또한 _top를 변경하는 다른 스레드는 없습니다.

즉 _top 인덱스 변수는 생산자 스레드에서 항상 최신의 데이터를 유지할 것이므로 생산자에서 _top을 읽을 때 memory_order_relaxed로 읽어도 무방할 것입니다.

그러나 큐가 꽉 찼는지 확인하기 위해 _bottom을 읽어야 하는데 다른 스레드(소비자)가 변경할 수 있는 데이터입니다. 따라서 _bottom은 최신 내용을 받으려면 memory_order_acquire로 읽어야 합니다.

pop() 소비자 입장에서 생각해보기

생산자와 정반대입니다. 소비자가 데이터를 소비하고 나서 변경해야 할 인덱스는 _bottom이고, 다른 소비자가 없기 때문에 소비자 입장에서 항상 최신화 되므로 memory_order_relaxed로 읽어도 무방합니다.

그러나 큐가 비었는지 확인하기 위해 확인해야 할 인덱스 _top은 생산자에 의해 변경될 수 있고 최신 내용을 받기 위해 memory_order_acquire로 읽습니다.

CPU 명령어 처리 또는 컴파일러 최적화로 인한 재배치, memory_order 관련 내용은 여기서 자세히 볼 수 있습니다.

🔬 테스트

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
TEST_CASE(SPSC_QUEUE_TEST) {
    Yami::TaskExecutor executor;
    Yami::Lockfree::SPSCQueue<std::size_t, 1024> q;
    const std::size_t kExpectSize = 1024;
    const std::size_t kRepeatCount = 10000000;
    std::size_t tmp = 0;
    
    // RoundUp test
    ASSERT_EQ(q.getCapacity(), kExpectSize);

    // Full test
    for(std::size_t i=0; i<kExpectSize; i++) {
        q.push(i);
    }
    ASSERT_TRUE(!q.push(1234));

    // Pop test
    for(std::size_t i=0; i<kExpectSize; i++) {
        ASSERT_TRUE(q.pop(tmp));
    }

    // Empty test
    ASSERT_TRUE(!q.pop(tmp));

    // SPSC test
    executor.run([&]() {
        // only produce
        std::size_t count = 0;
        while(count < kRepeatCount) {
            while(!q.push(count)) {
                std::this_thread::sleep_for(std::chrono::milliseconds(1));
            }
            count++;
        }
        
    });
    executor.run([&]() {
        // only consume
        std::size_t count = 0;
        std::size_t tmp = 0;
        while(count < kRepeatCount) {
            while(!q.pop(tmp)) {
                std::this_thread::sleep_for(std::chrono::milliseconds(1));
            }
            count++;
        }
    }).wait();
}