본문으로 건너뛰기
Folly Code Review · 46/89

folly::MPMCQueue — multi-producer multi-consumer

· Hawk · 6분 읽기

#한 줄 요약

folly::MPMCQueue<T>는 N명의 producer와 M명의 consumer를 가정한 ticket-based lock-free 큐다. push/pop마다 fetch_add로 ticket을 발급하고, 슬롯별 sequence 비교로 wait/notify를 흉내낸다. 한 슬롯당 mutex 하나 없이 strict FIFO를 보장한다.

#동기 — 왜 ticket인가

SPSC와 달리 MPMC는 본질적으로 경쟁이 있다. 단순히 mutex로 잠그면 producer 수와 throughput이 반비례한다. CAS 루프로 head/tail을 갱신하는 방식도 있지만, CAS 실패가 누적되면 thundering herd가 생긴다.

ticket-based 알고리즘은 다른 접근이다. 모든 producer가 fetch_add(1)로 자기 ticket(슬롯 인덱스)을 받고, 받은 슬롯에 단독으로 쓴다. 다른 producer와 같은 슬롯에서 충돌하지 않으므로 CAS 루프가 없다. 슬롯이 비어 있는지(consumer가 비웠는지)만 슬롯 내부의 sequence 카운터로 확인한다.

이 알고리즘은 Vyukov의 bounded MPMC queue로 잘 알려져 있고, Folly는 그것을 cache-line 정렬과 적응형 spin/block까지 더해 production-ready로 다듬었다.

#Producer/Consumer 그림

MPMC는 producer N + consumer M의 가장 일반적인 형태다.

Producer / consumer queue

여러 생산자/소비자가 공유 큐를 두고 경쟁한다. capacity가 차면 producer가 block, 비면 consumer가 block — 양방향 backpressure가 자연스럽게 생긴다.

#API

#include <folly/MPMCQueue.h>
folly::MPMCQueue<Task> q(1024);
// Producer (N명)
q.blockingWrite(task); // 큐 가득 차면 대기
q.write(task); // non-blocking, false on full
q.writeIfNotFull(task); // 같음
// Consumer (M명)
Task t;
q.blockingRead(t); // 비었으면 대기
q.read(t); // non-blocking, false on empty

folly::DynamicBoundedQueuefolly::UnboundedQueue도 같은 가족인데, MPMCQueue는 그 중 가장 단순한 bounded 버전이다.

#내부 구현 — slot의 sequence

MPMC ticket-based queue

struct Slot {
std::atomic<uint32_t> sequence;
T data;
};
Slot* slots_;
size_t capacity_;
alignas(cacheline) std::atomic<uint64_t> pushTicket_;
alignas(cacheline) std::atomic<uint64_t> popTicket_;

각 슬롯에 sequence가 있다. 초기값은 슬롯 인덱스다(slot[0].sequence = 0, slot[1].sequence = 1, …).

alignas(cacheline) 필요한지 짚어 보자. push 쪽과 pop 쪽이 같은 cache line에 있으면:

False sharing on cache line

producer가 pushTicket_++만 해도 consumer의 cache line이 invalidate되어 popTicket_ 읽기에 cache miss가 난다 — 두 변수가 독립인데도 그렇다. cache-line aligned로 두 변수를 분리하면 ping-pong이 사라진다.

#enqueue 알고리즘

void blockingWrite(T&& v) {
auto ticket = pushTicket_.fetch_add(1, std::memory_order_acq_rel);
auto idx = ticket % capacity_;
auto& slot = slots_[idx];
// 이 슬롯의 sequence가 ticket과 같아질 때까지 대기
while (slot.sequence.load(std::memory_order_acquire) != ticket) {
spin_or_park();
}
new (&slot.data) T(std::move(v));
slot.sequence.store(ticket + 1, std::memory_order_release);
}

