前言

        在 C++ 后端开发、线程池设计中,工作窃取(Work Stealing) 是目前高性能线程池的主流设计方案。 相比于传统固定线程池,它解决了一个核心痛点: 部分线程忙、部分线程闲 → 负载不均、CPU 利用率低。

        工作窃取线程池的核心思想:

  1. 每个线程拥有自己的任务队列
  2. 自己队列有任务 → 优先消费自己队列
  3. 自己队列空了 → 去 “偷” 其他线程队列的任务
  4. 最大限度利用多核 CPU,避免线程饥饿

        本篇文章带你从零实现一个可直接用于生产环境的工作窃取线程池,包含:

  • 多桶线程安全队列 WSyncQueue
  • 工作窃取调度核心 WorkStealingThreadPool
  • 完整测试用例
  • 所有 API、知识点逐行精讲

一、工作窃取线程池核心原理

1.1 传统线程池问题

  • 全局单队列 → 所有线程抢一把锁 → 锁竞争严重
  • 任务分配不均 → 有的线程累死,有的闲死

1.2 工作窃取原理

  1. 每个线程一个独立队列(多桶)
  2. 任务提交时轮询放入不同桶
  3. 线程先消费自己桶
  4. 空闲线程偷其他线程桶的任务
  5. 减少锁竞争,提升并行效率

二、整体架构设计

        本项目分为两个核心文件:

  1. SyncQueue_2.hpp 多桶线程安全队列(任务存储)

    • 支持多队列
    • 线程安全
    • 超时等待
    • 优雅停止
  2. WorkStealingPool.hpp 工作窃取线程池(调度核心)

    • 启动固定线程
    • 绑定线程与队列
    • 实现窃取逻辑
    • 提交任务 / 运行任务
  3. 测试用例 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::vector
  • std::deque 双端队列(任务队列首选)

5.4 线程安全

  • mutable:const 函数中修改锁
  • atomic<bool>:原子布尔
  • once_flag:只执行一次

5.5 工作窃取核心思想

  • 多队列分离
  • 减少锁竞争
  • 空闲线程主动偷任务
  • 负载均衡

六、工作窃取线程池优点

✅ 高 CPU 利用率 ✅ 低锁竞争 ✅ 负载均衡 ✅ 可伸缩性强 ✅ 适合高并发服务


七、总结

本文实现了一个工业级工作窃取线程池:

  • 多桶线程安全队列
  • 任务窃取调度
  • 超时等待
  • 优雅停止
  • 无参 / 有参任务都支持
  • 可直接用于项目

需要源码的同学可以直接复制文章内代码,可编译可运行!

更多推荐