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

线程池优化:避免std::condition_variable无效通知后任务执行异常排查

线程池优化后工作线程无法正确唤醒的问题

我正在编写一个线程池,主线程通过原子指针为工作线程分配任务。工作线程会检查该指针,若指针为nullptr,则等待std::condition_variable以避免不必要的CPU消耗;若指针非空,则执行任务。

主线程在分配任务后,即使工作线程当前正在执行任务而非等待状态(下一次循环会自动获取新任务),仍会调用工作线程条件变量的notify_one()。性能分析显示主线程在notify_one()调用上耗时显著,因此我尝试优化减少该开销,但优化后任务无法始终正确执行,工作线程未能在正确时机被唤醒。


原始代码

struct Worker
{
    bool wake_up {false};
    std::condition_variable condition;
    std::mutex mutex;

    std::atomic<Job*> next_job {nullptr};
    std::atomic<bool> stop {false};

    void run()
    {
        while (! stop.load())
        {
            auto* job = next_job.exchange(nullptr);

            if (job != nullptr)
                job->run();
            else
                wait_for_job();
        }
    }

    void wait_for_job()
    {
        std::unique_lock lock(mutex);
        if (! wake_up)
            condition.wait(lock, [this] { return wake_up; });
    }

    bool push_job (Job* job)
    {
        Job* expected = nullptr;

        if (next_job.compare_exchange_weak(expected, job))
        {
            std::lock_guard lock(mutex);
            wake_up = true;
            condition.notify_all();
            return true;
        }

        return false;
    }
};

int main()
{
    Worker w;
    std::thread t ([&w] { w.run(); });

    // 主线程运行事件循环,特定事件触发后台任务
    // 此处为事件循环的模拟
    for (int i = 0; i < 10000; ++i)
    {
        std::this_thread::sleep_for(std::chrono::milliseconds (50));
        w.push_job (chooseJob());
    }

    w.stop.store(true);
    t.join();
    return 0;
}

优化后的代码(main函数未改动)

enum class WaitState
{
    waiting,
    running,
    woken
};

struct Worker
{
    std::atomic<WaitState> waiting {WaitState::running};
    std::condition_variable condition;
    std::mutex mutex;

    std::atomic<Job*> next_job {nullptr};
    std::atomic<bool> stop {false};

    void run()
    {
        while (! stop.load())
        {
            auto* job = next_job.exchange(nullptr);

            if (job != nullptr)
                job->run();
            else
                wait_for_job();
        }
    }

    void wait_for_job()
    {
        std::unique_lock lock(mutex);
        auto w = WaitState::running;

        if (waiting.compare_exchange_strong(w, WaitState::waiting))
        {
            const auto stop_waiting = [this] { return waiting.load() == WaitState::woken; };
            condition.wait_for(lock, std::chrono::milliseconds(500), stop_waiting);
        }

        waiting.store(WaitState::running);
    }

    bool push_job (Job* job)
    {
        Job* expected = nullptr;

        if (next_job.compare_exchange_weak(expected, job))
        {
            std::lock_guard lock (mutex);

            auto w = WaitState::waiting;

            if (waiting.compare_exchange_strong (w, WaitState::woken))
                condition.notify_all();
        }
        // 此处遗漏了return语句
        return false;
    }
};

问题分析与修正

你的优化版本存在几个关键错误,导致工作线程无法被正确唤醒:

1. 变量名冲突与原子操作的锁滥用

在push_job中,局部变量w与结构体成员waiting命名混淆,且你在持有mutex的情况下对原子变量执行CAS操作——原子变量本身线程安全,加锁反而引发锁竞争,既浪费性能,又可能错过唤醒时机。

2. 工作线程的状态重置时机错误

wait_for_job中,无论条件变量是被唤醒还是超时,你直接将状态设为WaitState::running,忽略了主线程刚将状态改为woken的场景;同时,CAS修改状态的逻辑放在锁内部,导致工作线程持有锁等待,阻塞主线程的任务推送操作。

3. 唤醒逻辑的竞态条件

当工作线程处于从running切换到waiting的间隙时,主线程推送任务后若不执行唤醒,工作线程可能进入无意义的超时等待,无法及时获取新任务。

修正后的代码

enum class WaitState
{
    running,
    waiting,
    woken
};

struct Worker
{
    std::atomic<WaitState> state {WaitState::running};
    std::condition_variable condition;
    std::mutex mutex;

    std::atomic<Job*> next_job {nullptr};
    std::atomic<bool> stop {false};

    void run()
    {
        while (!stop.load(std::memory_order_acquire))
        {
            Job* job = next_job.exchange(nullptr, std::memory_order_acq_rel);
            if (job != nullptr)
            {
                job->run();
                state.store(WaitState::running, std::memory_order_release);
            }
            else
            {
                wait_for_job();
            }
        }
    }

    void wait_for_job()
    {
        std::unique_lock<std::mutex> lock(mutex);
        WaitState expected = WaitState::running;
        // 锁外执行CAS,减少锁持有时间
        if (state.compare_exchange_strong(expected, WaitState::waiting, std::memory_order_acq_rel))
        {
            // 等待时同时检查停止信号与唤醒状态
            condition.wait_for(lock, std::chrono::milliseconds(500), [this] {
                return stop.load(std::memory_order_acquire) || 
                       state.load(std::memory_order_acquire) == WaitState::woken;
            });

            // 被唤醒后重置状态
            if (state.load(std::memory_order_acquire) == WaitState::woken)
            {
                state.store(WaitState::running, std::memory_order_release);
            }
        }
    }

    bool push_job(Job* job)
    {
        Job* expected = nullptr;
        if (next_job.compare_exchange_weak(expected, job, std::memory_order_acq_rel))
        {
            WaitState target = WaitState::waiting;
            // 仅当工作线程处于等待状态时,才修改状态并唤醒
            if (state.compare_exchange_strong(target, WaitState::woken, std::memory_order_acq_rel))
            {
                std::lock_guard<std::mutex> lock(mutex);
                condition.notify_one();
            }
            return true;
        }
        return false;
    }
};

关键改动说明

  • 重命名waiting为state,避免变量名混淆。
  • 使用显式内存顺序(memory_order_acquire/release),确保状态同步的可见性。
  • 将CAS修改状态的逻辑移到锁外,减少锁持有时间;等待条件同时检查停止信号,避免线程无法退出。
  • 仅当成功将工作线程状态从waiting改为woken时,才调用notify_one,彻底减少不必要的系统调用开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 13:45:01