해석하자면.

  1. pushTicket_에서 자기 ticket을 받는다. 이건 fetch_add라 절대 충돌 안 한다.
  2. ticket % capacity로 슬롯을 정한다.
  3. 슬롯의 sequence == ticket이면 “내가 쓸 차례”. 그렇지 않으면 이전 round의 consumer가 아직 안 비웠다는 뜻이라 기다린다.
  4. 데이터 쓰고 sequence를 ticket+1로 올린다 → consumer에게 “여기 데이터 있음” 신호.

#dequeue

void blockingRead(T& v) {
auto ticket = popTicket_.fetch_add(1, std::memory_order_acq_rel);
auto idx = ticket % capacity_;
auto& slot = slots_[idx];
// sequence가 ticket+1이 될 때까지 대기 (producer 완료 신호)
while (slot.sequence.load(std::memory_order_acquire) != ticket + 1) {
spin_or_park();
}
v = std::move(slot.data);
slot.data.~T();
slot.sequence.store(ticket + capacity_, std::memory_order_release);
}

consumer는 ticket+1을 기다리고, 다 읽으면 sequence를 ticket+capacity로 올린다. 그러면 다음 round의 producer가 그 슬롯에 쓸 수 있다.

#적응형 spin

spin_or_park은 처음 몇 번은 spin(_mm_pause), 그 후에도 안 풀리면 Baton::wait로 park한다. 짧은 경합은 spin이 빠르고, 긴 경합은 park가 CPU를 살린다. tuning constant(SpinCount)는 워크로드별로 조정 가능하다.

#std / abseil 비교

동시성크기비고
std::queue + mutex + cvMPMCyes동적단순. baseline.
boost::lockfree::queueMPMCno둘 다CAS-based. ABA 카운터
folly::MPMCQueueMPMCnoboundedticket-based, 더 빠름
folly::DynamicBoundedQueueMPMCnoboundedMPMCQueue + 동적 capacity
folly::UnboundedQueueMPMCnounbounded다음 절

표준에는 lock-free MPMC가 없다. boost::lockfree::queue는 CAS-based여서 경합 시 retry가 발생한다. Folly의 ticket 방식은 retry가 없고, 슬롯별 spin/park만 있어 throughput이 안정적이다. Meta의 RPC 백엔드 큐는 거의 모두 MPMCQueue로 통일되어 있다.

#코드 리뷰 포인트

#1. capacity는 2의 거듭제곱일 필요는 없다

내부 구현은 ticket % capacity로 인덱스를 정한다. capacity가 2^n이면 컴파일러가 & (capacity-1)로 최적화하지만, 그게 아니라도 동작한다. 메모리 사용량이 더 중요하면 정확한 capacity를 잡는다.

#2. blockingWrite vs writeIfNotFull

// blockingWrite: 가득 차면 영원히 대기
q.blockingWrite(task);
// writeIfNotFull: 가득 차면 즉시 false
if (!q.writeIfNotFull(task)) {
rejectTask(task);
}

production 시스템에서 blockingWrite는 위험하다. consumer가 죽으면 producer가 그대로 멈춘다. backpressure 정책을 명시하고 가능하면 writeIfNotFull + 명시적 reject로 처리한다.

#3. 큐 이름과 capacity로 메모리 추정

folly::MPMCQueue<LargeStruct> q(100000);
// 메모리 = 100000 * (sizeof(LargeStruct) + sizeof(uint32_t)) + 알파

bounded 큐는 생성 시점에 전체 메모리를 확보한다. T가 크고 capacity가 크면 수백 MB를 한 번에 잡는다. PR 리뷰에서 capacity 숫자가 적절한지 확인한다.

#4. fairness — strict FIFO

ticket 알고리즘은 FIFO를 정확히 보장한다. 어떤 producer도 추월할 수 없고, 어떤 consumer도 새치기할 수 없다. 이게 장점이지만, ticket을 받은 후에 죽거나 멈춘 스레드가 있으면 그 슬롯에서 모든 후속 작업이 막힐 수 있다. consumer가 hang하지 않도록 watchdog 필요.

#안티패턴

#1. capacity를 작게 잡고 backpressure 없이 blockingWrite

// 회피
folly::MPMCQueue<Task> q(16); // 너무 작음
for (auto& task : tasks_1000) {
q.blockingWrite(std::move(task)); // 17번째부터 모두 block
}

