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

如何实现仅执行队列最新未执行任务的std::queue线程安全任务队列

仅执行最新待执行任务的单线程任务队列实现

原代码问题分析

原代码无法正常运行的核心原因集中在线程同步逻辑的错误:

  • 互斥量使用混乱:队列、触发标记等共享资源被多把互斥量交叉保护,存在线程安全漏洞,会出现数据竞争
  • 条件变量使用不符合规范:修改触发标记m_trigger时持有的是m_push_mutex,但条件变量等待时持有的是m_queue_mutex,违反了std::condition_variable的使用要求,会导致通知丢失、死锁等问题
  • 入队逻辑不符合需求:队列空时直接在调用PushCopyTask的线程执行任务,不符合“独立线程执行任务”的设计目标,任务调度逻辑混乱
  • 任务完成通知与队列状态不匹配:用户自定义的任务执行线程和队列调度线程没有状态同步,任务完成的通知无法正确触达队列调度逻辑

修正后实现

以下实现满足需求:仅由独立线程执行任务,每次调度时取出队列中最新的任务执行,取出后立即清空队列,不存在待执行任务时不会误清空队列,API简洁易用:

#include <atomic>
#include <functional>
#include <iostream>
#include <mutex>
#include <queue>
#include <thread>
#include <chrono>

using CopyTask = std::function<void(void)>;

class TaskQueue {
public:
    TaskQueue() {
        // 构造时直接启动调度线程
        m_running = true;
        m_worker_thread = std::thread(&TaskQueue::ThreadLoop, this);
    }

    ~TaskQueue() {
        {
            std::lock_guard<std::mutex> lock(m_mutex);
            m_running = false;
        }
        m_cv.notify_one();
        if (m_worker_thread.joinable()) {
            m_worker_thread.join();
        }
    }

    // 对外仅暴露Push接口,添加任务到队列
    void PushCopyTask(CopyTask task) {
        std::lock_guard<std::mutex> lock(m_mutex);
        m_queue.push(std::move(task));
        // 有新任务入队,通知调度线程
        m_cv.notify_one();
    }

private:
    std::queue<CopyTask> m_queue;
    std::mutex m_mutex;
    std::condition_variable m_cv;
    std::thread m_worker_thread;
    std::atomic<bool> m_running;

    void ThreadLoop() {
        while (m_running) {
            std::unique_lock<std::mutex> lock(m_mutex);
            // 等待直到有任务入队或者队列停止运行
            m_cv.wait(lock, [this] {
                return !m_running || !m_queue.empty();
            });

            if (!m_running) {
                break;
            }

            // 取出最新的任务
            CopyTask latest_task = std::move(m_queue.back());
            // 清空队列,丢弃所有旧任务
            m_queue = {};
            // 执行任务前先解锁,避免执行任务时阻塞其他线程Push任务
            lock.unlock();

            // 执行最新任务
            latest_task();
        }
    }
};

// 测试用例
int main() {
    std::cout << " -- 程序开始 -- " << std::endl;
    
    TaskQueue task_queue;

    // 连续推入4个测试任务,最终仅会执行最新的第4个任务
    task_queue.PushCopyTask([](){
        std::this_thread::sleep_for(std::chrono::milliseconds(70));
        std::cout << " 任务1执行完成 " << std::endl;
    });

    task_queue.PushCopyTask([](){
        std::this_thread::sleep_for(std::chrono::milliseconds(70));
        std::cout << " 任务2执行完成 " << std::endl;
    });

    task_queue.PushCopyTask([](){
        std::this_thread::sleep_for(std::chrono::milliseconds(70));
        std::cout << " 任务3执行完成 " << std::endl;
    });

    task_queue.PushCopyTask([](){
        std::this_thread::sleep_for(std::chrono::milliseconds(70));
        std::cout << " 任务4执行完成 " << std::endl;
    });

    // 等待任务执行完成
    std::this_thread::sleep_for(std::chrono::seconds(1));
    std::cout << " -- 程序结束 -- " << std::endl;
    return 0;
}

实现说明

  • 所有共享资源(任务队列、运行状态)都由同一把互斥量m_mutex保护,避免线程安全问题
  • 条件变量的等待和通知都绑定同一把互斥量,不会出现通知丢失的问题
  • 调度线程每次被唤醒后,直接取队列尾部的最新任务,清空队列后再执行任务,符合“仅执行最新任务”的需求
  • 执行任务前先解锁互斥量,避免任务执行过程中阻塞其他线程推入新任务
  • API仅保留构造、析构和PushCopyTask三个对外接口,使用简洁,不需要用户手动调用任务完成通知,调度逻辑内部自动处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 03:06:03