AI 摘要

本文提出一种支持返回值传递且具备异常安全性的现代C++线程池实现,深入任务队列、条件变量协同机制,并剖析丢失唤醒、可拷贝包装等关键并发陷阱,为面试与工程实践提供高可用范式。

C++ 线程池:任务提交、异常传递与安全关闭

前言

线程池通过复用一组工作线程,执行不断提交的任务。实现一个基础线程池,需要协调任务队列、线程等待、返回值传递和对象生命周期。

本文使用 C++17 实现固定线程数量的线程池,支持通过 std::future 获取结果或任务异常,并处理线程创建中途失败的清理问题。它采用“停止接收新任务,执行完已接收任务,再退出”的关闭策略。

示例用于说明核心机制,不包含队列容量限制、任务取消、优先级或工作窃取。是否适合实际服务,需要结合负载与关闭要求评估。


核心设计思路

线程池由以下部分组成:

  1. 任务队列:保存待执行的 std::function<void()>。每个任务先封装自己的参数与返回值处理,再转换成统一的无参数调用入口。
  2. 工作线程:等待任务,在锁内取出一个任务,然后在锁外执行,避免把任务执行串行化。
  3. 同步机制:用同一把互斥锁保护队列和关闭标志,使用带谓词的条件变量等待处理通知和虚假唤醒。
  4. 结果通道:用 std::packaged_task 将执行结果或抛出的异常保存到共享状态,由对应的 std::future 读取。

这里区分两个动作:

  • close():停止接收新任务并唤醒工作线程,不等待任务完成;可以重复调用。
  • 析构:调用 close(),再 join() 所有工作线程,等待已经接收的任务执行结束。

对象正常存活期间,多个线程可以并发调用 push()close()析构必须由外部协调:析构开始前,外部提交者和其他成员函数调用应已结束;任务也不能在析构期间重新调用该线程池。不得从线程池自身的工作线程销毁线程池。


完整代码实现

下面的实现按值保存可调用对象和参数:传入左值时通常复制,传入右值时可以移动;任务执行时,将保存的对象和参数按右值转交给 std::invoke,每个任务只执行一次。

这支持移动捕获的 Lambda、按值接收 std::unique_ptr 的函数,以及成员函数指针。需要传递外部对象的引用时,显式使用 std::refstd::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() 原子解锁步骤内部存在漏洞。

如果关闭线程不遵循同一把锁的协议,就可能出现:

  1. 工作线程持锁检查关闭标志,读到 false,队列也为空。
  2. 关闭线程绕过这把锁,修改原子关闭标志并发送通知。
  3. 工作线程才进入等待;此前的通知不会作为一个待消费事件保留。

若之后没有其他通知或虚假唤醒,该线程可能一直等待。仅把标志换成 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_errorstd::thread 构造标准说明 如果前几个线程已经启动,之后的线程创建失败,构造函数不会完成,线程池自身的析构函数也不会执行。

此时,成员会被销毁。若 vector 中仍有可连接的 std::thread,它们的析构会调用 std::terminate(),即使底层线程恰好已经执行结束,也不能省略 join()std::thread 析构标准说明

所以代码需要在构造函数的异常路径中设置关闭标志、唤醒并连接已经启动的线程,然后重新抛出原异常。提前 reserve() 可以减少创建过程中的容器分配,但不能代替这个清理过程。

这里处理的是正常同步条件下的构造回滚和资源清理,不承诺恢复互斥锁、线程连接等底层机制自身的异常故障。析构函数也不提供抛异常后的恢复接口;正常使用必须满足线程连接与生命周期前提。

5. 为什么不继续使用 decltype(func(args...))std::bind

原文中的 funcargs 都是有名字的变量,表达式中按左值参与调用。decltype(func(args...)) 因而不能准确描述异步任务中保存副本后的调用方式,也不能直接覆盖成员函数指针调用。

std::bind 对普通已绑定参数通常以左值形式传给目标函数。std::bind 标准说明 因此,即使在绑定阶段移动保存了 unique_ptr,执行阶段也不意味着会自动再次移动它,按值接收移动专用参数的函数便可能无法调用。

修订代码用 decay_t 明确保存在任务中的类型,用 invoke_result_t 推导对应右值调用的结果,再通过 tupleapplyinvoke 实现一致的调用语义。std::invoke 标准说明 这也意味着:只允许以左值调用的函数对象需要通过 std::ref 显式提交,不能继续宣称支持“任意参数、任意可调用对象”。

类型推导不能替代生命周期管理。指针、string_view、引用捕获和 std::ref 都不会自动延长所指对象的寿命;如果任务返回左值引用,该引用也必须在调用者使用期间有效,尤其不能返回即将销毁的任务私有副本中的引用。

6. 析构等待不等于能够保证及时退出

本例选择排空任务后退出,因此有以下限制:

  • 任务永久阻塞或不返回时,析构也可能一直等待。close() 不取消正在运行或已经排队的任务。
  • 如果一个工作线程提交子任务后同步等待结果,而所有工作线程都以同样方式等待,子任务可能无人执行。例如,单线程池中的任务提交另一个任务后立即 get(),会形成死锁。
  • 在工作线程中销毁所属线程池会涉及连接自身线程,不能通过简单地跳过该线程的 join() 来修复,因为工作线程仍可能访问正在销毁的成员。
  • 队列没有容量上限。持续提交快于执行时,内存和排队延迟可能不断增长;需要背压、拒绝策略或容量限制的场景,应扩展接口设计。
  • 队列按入队顺序取出任务,但多个工作线程的实际开始和完成顺序不保证一致。

需要显式关闭时,可以先在对象仍然存活期间调用 close(),协调所有提交者停止并结束成员函数调用,再由外部拥有者销毁线程池。互斥锁能够保护内部状态,但不能让访问已销毁对象变得安全。


参考资料

  1. C++ 工作草案:条件变量的等待、通知与谓词等待语义
  2. C++ 工作草案:std::function 对目标对象的要求
  3. C++ 工作草案:packaged_task 的结果与异常存储
  4. C++ 工作草案:线程创建失败及线程入口异常
  5. C++ 工作草案:std::thread 的析构行为
  6. C++ 工作草案:std::bind 的参数传递规则
  7. C++ 工作草案:std::invoke