线程池线程过载时join调用挂死问题排查求助
线程池join挂死问题排查与修复
问题描述
自行实现的线程池在大量循环调用场景下,出现waitFinished函数中join挂死的情况。已确认所有任务均执行完成,但部分线程无法正常退出,导致join卡住。
线程池实现代码
#include <queue> #include <mutex> #include <condition_variable> #include <functional> #include <atomic> #include <vector> #include <thread> #include <iostream> class ThreadPool { public: ThreadPool() { m_shutdown.store(false, std::memory_order_relaxed); createThreads(1); } ThreadPool(std::size_t numThreads) { m_shutdown.store(false, std::memory_order_relaxed); createThreads(numThreads); } void add_job(std::function<void()> new_job) { { std::scoped_lock<std::mutex> lock(m_jobMutex); m_jobQueue.push(new_job); } m_notifier.notify_one(); } void waitFinished() { { std::unique_lock<std::mutex> lock(m_jobMutex); m_finished.wait(lock, [this] {return m_jobQueue.empty(); }); //&& busy == 0 } m_shutdown.store(true, std::memory_order_relaxed); m_notifier.notify_all(); for (std::thread& th : m_threads) { th.join(); } m_threads.clear(); } private: using Job = std::function<void()>; std::vector<std::thread> m_threads; std::queue<Job> m_jobQueue; std::condition_variable m_notifier; std::condition_variable m_finished; std::mutex m_jobMutex; std::atomic<bool> m_shutdown; void createThreads(std::size_t numThreads) { // Settup threads m_threads.reserve(numThreads); for (int i = 0; i != numThreads; ++i) { m_threads.emplace_back(std::thread([this]() { // Infinite loop to consume tasks from queue and execute while (true) { Job job; { std::unique_lock<std::mutex> lock(m_jobMutex); m_notifier.wait(lock, [this] {return !m_jobQueue.empty() || m_shutdown.load(std::memory_order_relaxed); }); if (m_shutdown.load(std::memory_order_relaxed) || m_jobQueue.empty()) { break; } job = std::move(m_jobQueue.front()); m_jobQueue.pop(); } job(); m_finished.notify_one(); } })); } } };
测试代码
void threader (int x) { std::cout<<"In threaded function: "<<x<<std::endl; } int main() { //outer loop for (auto i = 0; i < 10000; i++) { //Thread pool int num_threads = std::thread::hardware_concurrency(); ThreadPool test_pool(num_threads); // Assign work for (int j = 0; j < 48; j++) { test_pool.add_job(std::bind(threader, j)); } test_pool.waitFinished(); std::cout<<"Thread Pool Done"<<std::endl; } }
核心错误分析
- 等待条件不完整:
waitFinished仅等待任务队列空,但未确认所有已取出队列的任务是否执行完成。此时触发shutdown逻辑,可能导致部分线程仍在执行任务,后续出现竞态。 - 线程退出逻辑错误:线程在队列空但未shutdown时直接退出,违背线程池线程应持续等待任务的设计,可能导致线程提前退出或卡在wait状态。
- 内存序使用不当:
m_shutdown使用memory_order_relaxed,无法保证线程间的可见性,部分线程可能无法及时读取到shutdown的更新,导致卡在m_notifier.wait中无法退出。 - 条件变量
m_finished误用:用它等待队列空完全错误,它应该配合活跃任务计数来通知所有任务执行完成。
修复方案
修改后的线程池代码
#include <queue> #include <mutex> #include <condition_variable> #include <functional> #include <atomic> #include <vector> #include <thread> #include <iostream> class ThreadPool { public: ThreadPool() : m_busy(0) { m_shutdown.store(false, std::memory_order_release); createThreads(1); } ThreadPool(std::size_t numThreads) : m_busy(0) { m_shutdown.store(false, std::memory_order_release); createThreads(numThreads); } void add_job(std::function<void()> new_job) { { std::scoped_lock<std::mutex> lock(m_jobMutex); m_jobQueue.push(std::move(new_job)); } m_notifier.notify_one(); } void waitFinished() { { std::unique_lock<std::mutex> lock(m_jobMutex); // 等待队列空且所有任务执行完成 m_finished.wait(lock, [this] { return m_jobQueue.empty() && m_busy == 0; }); } // 设置shutdown并保证可见性 m_shutdown.store(true, std::memory_order_release); m_notifier.notify_all(); for (std::thread& th : m_threads) { if (th.joinable()) { th.join(); } } m_threads.clear(); } private: using Job = std::function<void()>; std::vector<std::thread> m_threads; std::queue<Job> m_jobQueue; std::condition_variable m_notifier; std::condition_variable m_finished; std::mutex m_jobMutex; std::atomic<bool> m_shutdown; std::atomic<int> m_busy; // 记录正在执行任务的线程数 void createThreads(std::size_t numThreads) { m_threads.reserve(numThreads); for (std::size_t i = 0; i < numThreads; ++i) { m_threads.emplace_back(std::thread([this]() { while (true) { Job job; { std::unique_lock<std::mutex> lock(m_jobMutex); // 等待任务或shutdown信号,保证shutdown的可见性 m_notifier.wait(lock, [this] { return !m_jobQueue.empty() || m_shutdown.load(std::memory_order_acquire); }); // 仅当shutdown且队列空时才退出 if (m_shutdown.load(std::memory_order_acquire) && m_jobQueue.empty()) { break; } // 队列非空,取出任务 job = std::move(m_jobQueue.front()); m_jobQueue.pop(); m_busy++; // 活跃任务数+1 } // 执行任务 try { job(); } catch (...) { // 捕获任务异常,避免线程崩溃 std::cerr << "Job execution threw an exception" << std::endl; } { std::scoped_lock<std::mutex> lock(m_jobMutex); m_busy--; // 活跃任务数-1 } m_finished.notify_all(); // 通知任务完成 } })); } } };
关键修改点
- 添加
m_busy活跃任务计数器:精确追踪正在执行任务的线程数量,确保waitFinished等待所有任务执行完成后再触发shutdown。 - 修正
waitFinished的等待条件:改为等待m_jobQueue.empty() && m_busy == 0,保证队列中无待执行任务且所有已取出的任务都执行完毕。 - 调整线程退出逻辑:仅当
shutdown为true且队列空时才退出线程,避免线程提前退出。 - 修正内存序:
m_shutdown使用memory_order_release和memory_order_acquire,确保线程间的可见性,避免线程无法感知shutdown信号。 - 异常安全处理:捕获任务执行中的异常,防止线程崩溃导致线程池异常。
m_finished改为notify_all:确保所有等待的线程都能收到任务完成的通知,避免遗漏信号。
内容的提问来源于stack exchange,提问作者fuzzy_motion
相关产品推荐
相关产品推荐

