C++ 手写工作窃取线程池(WorkStealing)从原理到实现,超详细教程(多桶同步队列 + 线程安全 + 无锁竞争 + 可直接运行)
·
前言
在 C++ 后端开发、线程池设计中,工作窃取(Work Stealing) 是目前高性能线程池的主流设计方案。 相比于传统固定线程池,它解决了一个核心痛点: 部分线程忙、部分线程闲 → 负载不均、CPU 利用率低。
工作窃取线程池的核心思想:
- 每个线程拥有自己的任务队列
- 自己队列有任务 → 优先消费自己队列
- 自己队列空了 → 去 “偷” 其他线程队列的任务
- 最大限度利用多核 CPU,避免线程饥饿
本篇文章带你从零实现一个可直接用于生产环境的工作窃取线程池,包含:
- 多桶线程安全队列
WSyncQueue - 工作窃取调度核心
WorkStealingThreadPool - 完整测试用例
- 所有 API、知识点逐行精讲
一、工作窃取线程池核心原理
1.1 传统线程池问题
- 全局单队列 → 所有线程抢一把锁 → 锁竞争严重
- 任务分配不均 → 有的线程累死,有的闲死
1.2 工作窃取原理
- 每个线程一个独立队列(多桶)
- 任务提交时轮询放入不同桶
- 线程先消费自己桶
- 空闲线程偷其他线程桶的任务
- 减少锁竞争,提升并行效率
二、整体架构设计
本项目分为两个核心文件:
-
SyncQueue_2.hpp多桶线程安全队列(任务存储)- 支持多队列
- 线程安全
- 超时等待
- 优雅停止
-
WorkStealingPool.hpp工作窃取线程池(调度核心)- 启动固定线程
- 绑定线程与队列
- 实现窃取逻辑
- 提交任务 / 运行任务
-
测试用例
Test05_13WS.cpp验证线程池功能
三、核心模块实现详解
3.1 多桶线程安全队列 WSyncQueue
作用
管理 N 个任务队列,提供线程安全的 Put/Take。
核心成员变量
std::vector<std::deque<T>> m_taskQueues; // 多桶队列
size_t m_bucketSize; // 桶数量 = 线程数
size_t m_maxSize; // 每个队列最大任务数
mutable std::mutex m_mutex; // 互斥锁
std::condition_variable m_notEmpty; // 队列不为空
std::condition_variable m_notFull; // 队列不满
size_t m_waitTime; // 超时时间
bool m_needStop; // 停止标记
核心 API 讲解
1. IsFull / IsEmpty
判断指定索引队列是否满 / 空。
2. Add 任务添加(核心)
template <class F>
int Add(F &&task, const size_t index)
{
std::unique_lock<std::mutex> locker(m_mutex);
// 等待队列不满 或 停止
bool waitret = m_notFull.wait_for(
locker,
std::chrono::milliseconds(m_waitTime),
[this, index] {
return m_needStop || !IsFull(index);
});
if (m_needStop) return 1;
if (!waitret) return 2; // 超时
m_taskQueues[index].push_back(std::forward<F>(task));
m_notEmpty.notify_all();
return 0;
}
知识点:
unique_lock:灵活锁,支持等待 / 解锁wait_for:超时等待,避免永久阻塞forward:完美转发,支持左值 / 右值notify_all:唤醒等待任务的线程
3. Put 对外提交任务
int Put(const T &task, size_t index);
int Put(T &&task, size_t index);
支持左值、右值两种任务。
4. Take 取任务(两个版本)
int Take(std::deque<T> &tque, size_t index); // 取整个队列
int Take(T &task, size_t index); // 取单个任务
作用:
- 取整个队列 → 效率极高
- 取单个任务 → 适合窃取
5. WaitStop 优雅停止
void WaitStop();
等待任务执行完,再唤醒所有线程,安全退出。
3.2 工作窃取线程池 WorkStealingThreadPool
核心思想
线程自己队列空 → 偷别人队列的任务执行
核心成员
size_t m_numThreads; // 线程数
WSyncQueue<Task> m_queue; // 多桶队列
vector<shared_ptr<thread>> m_threadgroup; // 线程组
atomic<bool> m_running; // 运行标记
once_flag m_flag; // 只停止一次
核心 API 讲解
1. Start 启动线程池
创建 N 个线程,每个线程执行 RunInThread。
2. RunInThread(窃取核心)
void RunInThread(const size_t index)
{
while (m_running)
{
deque<Task> taskque;
// 1. 先拿自己队列
if (m_queue.Take(taskque, index) == 0)
{
for (auto &task : taskque) task();
}
// 2. 自己空了 → 偷别人的
else
{
size_t i = threadIndex();
if (i != index && m_queue.Take(taskque, i) == 0)
{
for (auto &task : taskque) task();
}
}
}
}
这就是工作窃取!
3. AddTask 提交任务
void AddTask(Task &&task);
void AddTask(const Task &task);
轮询分配任务到不同桶。
4. Stop 停止线程池
void Stop();
安全停止所有线程。
四、测试用例(无参任务)
#include <iostream>
#include <thread>
#include <chrono>
#include "WorkStealingPool.hpp"
using namespace std;
using namespace tulun;
void taskWork()
{
cout << "threadId: " << this_thread::get_id() << " run task\n";
this_thread::sleep_for(chrono::milliseconds(200));
}
int main()
{
WorkStealingThreadPool pool(500, 8);
for (int i = 0; i < 30; ++i) {
pool.AddTask(taskWork);
}
this_thread::sleep_for(chrono::seconds(3));
pool.Stop();
return 0;
}
五、涉及的 C++ 知识点与 API 总结
5.1 多线程
std::thread:线程std::mutex:互斥锁std::unique_lock:智能锁std::condition_variable:条件变量wait_for:超时等待notify_all:唤醒线程
5.2 模板与通用编程
template <class T>std::forward完美转发std::function包装函数std::bind绑定函数
5.3 容器
std::vectorstd::deque双端队列(任务队列首选)
5.4 线程安全
mutable:const 函数中修改锁atomic<bool>:原子布尔once_flag:只执行一次
5.5 工作窃取核心思想
- 多队列分离
- 减少锁竞争
- 空闲线程主动偷任务
- 负载均衡
六、工作窃取线程池优点
✅ 高 CPU 利用率 ✅ 低锁竞争 ✅ 负载均衡 ✅ 可伸缩性强 ✅ 适合高并发服务
七、总结
本文实现了一个工业级工作窃取线程池:
- 多桶线程安全队列
- 任务窃取调度
- 超时等待
- 优雅停止
- 无参 / 有参任务都支持
- 可直接用于项目
需要源码的同学可以直接复制文章内代码,可编译可运行!
更多推荐

所有评论(0)