第 13 篇讲了 jthread、semaphore、latch/barrier——都是同步工具。生产环境里更常见的问题是:任务来得比处理快,线程该复用还是现开?队列该无限长还是限流?队列满了怎么办?

这一篇用 demo ref/cpp_demo/concurrency/thread_pool/ 讲线程池架构、有界队列、背压策略,以及 future/promise 取异步结果;附带 HTTP REST Server/Client 实战。

这是「现代 C++ 实战」系列的第 14 篇。建议先读 第 13 篇:C++20 同步原语。

一、为什么需要线程池?

每来一个任务就 std::thread(...).detach():

问题 后果
创建开销 线程栈分配、内核调度——毫秒级,短任务不划算
线程数爆炸 1 万并发 → 1 万线程,内存与上下文切换拖垮系统
难以复用 无法统一限流、监控、优雅关闭

线程池:固定 N 个工作线程 + 任务队列,提交方只 submit,worker 循环取任务执行——摊薄创建成本、控制并发上限。

二、基本架构

1
2
3
4
submit() ──→ [ 任务队列 ] ──→ Worker 1
│ Worker 2
│ Worker 3
└─ condition_variable 唤醒

核心组件(demo src/thread_pool.h 中的 threadpool::ThreadPool 类):

组件 作用
workers_ std::vector<std::thread>,长期运行
tasks_ std::queue<std::function<void()>>
max_queue_size_ 队列容量上限,0 表示无限制
queue_mutex_ mutable std::mutex,保护队列
condition_ 队列空时 worker 阻塞,有任务时 notify_one
finished_condition_ 任务完成时 notify_all,供 waitAll() 等待
stop_ std::atomic<bool>,析构时置 true,唤醒 worker 退出
active_tasks_ std::atomic<size_t>,正在执行(已出队)的任务数

Worker 循环(节选自 src/thread_pool.cpp 构造函数里 workers_.emplace_back([this] { ... }) 的 lambda,去掉了注释):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
while (true) {
std::function<void()> task;
{
std::unique_lock<std::mutex> lock(queue_mutex_);
condition_.wait(lock, [this] {
return stop_ || !tasks_.empty();
});
if (stop_ && tasks_.empty()) {
return;
}
task = std::move(tasks_.front());
tasks_.pop();
++active_tasks_;
}
task(); // 在锁外执行,避免阻塞其他 worker
--active_tasks_;
finished_condition_.notify_all(); // 唤醒 waitAll()
}

在锁外执行任务——否则长任务会饿死队列操作和其他 worker。

三、有界队列:防止内存无限增长

无界队列:生产者快于消费者 → 队列里堆积成万上十万 std::function → OOM。

demo 构造函数支持 max_queue_size:

1
threadpool::ThreadPool pool(4, 100);  // 4 线程,队列最多 100 个待处理任务

签名是 explicit ThreadPool(size_t num_threads = std::thread::hardware_concurrency(), size_t max_queue_size = 0);。thread_pool_demo 只传了线程数(无界队列),task_server 才用 -q 把上限传进来。

submit 时(持有 queue_mutex_,先检查 stop_,再检查容量):

1
2
3
if (max_queue_size_ > 0 && tasks_.size() >= max_queue_size_) {
throw QueueFullException(tasks_.size(), max_queue_size_);
}
max_queue_size 行为
0 无限制(仅适合压测或任务量可控)
> 0 满则拒绝新任务(demo 抛 QueueFullException,可用 getQueueSize() / getMaxSize() 取数值)

生产环境常见变体:阻塞 submit(等队列有空位)、丢弃最旧任务、返回错误码——取决于业务能否丢任务。

四、背压(Backpressure)

背压:下游处理不过来时,向上游施加压力,让生产者减速或失败,而不是无限缓冲。

1
2
3
4
客户端 ──快──→ [队列满] ──慢──→ 线程池 worker
↑
拒绝 / 阻塞 / 503
(背压传回上游)
策略 适用
抛异常 / 返回 false API 层返回 503,客户端重试
阻塞 submit 同步调用方自然减速
丢弃 + 计数 可丢的 telemetry、采样
semaphore 限流 配合 第 13 篇 限制「在途任务」总数

demo 的 HTTP Server 在 handlePostTask 里捕获 QueueFullException,返回 503 Service Unavailable,并带上 Retry-After: 1 头和 queue_size / max_queue_size / retry_after 字段,即把背压暴露给 REST 客户端。

五、future / promise / packaged_task

