跳转到内容
新建笔记

C++ 线程池:任务所有权、有界队列与排空关闭

线程池把一组工作线程保留下来,让它们反复从任务队列取工作。它可以减少短任务反复创建线程的开销,也能限制同时执行的任务数;是否更快仍取决于任务粒度、阻塞比例和机器资源。工作线程不是越多越好,线程池也不保证实时响应。

本页采用 C++17,先定义关闭和所有权规则,再给出一个有界队列、返回结果、排空退出的完整例子。先读线程生命周期、互斥量和条件变量。

1. 四个组成部分和两个边界

跳转到“1. 四个组成部分和两个边界”
组成部分责任
线程池对象创建、关闭并回收工作线程,拥有队列与同步状态。
工作线程等待任务,取出后释放队列锁,再执行任务。
任务队列拥有尚未执行的可调用对象;排队任务数受上限约束。
submit接收任务并返回 std::future<R>,用它取得结果或异常。

原来常见的 append(Task*) 接口,如果只把裸指针排队、调用者紧接着 delete task,后台线程就可能访问已释放对象。必须明确由谁维持任务和任务输入的寿命。下面把任务封装成 packaged_task,由队列持有;任务捕获的引用或裸指针仍由调用者保证有效。

第二个边界是关闭:本例一旦关闭就拒绝新任务,但会完成已经接受的任务。关闭不是强杀线程,也不是丢弃队列。

接收任务 ── close() ──► 拒绝新任务,继续执行队列 ──► 队列为空,工作线程退出
│
shutdown()/析构等待全部线程结束

2. 完整实现:有界队列和排空退出

跳转到“2. 完整实现:有界队列和排空退出”

将下面整个代码块保存为一个 .cpp 文件。maxQueued 限制等待中的任务数,不包括已经执行的任务;满队列立即抛异常,调用方可以延后重试、丢弃不重要任务或向上游施加背压。本例不让 submit 阻塞,以免池内任务递归提交时又增加一条隐蔽等待链。

#include <cassert>
#include <condition_variable>
#include <cstddef>
#include <functional>
#include <future>
#include <memory>
#include <mutex>
#include <queue>
#include <stdexcept>
#include <thread>
#include <type_traits>
#include <utility>
#include <vector>
class ThreadPool {
public:
explicit ThreadPool(std::size_t workers, std::size_t maxQueued = 256)
: maxQueued_(maxQueued)
{
// 64只是本示例的资源预算,不是C++规定或通用最佳线程数。
if (workers == 0 || workers > 64 || maxQueued == 0)
throw std::invalid_argument("invalid pool limits");
workers_.reserve(workers);
try {
for (std::size_t i = 0; i < workers; ++i)
workers_.emplace_back([this] { run(); });
} catch (...) {
// 某个线程创建失败时,也必须回收之前成功创建的线程。
close();
joinWorkers();
throw;
}
}
ThreadPool(const ThreadPool&) = delete;
ThreadPool& operator=(const ThreadPool&) = delete;
~ThreadPool()
{
shutdown();
}
template<class F>
auto submit(F&& function)
-> std::future<std::invoke_result_t<std::decay_t<F>&>>
{
using Result = std::invoke_result_t<std::decay_t<F>&>;
auto task = std::make_shared<std::packaged_task<Result()>>(
std::forward<F>(function));
auto result = task->get_future();
{
std::lock_guard<std::mutex> lock(mutex_);
if (closed_)
throw std::runtime_error("pool is closed");
if (tasks_.size() >= maxQueued_)
throw std::runtime_error("task queue is full");
tasks_.emplace([task] { (*task)(); });
}
changed_.notify_one();
return result;
}
void close()
{
{
std::lock_guard<std::mutex> lock(mutex_);
closed_ = true;
}
changed_.notify_all();
}
// 由池外拥有者调用;不得与另一shutdown/析构并发。
void shutdown()
{
close();
joinWorkers();
}
private:
void run()
{
for (;;) {
std::function<void()> task;
{
std::unique_lock<std::mutex> lock(mutex_);
changed_.wait(lock, [this] {
return closed_ || !tasks_.empty();
});
if (tasks_.empty())
return; // 谓词成立且队列空,说明已经关闭。
task = std::move(tasks_.front());
tasks_.pop();
}
task(); // packaged_task将用户函数抛出的异常存入future。
}
}
void joinWorkers()
{
for (auto& worker : workers_)
if (worker.joinable())
worker.join();
}
const std::size_t maxQueued_;
std::mutex mutex_;
std::condition_variable changed_;
std::queue<std::function<void()>> tasks_;
bool closed_ = false; // 所有访问都在mutex_保护下。
std::vector<std::thread> workers_;
};
int main()
{
ThreadPool pool(2);
auto number = std::make_unique<int>(21);
auto answer = pool.submit([value = std::move(number)] {
return *value * 2;
});
auto failure = pool.submit([]() -> int {
throw std::runtime_error("example task failed");
});
auto stillWorks = pool.submit([] { return 7; });
pool.close();
bool rejected = false;
try {
(void)pool.submit([] {});
} catch (const std::runtime_error&) {
rejected = true;
}
pool.shutdown();
assert(!number && answer.get() == 42 && stillWorks.get() == 7);
assert(rejected);
bool receivedException = false;
try {
(void)failure.get();
} catch (const std::runtime_error&) {
receivedException = true;
}
assert(receivedException);
}