producer가 consumer를 추월하는 워크로드에서 capacity 부족 + blocking은 thread starvation을 만든다. 큐 사용량 모니터링 + writeIfNotFull로 reject path를 둔다.

#2. consumer 쪽에서 무거운 작업

// 회피
Task t;
q.blockingRead(t);
heavy_cpu_work(t); // 다음 read 차례인 ticket 보유자가 영원히 대기

consumer가 read 후 무거운 작업을 동기로 하면 다른 consumer가 자기 ticket 슬롯에 들어가지 못한다. read는 빠르게 마치고 작업은 별도 executor로 보낸다.

#3. MPMCQueue를 SPSC 자리에 쓰기

// 회피 — 동작은 하지만 5배 느림
folly::MPMCQueue<int> q(1024);
// producer 1명, consumer 1명만

SPSC 패턴이 확실하면 ProducerConsumerQueue가 훨씬 빠르다.

#정리

  • MPMCQueue는 ticket-based lock-free MPMC 큐다.
  • 슬롯별 sequence로 wait/notify를 흉내낸다. CAS 루프가 없다.
  • 짧은 경합은 spin, 긴 경합은 Baton park로 자동 전환.
  • bounded이므로 capacity 정책과 backpressure를 항상 함께 설계한다.
  • strict FIFO 보장. consumer hang은 큐 전체를 막을 수 있다.
  • SPSC 패턴이면 ProducerConsumerQueue로 다운그레이드한다.

#다음 편

Part 10-03 UnboundedQueue — capacity 고정 없는 동적 lock-free 큐. linked segment로 어떻게 성장하는지 본다.

#관련 항목