调用 submit 往往要拿返回值或异常——用标准库异步 trio。demo 的 submit 是模板,实现放在 src/thread_pool.h(去掉注释后):

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
template<typename F, typename... Args>
auto ThreadPool::submit(F&& f, Args&&... args)
-> std::future<typename std::invoke_result<F, Args...>::type>
{
using return_type = typename std::invoke_result<F, Args...>::type;

auto task = std::make_shared<std::packaged_task<return_type()>>(
std::bind(std::forward<F>(f), std::forward<Args>(args)...)
);

std::future<return_type> res = task->get_future();
{
std::unique_lock<std::mutex> lock(queue_mutex_);

if (stop_) {
throw std::runtime_error("无法向已停止的线程池提交任务");
}

if (max_queue_size_ > 0 && tasks_.size() >= max_queue_size_) {
throw QueueFullException(tasks_.size(), max_queue_size_);
}

tasks_.emplace([task]() { (*task)(); });
}
condition_.notify_one();
return res;
}
类型 角色
packaged_task 包装可调用对象,执行结果写入关联的 future
future 调用方 get() 阻塞取结果(或异常)
promise 手动设值(更底层;线程池一般用 packaged_task)

用法(src/main.cpp 里 computeTask(int task_id, int duration_ms) 返回 int):

1
2
3
4
std::vector<std::future<int>> results;
results.push_back(pool.submit(computeTask, i, duration));
// ...
task_results.push_back(results[i].get()); // 阻塞直到第 i 个任务完成

shared_ptr 包装 packaged_task:因为 std::function 要求可拷贝,而 packaged_task 只能移动。

六、HTTP REST:线程池的实际应用

demo 除了本地压测的 thread_pool_demo(src/main.cpp),还提供 Server 与 Client 两个可执行文件 task_server / task_client(src/server.cpp / src/client.cpp,基于 cpp-httplib v0.58.0 + nlohmann/json):

端点 行为
POST /task 请求体 {"duration_ms": N}(1–60000),提交到线程池,返回 201 与 task_id;队列满返回 503
GET /task/:id 查询任务状态(pending / running / completed / failed)
GET /status thread_count、pending_tasks、max_queue_size、queue_usage、queue_full、total_tasks_received 等
GET /health 健康检查,返回 {"status":"ok"}

流程:HTTP 请求线程只做解析 + submit,重活在线程池 worker 里跑——避免「一请求一线程」把连接数撑爆。队列满时 REST 层返回 503,客户端可退避重试。

1
2
3
4
5
6
7
8
9
10
11
12
13
cd ref/cpp_demo/concurrency/thread_pool
./build.sh # 编译出三个可执行文件
./build.sh --run thread_pool_demo # 本地 submit 压测
./build/thread_pool_demo -t 4 -T 20 # 4 线程、20 个任务

# 终端 1:启动 REST 服务(8 线程、队列上限 100)
./build/task_server -p 8080 -t 8 -q 100

# 终端 2:客户端
./build/task_client -p 8080 status # 查看服务器状态
./build/task_client -p 8080 submit -d 500 # 提交一个 500ms 的任务
./build/task_client -p 8080 query -i <task_id> # 查询任务
./build/task_client -p 8080 batch -c 10 --min 100 --max 500 # 批量提交并等待

不要对这个 demo 直接 ./build.sh --run:它会依次运行所有可执行文件,而 task_server 会一直阻塞在监听上。

七、性能调优:线程数设多少?

场景 经验
CPU 密集 ≈ std::thread::hardware_concurrency()(物理/逻辑核数)
I/O 密集 可大于核数(线程阻塞等 I/O 时不占 CPU)
混合 压测:从小往大调,看吞吐与延迟拐点

demo 构造函数默认参数就是 hardware_concurrency();显式传 0(命令行 -t 0,也是默认值)时同样取 hardware_concurrency(),若它返回 0 则回退 4。

别迷信公式:线程过多 → 上下文切换、缓存失效;过少 → CPU 闲置。用指标(队列长度、P99 延迟、CPU 利用率)迭代。

八、优雅关闭

析构顺序(demo):

  1. 持 queue_mutex_ 把 stop_ 置为 true(此后 submit 抛 runtime_error)
  2. condition_.notify_all() 唤醒所有 worker
  3. worker 先把队列里剩余任务执行完,发现 stop_ && tasks_.empty() 才退出
  4. 析构线程对每个 joinable() 的 worker 调用 join()

waitAll():在 finished_condition_ 上等「队列空且 active_tasks_ == 0」——适合「提交完一批不关心返回值的任务再退出」。thread_pool_demo 走的是另一条路:保存每个 future,按顺序 get() 收集结果。

九、小结

概念 要点
线程池 固定 worker + 任务队列,复用线程
有界队列 限制待处理任务数,防 OOM
背压 队列满时拒绝/阻塞/503,压力回传上游
future submit 返回异步结果
REST 请求线程轻量,计算进池

现代 C++ 实战系列第 14 篇完。下一篇 现代设计模式——Strategy、Observer、Builder 的 C++17 写法。

系列导航

篇号 标题 状态
13 C++20 同步原语 ✅
14 线程池与背压控制(本篇) ✅
15 现代设计模式 ✅