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

如何在任务执行线程中同步挂起线程并解锁互斥锁

解决方案:任务队列的线程等待与唤醒

问题分析

你原来的代码中,线程在解锁互斥锁(标记1)和调用SuspendThread(标记2)之间存在竞态窗口:如果此时addJob刚好取出空闲线程并调用ResumeThread,线程可能还没进入挂起状态,导致ResumeThread无效,后续线程依然会挂起,最终任务无法被处理,同时freeThreads里的线程编号也会和实际状态不一致。


Windows API 解决方案(使用事件对象)

Windows下不要用SuspendThread实现线程等待,这是不安全的操作。正确做法是用手动重置事件(Manual-Reset Event),让线程在无任务时等待事件信号,有新任务时触发事件唤醒线程。

核心思路

  1. 每个线程对应一个事件对象,初始为无信号状态。
  2. 线程无任务时,持有互斥锁的情况下加入freeThreads,释放锁后等待事件信号。
  3. addJob添加任务后,取出空闲线程,触发其事件唤醒线程,同时从freeThreads中移除该线程编号。

完整代码示例

#include <Windows.h>

#include <iostream>
#include <vector>
#include <mutex>
#include <memory>
#include <stack>
#include <queue>

class Job {};

std::queue<Job> jobs;
std::mutex mutex;
std::vector<HANDLE> threads;
std::vector<HANDLE> threadEvents;
std::stack<std::size_t> freeThreads;
std::atomic_bool stopRequested { false };

void DoJob(Job job = {}) {
    // 实际执行HTTP请求等任务
}

DWORD WINAPI ThreadedMethod(LPVOID param) {
    std::size_t threadNumber = *static_cast<std::size_t*>(param);
    delete static_cast<std::size_t*>(param); // 释放传入的参数内存

    while (!stopRequested) {
        std::lock_guard<std::mutex> lock(mutex);
        if (!jobs.empty()) {
            Job job = std::move(jobs.front());
            jobs.pop();
            // 手动提前释放锁,避免任务执行时持有锁
            lock.~lock_guard();
            DoJob(std::move(job));
        } else {
            freeThreads.push(threadNumber);
            // 释放锁后等待事件
            lock.~lock_guard();
            // 等待事件信号或停止请求
            DWORD waitResult = WaitForSingleObject(threadEvents[threadNumber], INFINITE);
            if (waitResult == WAIT_OBJECT_0) {
                // 被唤醒后重置事件,等待下一次触发
                ResetEvent(threadEvents[threadNumber]);
            }
        }
    }
    return 0;
}

void addJob(Job job) {
    std::lock_guard<std::mutex> lock(mutex);
    jobs.push(std::move(job));

    if (!freeThreads.empty()) {
        std::size_t threadNum = freeThreads.top();
        freeThreads.pop();
        // 触发事件唤醒线程
        SetEvent(threadEvents[threadNum]);
    }
    // 如果没有空闲线程,任务会留在队列,线程处理完当前任务后会自动取走
}

// 线程初始化示例
void InitThreads(std::size_t count) {
    threads.resize(count);
    threadEvents.resize(count);

    for (std::size_t i = 0; i < count; ++i) {
        // 创建手动重置事件,初始无信号
        threadEvents[i] = CreateEvent(nullptr, TRUE, FALSE, nullptr);
        // 传入线程编号(动态分配避免栈内存失效)
        std::size_t* threadNum = new std::size_t(i);
        threads[i] = CreateThread(nullptr, 0, ThreadedMethod, threadNum, 0, nullptr);
    }
}

// 线程清理示例
void CleanupThreads() {
    stopRequested = true;
    // 唤醒所有等待的线程
    for (HANDLE event : threadEvents) {
        SetEvent(event);
    }
    // 等待所有线程退出
    for (HANDLE thread : threads) {
        WaitForSingleObject(thread, INFINITE);
        CloseHandle(thread);
    }
    // 关闭事件句柄
    for (HANDLE event : threadEvents) {
        CloseHandle(event);
    }
}

标准C++ 解决方案(使用std::condition_variable)

C++11及以上标准提供了std::condition_variable,专门解决生产者-消费者的等待唤醒问题,跨平台且安全。

核心思路

  1. 用std::condition_variable代替Windows事件,线程在无任务时等待条件变量。
  2. addJob添加任务后,通知一个等待的线程(或所有线程)。
  3. 线程被唤醒后重新检查队列状态(避免虚假唤醒)。

完整代码示例

#include <iostream>
#include <vector>
#include <mutex>
#include <queue>
#include <thread>
#include <condition_variable>
#include <atomic>

class Job {};

std::queue<Job> jobs;
std::mutex mutex;
std::condition_variable cv;
std::vector<std::thread> threads;
std::atomic_bool stopRequested { false };

void DoJob(Job job = {}) {
    // 实际执行HTTP请求等任务
}

void ThreadedMethod() {
    while (!stopRequested) {
        std::unique_lock<std::mutex> lock(mutex);
        // 等待条件:队列非空或停止请求
        cv.wait(lock, []{ return !jobs.empty() || stopRequested; });

        if (stopRequested) {
            break;
        }

        Job job = std::move(jobs.front());
        jobs.pop();
        // 释放锁后执行任务
        lock.unlock();
        DoJob(std::move(job));
    }
}

void addJob(Job job) {
    std::lock_guard<std::mutex> lock(mutex);
    jobs.push(std::move(job));
    // 通知一个等待的线程
    cv.notify_one();
}

// 线程初始化示例
void InitThreads(std::size_t count) {
    threads.reserve(count);
    for (std::size_t i = 0; i < count; ++i) {
        threads.emplace_back(ThreadedMethod);
    }
}

// 线程清理示例
void CleanupThreads() {
    stopRequested = true;
    // 通知所有线程退出
    cv.notify_all();
    // 等待所有线程完成
    for (auto& t : threads) {
        if (t.joinable()) {
            t.join();
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 11:44:52