C++ 작업 큐 설계: 워크 스틸링으로 스레드 풀 성능 끌어올리기
들어가며: 메인 스레드에서 무거운 일을 하면
이미지 편집 앱에서 “썸네일 100장 생성” 버튼을 누르면 메인 스레드가 이미지 변환을 시작한다고 해 봅시다. 변환이 끝날 때까지 이벤트 루프가 돌지 못하므로 클릭·스크롤·드래그가 모두 무시되고, 운영체제는 창에 “응답 없음”을 표시합니다. 이미지 변환뿐 아니라 로그 파일 쓰기, HTTP 요청, 게임 서버의 경로 탐색처럼 시간이 걸리는 작업은 모두 같은 문제를 일으킵니다. 서버라면 요청 하나를 처리하는 스레드가 디스크 I/O를 기다리는 동안 다른 요청을 받지 못해 처리량이 떨어집니다.
해결책은 작업 큐와 스레드 풀입니다. 메인 스레드는 “할 일”을 큐에 넣고 바로 돌아오고, 미리 만들어 둔 몇 개의 워커 스레드가 큐에서 작업을 꺼내 백그라운드에서 실행합니다.
flowchart LR
subgraph 메인스레드
A[사용자 클릭] --> B[작업 큐에 push]
B --> C[즉시 반환]
end
subgraph 워커스레드들
D[워커 1] --> E[pop & 실행]
F[워커 2] --> E
G[워커 3] --> E
end
B -.->|notify| D
B -.->|notify| F
B -.->|notify| G
작업 큐와 스레드 풀은 7편 condition_variable의 생산자-소비자 패턴을 그대로 활용합니다. 이 글에서는 std::function과 람다로 작업을 담는 스레드 안전 큐를 만들고, std::thread·std::mutex·std::condition_variable로 스레드 풀을 구성한 뒤, std::future로 결과를 돌려받는 방법, 부하 불균형을 줄이는 워크 스틸링, 안전한 종료 처리까지 이어서 구현합니다.
메인 스레드 직접 실행에서 작업 큐로
Before: 메인 스레드에서 직접 실행
// 메인 스레드가 끝날 때까지 블로킹됨
void onGenerateThumbnailsClicked() {
for (const auto& path : imagePaths) {
auto thumbnail = resizeImage(path, 128, 128); // 장당 수십 ms라고 가정
saveThumbnail(thumbnail);
}
// 장당 50ms라면 100장 × 50ms = 5초 동안 UI가 멈춤
}
After: 작업 큐 + 스레드 풀
// 작업을 큐에 넣고 즉시 반환
void onGenerateThumbnailsClicked() {
for (const auto& path : imagePaths) {
pool.submit([path]() {
auto thumbnail = resizeImage(path, 128, 128);
saveThumbnail(thumbnail);
});
}
}
메인 스레드는 submit만 하고 바로 반환하므로 UI가 계속 반응합니다. 워커 여러 개가 동시에 처리하므로, CPU 코어가 남아 있다면 전체 완료 시간도 줄어듭니다. 다만 작업이 끝났다는 사실을 UI에 알리려면 결과를 다시 메인 스레드로 넘기는 장치(GUI 프레임워크의 이벤트 포스트 함수 등)가 필요합니다. 워커 스레드에서 UI 객체를 직접 건드리면 대부분의 GUI 프레임워크에서 미정의 동작입니다.
할 일을 함수 객체로 표현하는 작업 큐
작업 타입을 std::function<void()>로 두면, 인자와 반환값 없이 “한 번 실행하면 끝나는 작업”을 람다나 함수 객체로 표현할 수 있습니다. 필요한 데이터는 람다 캡처로 함께 담습니다. 워커는 큐에서 작업을 꺼내 task()로 호출하기만 하면 됩니다.
sequenceDiagram participant M as 메인 스레드 participant Q as TaskQueue participant W1 as 워커 1 participant W2 as 워커 2 M->>Q: push(task1) M->>Q: push(task2) Q->>W1: notify_one W1->>Q: waitAndPop W1->>W1: task1() M->>Q: push(task3) Q->>W2: notify_one W2->>Q: waitAndPop W2->>W2: task2()
뮤텍스와 조건 변수로 만든 스레드 안전 큐
TaskQueue는 여러 스레드가 동시에 push·pop해도 안전한 큐입니다. push는 뮤텍스로 큐를 보호한 채 작업을 넣고, 락을 푼 뒤 notify_one()으로 대기 중인 워커 하나를 깨웁니다. waitAndPop은 condition_variable::wait에 “큐가 비어 있지 않거나 종료되었다”는 조건을 넘겨, 작업이 들어오거나 shutdown이 호출될 때만 깨어납니다. 조건자를 넘기는 형태의 wait는 깨어날 때마다 조건을 다시 검사하므로, 이유 없이 깨어나는 spurious wakeup도 자동으로 처리됩니다.
#include <queue>
#include <mutex>
#include <condition_variable>
#include <functional>
using Task = std::function<void()>;
class TaskQueue {
public:
void push(Task task) {
{
std::lock_guard<std::mutex> lock(mutex_);
if (done_) return; // shutdown 후 push 무시
queue_.push(std::move(task));
}
cv_.notify_one();
}
bool tryPop(Task& out) { // 기다리지 않는 pop
std::lock_guard<std::mutex> lock(mutex_);
if (queue_.empty()) return false;
out = std::move(queue_.front());
queue_.pop();
return true;
}
bool waitAndPop(Task& out) {
std::unique_lock<std::mutex> lock(mutex_);
cv_.wait(lock, [this] { return !queue_.empty() || done_; });
if (queue_.empty()) { // 여기 왔다면 done_ == true
return false;
}
out = std::move(queue_.front());
queue_.pop();
return true;
}
void shutdown() {
{
std::lock_guard<std::mutex> lock(mutex_);
done_ = true;
}
cv_.notify_all();
}
private:
std::queue<Task> queue_;
std::mutex mutex_;
std::condition_variable cv_;
bool done_ = false;
};
waitAndPop은 done_이 켜져도 큐에 작업이 남아 있으면 계속 꺼내 주고, 큐가 비었을 때만 false를 돌려줍니다. 그래서 이 큐를 쓰는 풀은 종료 시 이미 넣은 작업을 모두 처리한 뒤 끝납니다. done_은 항상 뮤텍스를 잡은 상태에서만 읽고 쓰므로 atomic일 필요가 없습니다. 반대로 락 없이 done_을 읽는 isDone() 같은 함수를 추가하면 데이터 레이스가 되므로 주의해야 합니다.
워커가 큐에서 꺼내 실행하는 스레드 풀
ThreadPool은 생성 시 지정한 개수만큼 스레드를 만들고, 각 스레드는 waitAndPop으로 작업을 기다렸다가 실행하는 루프를 돕니다. 소멸자는 큐를 종료시킨 뒤 모든 워커를 join합니다.
#include <thread>
#include <vector>
class ThreadPool {
public:
explicit ThreadPool(size_t numThreads) {
workers_.reserve(numThreads);
for (size_t i = 0; i < numThreads; ++i) {
workers_.emplace_back([this]() { workerLoop(); });
}
}
~ThreadPool() {
queue_.shutdown();
for (auto& w : workers_) {
if (w.joinable()) w.join();
}
}
void submit(Task task) {
queue_.push(std::move(task));
}
size_t workerCount() const { return workers_.size(); }
private:
void workerLoop() {
Task task;
while (queue_.waitAndPop(task)) {
try {
task();
} catch (...) {
// 예외가 워커 스레드 밖으로 나가면 std::terminate가 호출됨
// 실무에서는 여기서 로깅
}
}
}
TaskQueue queue_;
std::vector<std::thread> workers_;
};
멤버 선언 순서도 중요합니다. 생성자 본문에서 워커가 queue_를 쓰기 시작하므로 queue_가 workers_보다 먼저 선언(먼저 생성)되어 있어야 합니다. try/catch는 선택이 아닙니다. 스레드 함수 밖으로 예외가 빠져나가면 워커 하나가 죽는 정도가 아니라 프로세스 전체가 std::terminate로 종료됩니다.
실행 가능한 전체 예제
아래 코드는 TaskQueue와 ThreadPool을 포함한 단일 파일로, g++ -std=c++17 -O2 -pthread -o taskqueue taskqueue.cpp로 컴파일해 실행할 수 있습니다.
// taskqueue.cpp - 작업 큐 + 스레드 풀 예제
#include <iostream>
#include <queue>
#include <mutex>
#include <condition_variable>
#include <functional>
#include <thread>
#include <vector>
#include <chrono>
using Task = std::function<void()>;
class TaskQueue {
public:
void push(Task task) {
{
std::lock_guard<std::mutex> lock(mutex_);
if (done_) return;
queue_.push(std::move(task));
}
cv_.notify_one();
}
bool waitAndPop(Task& out) {
std::unique_lock<std::mutex> lock(mutex_);
cv_.wait(lock, [this] { return !queue_.empty() || done_; });
if (queue_.empty()) return false;
out = std::move(queue_.front());
queue_.pop();
return true;
}
void shutdown() {
{
std::lock_guard<std::mutex> lock(mutex_);
done_ = true;
}
cv_.notify_all();
}
private:
std::queue<Task> queue_;
std::mutex mutex_;
std::condition_variable cv_;
bool done_ = false;
};
class ThreadPool {
public:
explicit ThreadPool(size_t n) {
for (size_t i = 0; i < n; ++i) {
workers_.emplace_back([this]() {
Task task;
while (queue_.waitAndPop(task)) {
try { task(); } catch (...) { /* 로깅 */ }
}
});
}
}
~ThreadPool() {
queue_.shutdown();
for (auto& w : workers_) if (w.joinable()) w.join();
}
void submit(Task task) { queue_.push(std::move(task)); }
private:
TaskQueue queue_;
std::vector<std::thread> workers_;
};
std::mutex coutMutex;
int main() {
auto start = std::chrono::steady_clock::now();
{
ThreadPool pool(4);
for (int i = 0; i < 20; ++i) {
pool.submit([i]() {
std::this_thread::sleep_for(std::chrono::milliseconds(50));
std::lock_guard<std::mutex> lock(coutMutex);
std::cout << "Task " << i << " on " << std::this_thread::get_id() << "\n";
});
}
} // 소멸자: 남은 작업을 모두 처리한 뒤 join
auto ms = std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now() - start).count();
std::cout << "Total: " << ms << " ms\n";
return 0;
}
작업 20개를 각각 50ms씩 순차 실행하면 1000ms가 걸리지만, 워커 4개가 나눠 처리하면 이론적으로 5라운드, 즉 250ms 남짓에 끝납니다. 시간 측정은 풀의 소멸자가 모든 작업을 마치고 join한 뒤에 하므로 실제 처리 시간을 반영합니다. std::cout 출력은 여러 스레드에서 섞이지 않도록 별도 뮤텍스로 보호했습니다.
std::promise와 std::future로 결과 받기
작업이 값을 반환해야 할 때는 std::promise<T>와 std::future<T>를 씁니다. submitWithResult는 promise를 만들어 그 future를 호출자에게 돌려주고, 워커가 실행하는 람다 안에서 결과를 set_value로, 예외를 set_exception으로 전달합니다. 호출자는 fut.get()으로 결과를 기다렸다가 받거나, 작업에서 난 예외를 그 자리에서 다시 받습니다.
#include <future>
#include <memory>
#include <type_traits>
// ThreadPool의 public 멤버로 추가
template<typename F>
auto submitWithResult(F func) -> std::future<std::invoke_result_t<F>> {
using R = std::invoke_result_t<F>;
auto prom = std::make_shared<std::promise<R>>();
auto fut = prom->get_future();
submit([prom, func = std::move(func)]() mutable {
try {
if constexpr (std::is_void_v<R>) {
func();
prom->set_value();
} else {
prom->set_value(func());
}
} catch (...) {
prom->set_exception(std::current_exception());
}
});
return fut;
}
// 사용
ThreadPool pool(4);
auto fut = pool.submitWithResult([]() { return 42; });
std::cout << fut.get() << "\n"; // 42
std::promise는 이동만 가능한 타입인데, std::function은 담는 함수 객체가 복사 가능해야 합니다. 그래서 promise를 shared_ptr로 감싸 람다가 복사 가능하도록 만들었습니다. 반환 타입이 void인 작업은 set_value(func())로 쓸 수 없으므로 if constexpr로 나눴습니다. 같은 일을 하는 표준 도구로 std::packaged_task도 있지만, 역시 이동 전용이라 std::function 기반 큐에 넣으려면 같은 방식으로 shared_ptr에 담아야 합니다.
부하 불균형을 줄이는 워크 스틸링
문제: 단일 큐의 경합과 불균형
모든 워커가 하나의 큐와 하나의 뮤텍스를 공유하면, 작업이 아주 짧고 많을 때 워커들이 같은 락을 두고 경쟁하는 시간이 실제 작업 시간보다 커질 수 있습니다. 또 작업이 다른 작업을 만들어 넣는 재귀적 분할(병렬 정렬, 트리 탐색 등)에서는 방금 만든 하위 작업을 같은 워커가 바로 이어서 처리하는 편이 캐시에 유리한데, 단일 FIFO 큐는 이런 지역성을 살리지 못합니다.
해결: 워커별 deque + 워크 스틸링
워크 스틸링(work stealing)에서는 각 워커가 자기 deque를 갖습니다. 워커는 자기 deque의 한쪽 끝(보통 뒤쪽)에 작업을 넣고 같은 쪽에서 꺼냅니다. 즉 자기 작업은 LIFO로 처리해, 방금 만든 작업의 데이터가 캐시에 남아 있을 때 처리합니다. 자기 deque가 비면 다른 워커 deque의 반대쪽 끝(앞쪽)에서 작업을 훔쳐 옵니다. 앞쪽에 있는 것은 가장 오래된 작업이고, 재귀 분할에서는 대개 아직 쪼개지지 않은 큰 작업이므로 한 번 훔쳐 와서 오래 일할 수 있고, 주인과 도둑이 deque의 서로 다른 끝을 건드려 충돌도 줄어듭니다.
flowchart TB
subgraph 워커1
Q1[deque: 오래된 작업 ... 최근 작업]
W1[뒤에서 pop]
end
subgraph 워커2
Q2[deque: 비어 있음]
W2[다른 deque 앞에서 steal]
end
W2 -.->|steal from front| Q1
워크 스틸링 스레드 풀 (구조를 보여 주는 간소화 버전)
아래 구현은 구조를 이해하기 위한 것으로, 모든 deque를 뮤텍스 하나로 보호합니다. 그래서 락 경합은 단일 큐와 똑같고, 외부에서 넣는 작업은 라운드 로빈으로 각 deque에 나눠 넣습니다. 실제로 경합을 줄이려면 아래의 “락 구조 선택”에서 설명하는 워커별 락이나 lock-free deque가 필요합니다.
#include <deque>
#include <random>
#include <future>
class WorkStealingThreadPool {
public:
explicit WorkStealingThreadPool(size_t numThreads)
: numThreads_(numThreads), queues_(numThreads) {
workers_.reserve(numThreads);
for (size_t i = 0; i < numThreads; ++i) {
workers_.emplace_back([this, i]() { workerLoop(i); });
}
}
~WorkStealingThreadPool() {
{
std::lock_guard<std::mutex> lock(mutex_); // 락 안에서 바꿔야 깨움 신호를 놓치지 않음
done_ = true;
}
cv_.notify_all();
for (auto& w : workers_) {
if (w.joinable()) w.join();
}
}
void submit(Task task) {
{
std::lock_guard<std::mutex> lock(mutex_);
size_t idx = nextSubmitIdx_++ % numThreads_;
queues_[idx].push_back(std::move(task));
}
cv_.notify_one();
}
private:
bool hasWork() const {
for (const auto& q : queues_) if (!q.empty()) return true;
return false;
}
void workerLoop(size_t myIdx) {
std::mt19937 rng{std::random_device{}()};
for (;;) {
Task task;
{
std::unique_lock<std::mutex> lock(mutex_);
cv_.wait(lock, [this] { return done_ || hasWork(); });
if (done_ && !hasWork()) break; // 남은 작업을 모두 처리한 뒤 종료
auto& mine = queues_[myIdx];
if (!mine.empty()) {
task = std::move(mine.back()); // 내 작업: 뒤에서 (LIFO)
mine.pop_back();
} else {
// 무작위 위치부터 다른 deque를 한 바퀴 돌며 훔칠 작업을 찾음
size_t start = rng() % numThreads_;
for (size_t k = 0; k < numThreads_; ++k) {
auto& victim = queues_[(start + k) % numThreads_];
if (!victim.empty()) {
task = std::move(victim.front()); // 남의 작업: 앞에서 (FIFO)
victim.pop_front();
break;
}
}
}
}
if (task) {
try { task(); } catch (...) { /* 로깅 */ }
}
}
}
size_t numThreads_;
bool done_ = false;
size_t nextSubmitIdx_ = 0;
std::vector<std::deque<Task>> queues_;
std::mutex mutex_;
std::condition_variable cv_;
std::vector<std::thread> workers_;
};
훔칠 대상을 무작위 하나만 골라 보고 비어 있으면 그냥 넘어가면, 다른 deque에는 작업이 있는데 이 워커는 대기 조건이 계속 참이라 잠들지 못하고 헛돌게 됩니다. 그래서 무작위 시작점부터 모든 deque를 한 바퀴 확인합니다. done_도 반드시 뮤텍스를 잡고 바꿔야 합니다. 락 없이 바꾸면 워커가 조건을 검사하고 잠들기 직전 사이에 신호가 지나가 버려, 영원히 깨어나지 않는 lost wakeup이 생길 수 있습니다.
락 구조 선택
뮤텍스 하나로 모든 deque를 보호하는 위 방식은 구현이 단순하지만 경합은 줄지 않습니다. 워커별로 뮤텍스를 두면 자기 deque에서 꺼낼 때는 자기 락만, 훔칠 때는 대상 워커의 락만 잡으므로 경합이 분산됩니다. 대신 대기·깨움 처리가 복잡해집니다. 더 나아가 Chase-Lev deque처럼 주인 쪽 연산은 대부분 원자적 읽기·쓰기만으로 처리하고 훔치기만 CAS를 쓰는 lock-free 자료구조를 쓰면 경합을 가장 줄일 수 있지만, 메모리 순서를 정확히 다루기 어렵습니다. 직접 구현하기보다 Intel oneTBB나 Taskflow 같은 검증된 라이브러리를 쓰는 것이 현실적입니다.
단일 큐와 워크 스틸링 선택
| 상황 | 권장 방식 |
|---|---|
| 작업 수가 적고 크기가 비슷함 | 단일 큐 (구현 단순) |
| 작업이 아주 짧고 많아 락 경합이 보임 | 워커별 큐 + 워크 스틸링 |
| 작업이 하위 작업을 재귀적으로 만들어 넣음 | 워크 스틸링 (LIFO 지역성) |
| I/O 대기가 많음 | 단일 큐 + 워커 수를 코어 수보다 많게 |
어느 쪽이 빠른지는 작업 크기와 개수에 따라 달라지므로, 실제 워크로드로 처리량과 락 대기 시간을 측정해 보고 결정해야 합니다.
자주 만나는 문제
notify를 락 안에서 호출해도 되나
notify_one·notify_all을 뮤텍스를 잡은 채로 호출해도 데드락은 생기지 않으며, 표준상 올바른 코드입니다. 다만 깨어난 스레드가 곧바로 락을 얻지 못해 다시 기다리는 비효율이 생길 수 있어, 위 코드처럼 공유 상태를 바꾼 뒤 락을 풀고 notify하는 형태를 많이 씁니다. 반대로 꼭 지켜야 하는 것은 공유 상태(done_, 큐)를 바꾸는 일을 반드시 락 안에서 하는 것입니다. 상태 변경을 락 밖에서 하면 앞에서 본 lost wakeup이 생깁니다.
shutdown 후 push
shutdown() 뒤에 들어온 작업은 실행될 워커가 없으므로, push에서 done_을 확인해 거부해야 합니다. 위 TaskQueue::push는 조용히 무시하지만, 호출자가 알아야 한다면 bool을 반환하거나 예외를 던지게 바꿉니다.
람다의 참조 캡처
// 잘못된 예: path는 paths의 원소에 대한 참조
void processFiles(const std::vector<std::string>& paths) {
for (const auto& path : paths) {
pool.submit([&path]() {
loadFile(path); // 실행될 때 paths가 이미 사라졌을 수 있음
});
}
}
// 올바른 예: 값으로 캡처
for (const auto& path : paths) {
pool.submit([path]() { loadFile(path); });
}
작업은 submit이 반환된 뒤 언제든 실행될 수 있으므로, 람다가 참조하는 객체는 작업이 끝날 때까지 살아 있어야 합니다. 확신이 없으면 값으로 캡처하거나 shared_ptr로 수명을 함께 묶습니다.
워커 안에서 future.get() 호출
// 잘못된 예: 워커 안에서 같은 풀의 결과를 기다림
pool.submit([&pool]() {
auto fut = pool.submitWithResult([]() { return 42; });
int x = fut.get(); // 모든 워커가 이렇게 기다리면 아무도 안쪽 작업을 실행하지 못함
});
워커가 같은 풀에 넣은 작업의 결과를 get()으로 기다리면, 그동안 그 워커는 다른 작업을 처리하지 못합니다. 이런 작업이 워커 수만큼 동시에 실행되면 모든 워커가 기다리기만 하고 안쪽 작업을 실행할 워커가 없어 교착 상태가 됩니다. 결과는 바깥(메인 스레드 등)에서 기다리거나, 다음 단계를 새 작업으로 submit하는 방식으로 연결합니다.
작업이 무한히 쌓이는 경우
submit 속도가 처리 속도보다 계속 빠르면 큐가 끝없이 커져 결국 메모리가 부족해집니다. 큐 길이에 상한을 두고, 가득 차면 생산자를 기다리게 하거나 거부하는 백프레셔(backpressure)를 적용합니다.
// 큐 길이 제한 (간단한 백프레셔)
void push(Task task) {
std::unique_lock<std::mutex> lock(mutex_);
cvNotFull_.wait(lock, [this] { return queue_.size() < maxSize_ || done_; });
if (done_) return;
queue_.push(std::move(task));
lock.unlock();
cv_.notify_one();
}
// waitAndPop에서 pop한 뒤 cvNotFull_.notify_one() 호출
UI 스레드처럼 블로킹되면 안 되는 생산자라면 기다리는 대신 false를 반환해 작업을 버리거나 합치는 정책이 더 적합합니다.
종료 순서
종료는 “종료 플래그를 락 안에서 설정 → notify_all → 모든 워커 join” 순서여야 합니다. join을 먼저 호출하면 워커가 아직 대기 중이라 영원히 반환되지 않고, notify_all을 플래그 설정보다 먼저 하면 깨어난 워커가 플래그를 보지 못하고 다시 잠듭니다. 위 ThreadPool의 소멸자는 queue_.shutdown() 안에서 앞의 두 단계를 처리한 뒤 join합니다.
풀 크기, 우선순위, 종료 정책
풀 크기
CPU 바운드 작업은 코어 수만큼의 워커면 충분합니다. 더 많으면 컨텍스트 스위칭만 늘어납니다. I/O 바운드 작업은 워커가 대기하는 동안 다른 워커가 일할 수 있으므로 코어 수보다 많이 두는 경우가 많고, 적정값은 대기 시간 비율에 따라 달라 측정으로 정합니다. 두 성격이 섞여 있다면 풀을 분리합니다.
// hardware_concurrency()는 알 수 없을 때 0을 반환할 수 있음
unsigned hw = std::max(1u, std::thread::hardware_concurrency());
unsigned cpuPoolSize = std::max(1u, hw - 1); // 메인 스레드 몫 하나를 남김
unsigned ioPoolSize = hw * 2;
hardware_concurrency() - 1을 그대로 쓰면 반환값이 0일 때 부호 없는 정수가 감싸져 수십억 개의 스레드를 만들려 하게 되므로, 먼저 최솟값을 1로 맞춥니다.
CPU 풀과 I/O 풀 분리
CPU 작업과 I/O 대기 작업을 같은 풀에 넣으면, I/O를 기다리는 워커가 자리를 차지해 CPU 작업이 코어를 다 쓰지 못합니다. 두 풀을 분리하고 단계마다 맞는 풀에 넣습니다.
ThreadPool ioPool(ioPoolSize); // 먼저 선언 → 나중에 파괴
ThreadPool cpuPool(cpuPoolSize);
cpuPool.submit([&ioPool, img, path]() {
auto thumb = resize(img, 128, 128); // CPU 작업
ioPool.submit([thumb, path]() { // 저장은 I/O 풀로 넘김
saveToFile(path, thumb);
});
});
여기서는 ioPool이 cpuPool보다 오래 살아 있어야 합니다. 지역 변수는 선언의 역순으로 파괴되므로 ioPool을 먼저 선언했습니다. 순서가 반대면 ioPool이 먼저 파괴되고, cpuPool에 남은 작업이 이미 사라진 ioPool에 접근하게 됩니다.
우선순위 큐
긴급한 작업을 먼저 처리하려면 std::queue 대신 std::priority_queue를 씁니다. 같은 우선순위끼리의 순서는 보장되지 않으므로, 넣은 순서도 유지하려면 증가하는 순번을 함께 저장해 비교합니다.
struct PrioritizedTask {
int priority;
uint64_t seq; // 같은 우선순위면 먼저 들어온 것 먼저
Task task;
bool operator<(const PrioritizedTask& o) const {
if (priority != o.priority) return priority < o.priority; // 큰 숫자가 먼저
return seq > o.seq;
}
};
std::priority_queue<PrioritizedTask> queue_;
priority_queue::top()은 const 참조를 돌려주므로 작업을 이동해 꺼낼 수 없어 복사가 일어납니다. std::function 복사가 부담된다면 std::vector와 std::push_heap·std::pop_heap으로 직접 힙을 관리하면 이동으로 꺼낼 수 있습니다.
배치 완료 대기
여러 작업을 넣고 모두 끝날 때까지 기다리려면 future를 모아 두고 차례로 get()합니다. get()은 작업에서 난 예외도 다시 던져 주므로 오류를 놓치지 않습니다.
std::vector<std::future<void>> futures;
for (const auto& item : items) {
futures.push_back(pool.submitWithResult([&item]() { process(item); }));
}
for (auto& f : futures) f.get(); // 여기서 모두 기다리므로 &item 캡처가 안전
종료 정책: 남은 작업을 처리할지 버릴지
이 글의 TaskQueue는 종료 후에도 큐에 남은 작업을 모두 꺼내 주므로, 풀을 파괴하면 이미 넣은 작업이 다 끝날 때까지 기다립니다. 작업이 오래 걸릴 수 있는데 빨리 종료해야 한다면, shutdown에서 큐를 비우는 shutdownNow()를 따로 두고, 실행 중인 작업에는 취소 플래그로 중단을 알려야 합니다. 어느 쪽이든 정책을 명시하지 않으면 “종료 시 작업이 조용히 사라지는” 버그가 됩니다.
타임아웃
future::wait_for로 결과를 일정 시간만 기다릴 수 있습니다. 하지만 타임아웃이 나도 작업 자체는 워커에서 계속 실행되며, C++ 표준에는 실행 중인 스레드를 강제로 멈추는 방법이 없습니다. 작업 안에서 취소 플래그를 확인하거나, 네트워크 라이브러리의 자체 타임아웃을 함께 설정해야 합니다.
auto fut = pool.submitWithResult([url]() { return fetchUrl(url); });
if (fut.wait_for(std::chrono::seconds(5)) == std::future_status::ready) {
std::string result = fut.get();
} else {
log("Request timed out"); // 작업은 백그라운드에서 계속 실행 중
}
활용 예제
여러 URL 동시 요청
URL마다 본문을 반환하는 작업을 넣고 future를 모은 뒤, 순서대로 get()으로 결과를 받습니다. 풀의 워커가 4개라면 최대 4개 요청이 동시에 진행됩니다.
ThreadPool pool(4);
std::vector<std::future<std::string>> results;
for (const auto& url : urls) {
results.push_back(pool.submitWithResult([url]() {
SimpleHttpClient client;
std::string body;
client.get(host(url), path(url), body);
return body;
}));
}
for (auto& f : results) {
std::cout << f.get().substr(0, 80) << "\n";
}
비동기 로그 쓰기
메인 스레드는 로그 작업을 큐에 넣기만 하고, 전용 워커 하나가 파일에 씁니다. 워커가 하나뿐이라 로그 순서가 유지되고 파일 접근에 별도 락이 필요 없습니다. 종료 시 shutdown() 후 join하면 큐에 남은 로그를 모두 쓴 뒤 워커가 끝납니다.
TaskQueue logQueue;
std::thread logWorker([&]() {
Task t;
while (logQueue.waitAndPop(t)) t();
});
// 메인 스레드
logQueue.push([&logFile, msg = std::string("message")]() {
logFile << msg << '\n';
});
// 종료 시
logQueue.shutdown();
logWorker.join();
다단계 파이프라인
다운로드 → 파싱 → 저장처럼 단계가 이어지는 처리는, 앞 단계 작업이 끝나면서 다음 단계를 새 작업으로 submit하는 방식으로 연결할 수 있습니다. submit은 큐에 넣고 바로 반환하므로, 워커 안에서 get()으로 기다릴 때 생기는 교착 상태가 없습니다.
void processBatch(ThreadPool& pool, const std::vector<std::string>& urls) {
for (const auto& url : urls) {
pool.submit([&pool, url]() {
auto data = fetchUrl(url);
auto parsed = parseJson(data);
pool.submit([parsed = std::move(parsed)]() {
saveToDatabase(parsed);
});
});
}
}
진행률과 취소
긴 배치 작업에서 진행률을 보여 주거나 취소를 지원하려면 원자 변수를 공유합니다. 작업이 비동기로 실행되므로 이 상태는 지역 변수가 아니라 shared_ptr로 작업들과 수명을 함께해야 합니다. 취소된 작업도 완료 수에 포함해야 마지막 작업이 “모두 끝남”을 알릴 수 있습니다.
struct BatchState {
std::atomic<bool> cancelled{false};
std::atomic<int> finished{0};
int total = 0;
};
auto state = std::make_shared<BatchState>();
state->total = 100;
for (int i = 0; i < state->total; ++i) {
pool.submit([state, i]() {
if (!state->cancelled) doWork(i);
if (++state->finished == state->total) {
notifyUiComplete(state->cancelled); // UI 스레드로 이벤트를 넘기는 함수
}
});
}
// 취소 버튼: state->cancelled = true;
같이 보면 좋은 글
- C++ condition_variable로 작업 큐 만들기: 허위 깨움 처리와 스레드 풀
- C++ 고급 멀티스레딩 | 스레드 풀·Work Stealing
- C++ std::thread 입문 | join 누락·디태치 남용 등 자주 하는 실수 3가지와 해결법
- C++ HTTP 클라이언트 직접 만들기
- C++ JSON 파싱
- C++ REST API 클라이언트
- C++ 디자인 패턴 | Observer·Strategy
- C++ RAII: 생성자에서 획득하고 소멸자에서 해제하는 자원 관리 클래스 만들기