💡 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();
}
|