C++ 线程池:任务提交、异常传递与安全关闭
前言
线程池通过复用一组工作线程,执行不断提交的任务。实现一个基础线程池,需要协调任务队列、线程等待、返回值传递和对象生命周期。
本文使用 C++17 实现固定线程数量的线程池,支持通过 std::future 获取结果或任务异常,并处理线程创建中途失败的清理问题。它采用“停止接收新任务,执行完已接收任务,再退出”的关闭策略。
示例用于说明核心机制,不包含队列容量限制、任务取消、优先级或工作窃取。是否适合实际服务,需要结合负载与关闭要求评估。
核心设计思路
线程池由以下部分组成:
- 任务队列:保存待执行的
std::function<void()>。每个任务先封装自己的参数与返回值处理,再转换成统一的无参数调用入口。 - 工作线程:等待任务,在锁内取出一个任务,然后在锁外执行,避免把任务执行串行化。
- 同步机制:用同一把互斥锁保护队列和关闭标志,使用带谓词的条件变量等待处理通知和虚假唤醒。
- 结果通道:用
std::packaged_task将执行结果或抛出的异常保存到共享状态,由对应的std::future读取。
这里区分两个动作:
close():停止接收新任务并唤醒工作线程,不等待任务完成;可以重复调用。- 析构:调用
close(),再join()所有工作线程,等待已经接收的任务执行结束。
对象正常存活期间,多个线程可以并发调用 push() 和 close()。析构必须由外部协调:析构开始前,外部提交者和其他成员函数调用应已结束;任务也不能在析构期间重新调用该线程池。不得从线程池自身的工作线程销毁线程池。
完整代码实现
下面的实现按值保存可调用对象和参数:传入左值时通常复制,传入右值时可以移动;任务执行时,将保存的对象和参数按右值转交给 std::invoke,每个任务只执行一次。
这支持移动捕获的 Lambda、按值接收 std::unique_ptr 的函数,以及成员函数指针。需要传递外部对象的引用时,显式使用 std::ref 或 std::cref,并由调用者保证生命周期和同步。
#include <condition_variable>
#include <cstddef>
#include <functional>
#include <future>
#include <memory>
#include <mutex>
#include <queue>
#include <stdexcept>
#include <thread>
#include <tuple>
#include <type_traits>
#include <utility>
#include <vector>
class ThreadPool {
private:
std::mutex mutex_;
std::condition_variable condition_;
std::queue<std::function<void()>> tasks_;
bool closed_ = false; // 所有访问均由 mutex_ 保护
std::vector<std::thread> workers_;
void worker_loop() {
for (;;) {
std::function<void()> task;
{
std::unique_lock<std::mutex> lock(mutex_);
condition_.wait(lock, [this] {
return closed_ || !tasks_.empty();
});
if (closed_ && tasks_.empty()) {
return;
}
task = std::move(tasks_.front());
tasks_.pop();
}
// 用户任务的异常由 packaged_task 保存到 future。
// 队列只接收下面 push() 创建的有效、仅执行一次的包装任务。
task();
}
}
// 仅由构造失败清理路径或析构函数调用,不允许并发调用。
void join_workers() {
for (auto& worker : workers_) {
if (worker.joinable()) {
worker.join();
}
}
}
public:
explicit ThreadPool(std::size_t thread_count) {
if (thread_count == 0) {
throw std::invalid_argument("thread_count must be greater than zero");
}
// 若 reserve 失败,此时还没有启动任何工作线程。
workers_.reserve(thread_count);
try {
for (std::size_t i = 0; i < thread_count; ++i) {
workers_.emplace_back([this] {
worker_loop();
});
}
} catch (...) {
// 构造失败不会调用 ThreadPool 的析构函数,必须在这里收尾。
close();
join_workers();
throw;
}
}
template <class F, class... Args>
auto push(F&& func, Args&&... args)
-> std::future<std::invoke_result_t<
std::decay_t<F>, std::decay_t<Args>...>> {
using R = std::invoke_result_t<
std::decay_t<F>, std::decay_t<Args>...>;
// 本接口支持值、void 和左值引用结果,不支持 future<T&&>。
static_assert(!std::is_rvalue_reference_v<R>,
"tasks returning rvalue references are not supported");
auto invocation =
[fn = std::decay_t<F>(std::forward<F>(func)),
saved = std::tuple<std::decay_t<Args>...>(
std::forward<Args>(args)...)]() mutable -> R {
return std::apply(
[&fn](auto&&... values) -> R {
return std::invoke(
std::move(fn),
std::forward<decltype(values)>(values)...);
},
std::move(saved));
};
auto task = std::make_shared<std::packaged_task<R()>>(
std::move(invocation));
auto result = task->get_future();
{
std::lock_guard<std::mutex> lock(mutex_);
if (closed_) {
throw std::runtime_error("cannot push to a closed thread pool");
}
tasks_.emplace( {
(*task)();
});
}
condition_.notify_one();
return result;
}
// 停止接收新任务;已经接收的任务继续执行。
void close() {
{
std::lock_guard<std::mutex> lock(mutex_);
closed_ = true;
}
condition_.notify_all();
}
// 使用前提:不能由本池工作线程执行,也不能与外部成员调用并发。
~ThreadPool() {
close();
join_workers();
}
ThreadPool(const ThreadPool&) = delete;
ThreadPool& operator=(const ThreadPool&) = delete;
ThreadPool(ThreadPool&&) = delete;
ThreadPool& operator=(ThreadPool&&) = delete;
};
push() 与 close() 通过同一把锁决定任务是否被接收:如果入队先完成,任务会继续执行;如果关闭标志先被设置,提交会抛异常。
任务包装在锁外完成,以减少持锁时间。因此,即使最终因线程池已关闭而拒绝提交,传入的右值也可能已经被移动。这里保证的是任务没有入队,不承诺恢复调用者参数的原始状态。
使用示例与基础检查
将下面的代码接在类定义之后,可以检查普通返回值、任务异常、移动参数、成员函数调用、引用参数、关闭后拒绝提交以及析构时排空任务。
#include <atomic>
#include <cassert>
#include <iostream>
struct Calculator {
int add(int a, int b) const {
return a + b;
}
};
int main() {
// 零线程会导致任务无人执行,因此直接拒绝。
bool zero_rejected = false;
try {
ThreadPool invalid(0);
} catch (const std::invalid_argument&) {
zero_rejected = true;
}
assert(zero_rejected);
ThreadPool pool(4);
auto sum = pool.push([](int a, int b) {
return a + b;
}, 10, 20);
auto failure = pool.push([]() -> int {
throw std::runtime_error("task failed");
});
auto moved_argument = pool.push([](std::unique_ptr<int> value) {
return *value;
}, std::make_unique<int>(42));
auto moved_callable = pool.push(
[value = std::make_unique<int>(7)] {
return *value;
});
// 按值保存 Calculator,不依赖外部对象的生命周期。
auto member_result = pool.push(&Calculator::add, Calculator{}, 2, 3);
int counter = 0;
auto by_reference = pool.push([](int& value) {
++value;
}, std::ref(counter));
const int sum_value = sum.get();
const int moved_value = moved_argument.get();
const int callable_value = moved_callable.get();
const int member_value = member_result.get();
by_reference.get(); // 等任务完成后,主线程才读取 counter。
assert(sum_value == 30);
assert(moved_value == 42);
assert(callable_value == 7);
assert(member_value == 5);
assert(counter == 1);
bool task_exception_received = false;
try {
(void)failure.get();
} catch (const std::runtime_error&) {
task_exception_received = true;
}
assert(task_exception_received);
// 普通任务异常不会使工作线程退出。
auto after_failure = pool.push([] { return 99; });
const int next_value = after_failure.get();
assert(next_value == 99);
pool.close();
pool.close(); // 重复关闭允许发生。
bool submission_rejected = false;
try {
pool.push([] {});
} catch (const std::runtime_error&) {
submission_rejected = true;
}
assert(submission_rejected);
// 不依赖逐个 future.get(),由析构等待所有已接收任务执行完。
std::atomic<int> completed{0};
{
ThreadPool draining_pool(2);
for (int i = 0; i < 100; ++i) {
draining_pool.push([&completed] {
completed.fetch_add(1, std::memory_order_relaxed);
});
}
}
assert(completed.load(std::memory_order_relaxed) == 100);
std::cout << "Basic checks passed\n";
}
使用 GCC 或 Clang 时,可以保存为 thread_pool.cpp,以 C++17 模式编译,例如:
g++ -std=c++17 -O2 -Wall -Wextra -Wpedantic -pthread thread_pool.cpp -o thread_pool
这些是基础行为检查,不是并发正确性的完整证明。进一步验证应覆盖多提交者并发、提交与关闭竞争,以及线程创建中途失败的故障注入;也可结合平台支持的竞态检测工具检查数据竞争。
关键问题解析
1. 条件变量为什么要配合同一把互斥锁?
工作线程等待的不是“通知次数”,而是 closed_ || !tasks_.empty() 这个共享状态。带谓词的等待等价于循环检查条件,在条件不满足时调用 wait(),因此也能处理虚假唤醒。条件变量标准说明
wait() 的“释放锁并进入等待”是原子步骤。原文中关于丢失唤醒的风险方向正确,但需要把竞争窗口描述清楚:问题可以发生在谓词已经检查为假、尚未调用底层等待操作之间,而不是 wait() 原子解锁步骤内部存在漏洞。
如果关闭线程不遵循同一把锁的协议,就可能出现:
- 工作线程持锁检查关闭标志,读到
false,队列也为空。 - 关闭线程绕过这把锁,修改原子关闭标志并发送通知。
- 工作线程才进入等待;此前的通知不会作为一个待消费事件保留。
若之后没有其他通知或虚假唤醒,该线程可能一直等待。仅把标志换成 std::atomic<bool>,只能解决标志本身的数据竞争,不能自动修复这个等待协议。
本实现的检查和修改都受同一把锁保护:关闭线程要么先修改状态,使工作线程不再等待;要么等工作线程原子地解锁并进入等待后才能修改状态,再通知它。
notify_one() 和 notify_all() 不必在持锁状态下调用。本例在解锁后通知;对象生命周期则由前述使用约束保证。也存在其他正确的等待设计,不能将本实现的协议推广为“所有条件变量代码都只能这样写”。
2. 为什么将 packaged_task 放进 shared_ptr?
std::packaged_task 仅支持移动,而 C++17 的 std::function 要求保存的目标可复制构造。因此,不能把一个直接持有 packaged_task 的仅移动 Lambda 放入本例的 std::function<void()> 队列。std::function 标准说明
使用 shared_ptr 后,外层 Lambda 可以复制;真正的任务状态仍由同一个 packaged_task 管理。这里的共享所有权用于满足包装器的类型要求,不表示允许多个线程重复执行同一个任务。
这不是唯一方案:可以使用自定义的仅移动任务包装器,或改用 std::queue<std::packaged_task<void()>>,在这个统一签名的外层任务中持有不同返回类型的内层任务。选择不同的队列元素类型,就可以不依赖 shared_ptr 解决可复制性问题。
同样,返回 future 也不是必须使用 packaged_task;还可以使用 promise 等机制。packaged_task 的优势是已经封装了调用、结果存储和异常传递。
3. 任务异常究竟由谁处理?
std::packaged_task::operator() 调用用户任务。如果用户任务正常返回,它保存结果;如果任务抛出异常,它保存异常,并让关联的共享状态就绪。之后调用 future.get(),异常会在等待结果的线程中重新抛出。packaged_task 标准说明
因此,原文工作线程外层的 try-catch 通常捕获不到用户任务抛出的异常,也不会自动把这类异常写入日志。对于本例中有效且只调用一次的 packaged_task,用户任务异常已经由结果通道处理。
修订实现没有额外吞掉工作线程内部错误。若以后允许直接提交未经 packaged_task 包装的任务,应另外定义异常处理策略,因为逃出线程入口函数的异常会触发 std::terminate()。std::thread 构造标准说明 packaged_task 被重复执行或没有有效共享状态等错误,也不属于正常任务异常通道的使用范围。
还应区分 C++ 异常与程序故障:vector::at() 越界会抛异常,而数组越界、无效指针访问等未定义行为不能靠普通 catch (...) 保证恢复;任务直接调用 std::terminate() 也无法这样拦截。捕获异常不能作为“服务高可用”的保证。
如果调用者丢弃 future,任务仍会执行,但保存的异常通常不会被观察或自动记录。需要日志时,应明确由谁读取结果,或在任务内部记录后重新抛出。
4. 为什么线程池构造也要处理异常?
创建线程可能因资源不足而抛出 std::system_error。std::thread 构造标准说明 如果前几个线程已经启动,之后的线程创建失败,构造函数不会完成,线程池自身的析构函数也不会执行。
此时,成员会被销毁。若 vector 中仍有可连接的 std::thread,它们的析构会调用 std::terminate(),即使底层线程恰好已经执行结束,也不能省略 join()。std::thread 析构标准说明
所以代码需要在构造函数的异常路径中设置关闭标志、唤醒并连接已经启动的线程,然后重新抛出原异常。提前 reserve() 可以减少创建过程中的容器分配,但不能代替这个清理过程。
这里处理的是正常同步条件下的构造回滚和资源清理,不承诺恢复互斥锁、线程连接等底层机制自身的异常故障。析构函数也不提供抛异常后的恢复接口;正常使用必须满足线程连接与生命周期前提。
5. 为什么不继续使用 decltype(func(args...)) 与 std::bind?
原文中的 func 和 args 都是有名字的变量,表达式中按左值参与调用。decltype(func(args...)) 因而不能准确描述异步任务中保存副本后的调用方式,也不能直接覆盖成员函数指针调用。
std::bind 对普通已绑定参数通常以左值形式传给目标函数。std::bind 标准说明 因此,即使在绑定阶段移动保存了 unique_ptr,执行阶段也不意味着会自动再次移动它,按值接收移动专用参数的函数便可能无法调用。
修订代码用 decay_t 明确保存在任务中的类型,用 invoke_result_t 推导对应右值调用的结果,再通过 tuple、apply 和 invoke 实现一致的调用语义。std::invoke 标准说明 这也意味着:只允许以左值调用的函数对象需要通过 std::ref 显式提交,不能继续宣称支持“任意参数、任意可调用对象”。
类型推导不能替代生命周期管理。指针、string_view、引用捕获和 std::ref 都不会自动延长所指对象的寿命;如果任务返回左值引用,该引用也必须在调用者使用期间有效,尤其不能返回即将销毁的任务私有副本中的引用。
6. 析构等待不等于能够保证及时退出
本例选择排空任务后退出,因此有以下限制:
- 任务永久阻塞或不返回时,析构也可能一直等待。
close()不取消正在运行或已经排队的任务。 - 如果一个工作线程提交子任务后同步等待结果,而所有工作线程都以同样方式等待,子任务可能无人执行。例如,单线程池中的任务提交另一个任务后立即
get(),会形成死锁。 - 在工作线程中销毁所属线程池会涉及连接自身线程,不能通过简单地跳过该线程的
join()来修复,因为工作线程仍可能访问正在销毁的成员。 - 队列没有容量上限。持续提交快于执行时,内存和排队延迟可能不断增长;需要背压、拒绝策略或容量限制的场景,应扩展接口设计。
- 队列按入队顺序取出任务,但多个工作线程的实际开始和完成顺序不保证一致。
需要显式关闭时,可以先在对象仍然存活期间调用 close(),协调所有提交者停止并结束成员函数调用,再由外部拥有者销毁线程池。互斥锁能够保护内部状态,但不能让访问已销毁对象变得安全。
参考资料
- C++ 工作草案:条件变量的等待、通知与谓词等待语义。
- C++ 工作草案:std::function 对目标对象的要求。
- C++ 工作草案:packaged_task 的结果与异常存储。
- C++ 工作草案:线程创建失败及线程入口异常。
- C++ 工作草案:std::thread 的析构行为。
- C++ 工作草案:std::bind 的参数传递规则。
- C++ 工作草案:std::invoke。

Comments NOTHING