如何实现仅执行队列最新未执行任务的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
相关产品推荐
相关产品推荐