Folly Code Review · 47 of 89

  1. 1 Folly Code Review — Meta의 production-grade C++ 라이브러리 코드 분석
  2. 2 Folly 개요 — Meta가 production에서 검증한 utility 모음 분석
  3. 3 Folly vs Abseil 철학 비교 — performance-first vs std-compatible
  4. 4 Folly 빌드와 fbcode 환경 — monorepo의 그림자
  5. 5 Folly API stability 정책 — 어떤 보장도 없다는 솔직함
  6. 6 Folly production validation 문화 — peta-scale에서 단련된 코드
  7. 7 folly::Future 분석 — std::future의 한계를 넘는 composable async
  8. 8 folly::Promise·makeFuture — Future를 만드는 두 길
  9. 9 folly::SemiFuture vs Future — executor binding의 명시화
  10. 10 folly::Future thenValue·thenError·thenTry — continuation 체인 분석
  11. 11 folly::collect·collectAll·collectAny — fan-in 패턴 분석
  12. 12 folly::Future retry·window·via — 제어 흐름 조합자
  13. 13 folly::fibers 분석 — M:N stackful coroutine
  14. 14 folly::InlineExecutor — 호출자 thread에서 즉시 실행
  15. 15 folly::CPUThreadPoolExecutor — CPU-bound 작업의 표준 thread pool
  16. 16 folly::IOThreadPoolExecutor — libevent 기반 I/O pool
  17. 17 folly::ManualExecutor — 결정적 테스트를 위한 수동 진행
  18. 18 folly::EventBase 분석 — libevent 이벤트 루프의 핵심
  19. 19 folly::IOBuf 분석 — zero-copy buffer chain의 기본 단위
  20. 20 folly::IOBufQueue — chain의 push/pull 추상화
  21. 21 folly::io::Cursor·RWCursor — chain 위의 stream
  22. 22 folly Zero-copy 패턴 — IOBuf로 ScatterGather I/O 표현
  23. 23 folly::IOBuf shared semantics — clone·unshare·takeOwnership
  24. 24 folly::FBString 분석 — SSO + COW 구현
  25. 25 folly의 fmt::format 통합 — 모던 포맷팅 채택
  26. 26 folly::StringPiece — string_view 호환 분석
  27. 27 folly Join·Split utilities — 문자열 분해와 결합
  28. 28 folly::to·tryTo — text↔num 변환 분석
  29. 29 folly Conv Customization — 사용자 타입 지원
  30. 30 folly Conv 성능 비교 — sprintf·stringstream 대비
  31. 31 folly::F14ValueMap vs std::unordered_map
  32. 32 folly::F14NodeMap — stable pointer가 필요할 때
  33. 33 folly::F14VectorMap — cache-friendly iteration
  34. 34 folly::F14FastMap — auto-select 동작
  35. 35 folly F14 internals — SIMD probing 메커니즘
  36. 36 folly::small_vector — inline storage 분석
  37. 37 folly::FixedString — compile-time string
  38. 38 folly::AtomicHashMap — lock-free read 분석
  39. 39 folly::ConcurrentHashMap — sharded 동시 해시 맵
  40. 40 folly::EvictingCacheMap — LRU 구현 분석
  41. 41 folly::Synchronized — lock wrapper 패턴
  42. 42 folly::SharedMutex 분석
  43. 43 folly::Baton — one-shot wait 동기화
  44. 44 folly::RWSpinLock 분석
  45. 45 folly::PicoSpinLock — 1-byte spinlock
  46. 46 folly::ProducerConsumerQueue — SPSC 큐 분석
  47. 47 folly::MPMCQueue — multi-producer multi-consumer
  48. 48 folly::UnboundedQueue — 동적 크기 lock-free
  49. 49 folly::fibers::Channel — Go-like channel
  50. 50 folly::dynamic — JSON-like dynamic type 분석
  51. 51 folly JSON conversion — toJson·parseJson
  52. 52 folly dynamic ↔ struct — manual marshaling
  53. 53 folly dynamic Visitor pattern — type별 분기
  54. 54 folly::Singleton vs Meyers/static — 왜 Folly의 Singleton인가
  55. 55 folly::SingletonVault 분석 — 등록·소멸·의존성
  56. 56 folly::Singleton try_get·try_get_fast — TLS-cached 접근
  57. 57 folly::ExceptionWrapper — type-erased exception holder
  58. 58 folly::ScopeGuard·SCOPE_EXIT — RAII cleanup
  59. 59 folly::Optional vs std::optional
  60. 60 folly::Function vs std::function
  61. 61 folly::Lazy — 지연 초기화 wrapper
  62. 62 folly Meta 스타일 code review 패턴
  63. 63 folly anti-patterns — 잘못 쓰면 std보다 느림
  64. 64 folly vs std 선택 기준 분석
  65. 65 folly::coro 개요 — production C++20 코루틴 어댑터
  66. 66 folly::coro::Task — lazy single-shot 코루틴
  67. 67 folly::coro::AsyncGenerator — 비동기 스트림
  68. 68 folly coro blockingWait·collectAll — 동기 경계와 fan-in
  69. 69 folly::coro::Baton·Mutex — 코루틴-aware 동기화
  70. 70 folly::Expected — 결과 또는 오류
  71. 71 folly::Try — Future 결과 wrapper
  72. 72 folly::Try vs Expected 선택 기준
  73. 73 folly::Range — 일반 iterator pair
  74. 74 folly::Uri — URL 파서
  75. 75 folly Fingerprint64·128 — 분산 hash
  76. 76 folly SpookyHashV2 — fast non-crypto hash
  77. 77 folly::Init — main() 부트스트랩
  78. 78 folly::Indestructible — global lifetime 패턴
  79. 79 folly::MicroLock — 1-byte 락
  80. 80 folly::MicroSpinLock — 가장 좁은 spin lock
  81. 81 folly::format — legacy formatter 분석
  82. 82 folly::demangle — typeid 디망글링
  83. 83 folly::DynamicConverter — dynamic ↔ struct
  84. 84 folly::RecordIO — append-only 로그 파일 포맷
  85. 85 folly::io::Compression — zstd·lz4·snappy wrapper
  86. 86 folly::AsyncIO — io_uring·Linux AIO
  87. 87 folly::CancellationToken — 코루틴·Future 취소 전파
  88. 88 folly::observer — hot config의 atomic refresh
  89. 89 fbcode 패턴 모음 — folly 사용의 실전