You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

线程池线程过载时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;
  }
}

核心错误分析

  1. 等待条件不完整:waitFinished仅等待任务队列空,但未确认所有已取出队列的任务是否执行完成。此时触发shutdown逻辑,可能导致部分线程仍在执行任务,后续出现竞态。
  2. 线程退出逻辑错误:线程在队列空但未shutdown时直接退出,违背线程池线程应持续等待任务的设计,可能导致线程提前退出或卡在wait状态。
  3. 内存序使用不当:m_shutdown使用memory_order_relaxed,无法保证线程间的可见性,部分线程可能无法及时读取到shutdown的更新,导致卡在m_notifier.wait中无法退出。
  4. 条件变量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(); // 通知任务完成
                }
            }));
        }
    }
};

关键修改点

  1. 添加m_busy活跃任务计数器:精确追踪正在执行任务的线程数量,确保waitFinished等待所有任务执行完成后再触发shutdown。
  2. 修正waitFinished的等待条件:改为等待m_jobQueue.empty() && m_busy == 0,保证队列中无待执行任务且所有已取出的任务都执行完毕。
  3. 调整线程退出逻辑:仅当shutdown为true且队列空时才退出线程,避免线程提前退出。
  4. 修正内存序:m_shutdown使用memory_order_release和memory_order_acquire,确保线程间的可见性,避免线程无法感知shutdown信号。
  5. 异常安全处理:捕获任务执行中的异常,防止线程崩溃导致线程池异常。
  6. m_finished改为notify_all:确保所有等待的线程都能收到任务完成的通知,避免遗漏信号。

内容的提问来源于stack exchange,提问作者fuzzy_motion

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.24 10:27:28