支持 POSIX 线程的 GCC/Clang 工具链通常使用 -std=c++17 -pthread;具体线程库链接选项按工具链配置。例子没有依赖打印顺序,断言验证返回值、移动捕获、异常传递和停止后的拒绝提交。

3. 为什么这些位置必须这样写

跳转到“3. 为什么这些位置必须这样写”

mutex_ 同时保护队列和 closed_,工作线程不能在锁外先做 while (!closed_)。否则普通布尔变量可能发生数据竞争,而且刚关闭时会直接退出、丢掉已经接受的任务。

条件变量等待的是“已关闭或有任务”这一状态。通知只促使等待者重新检查;先提交、后开始等待也没有问题,因为队列状态保留在锁保护下。关闭时用 notify_all,让所有空闲线程都能退出。任务完成顺序不保证与提交顺序相同:队列先进先出取任务,多个线程仍可先后完成不同任务。

取出任务后释放队列锁,用户函数才开始执行。这样不同工作线程可同时处理任务,也避免用户函数在持有内部锁时再次操作池。packaged_task保存函数返回值或异常,future::get()取结果、必要时等待,并重新抛出已保存的异常。一个普通 future的结果只取一次。

队列中的 std::function要求可复制目标,所以这里捕获指向 packaged_task 的 shared_ptr。这不要求用户函数自身可复制:例子已通过移动捕获持有 unique_ptr。共享的是任务包装的所有权,不代表任务会被执行多次。

原来的 Task::process() 形式也可适配为 pool.submit([task = std::make_unique<Task>()] { task->process(); }),其中 Task必须先定义;不要再在提交后手动释放任务。

4. 生命周期和仍然存在的限制

跳转到“4. 生命周期和仍然存在的限制”
  • 多个调用者可在池对象有效时并发 submit/close。shutdown及析构由池外的单一拥有者管理;进入析构前,要确保外部调用方不再访问这个对象。
  • 不得从本池工作线程调用 shutdown或析构池,否则会试图等待自己。任务可以请求 close,实际等待由外部拥有者完成。
  • 已接受的任务必须能结束。永久阻塞的I/O、无限循环、互相等待的任务都会使排空一直等待;标准C++不提供通用的安全强杀线程。
  • 不要让所有工作线程都提交新任务后,阻塞等待同一池中的新任务结果。这会耗尽可执行这些依赖任务的线程,即使队列未满也会死锁。
  • 捕获 this、引用和裸指针不延长对象寿命。池拥有任务包装,不自动拥有它引用的全部对象。
  • 这里只限制线程数和排队任务数量,不估算每个任务捕获的内存、不支持优先级、超时取消或公平性保证。生产服务还需按负载设计这些策略。

这是一份用于解释所有权、同步和关闭协议的教学实现。编译和有限场景断言能发现实际错误,但不能代替对全部线程交错的证明或真实负载测试。