Skip to content

Folders and files

NameName
Last commit message
Last commit date

Latest commit

 

History

17 Commits
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

cpp_thread_pool

English | 简体中文

一个基于 C++17 标准库实现的轻量级线程池,用于学习任务调度、并发同步和线程生命周期管理。

支持固定线程模式、Cached 动态扩容、std::future 返回值、任务异常传递、有界队列、超时拒绝、优雅关闭和空闲线程回收。

C++17 CMake Platform CI License Stars


📖 项目简介

cpp_thread_pool 是一个面向学习与实践的 C++ 线程池。项目围绕一条完整的并发执行链路展开:

submit() 提交任务 → 有界队列缓冲任务 → worker 线程取出并执行 → std::future 返回结果或异常 → shutdown() 完成关闭与线程回收。

项目重点处理了队列满、并发关闭、Cached 线程退出、线程创建失败和重复启动等边界情况。当前版本适合学习、面试展示和小型实验;用于生产环境前仍需结合实际负载补充压力测试、性能评估和故障处理。

✨ 核心特性

  • 仅依赖 C++17 标准库和系统线程库;
  • 支持 MODE_FIXED 固定线程模式;
  • 支持 MODE_CACHED 动态扩容与空闲线程超时回收;
  • 支持 Lambda、普通函数、函数对象和带参数任务;
  • 使用 std::packaged_task 执行任务,通过 std::future<T> 返回结果;
  • 支持 void 返回值,并将任务异常传递到 future.get()
  • 使用有界任务队列限制积压量;
  • 队列满时按可配置超时时间等待,超时后拒绝任务;
  • 支持显式、幂等的 shutdown(),析构时自动关闭;
  • 已入队任务会在关闭时继续执行,关闭完成后允许再次 start()
  • 提供行为测试,并通过 GitHub Actions 在 Windows、Linux 和 macOS 上构建验证。

🏗️ 项目架构

flowchart LR
    Caller[调用方] -->|submit| Submit[任务提交层]
    Submit -->|封装 packaged_task| Queue[有界任务队列]
    Queue -->|notEmpty 唤醒| Workers[工作线程组]
    Workers -->|取出并执行| Task[用户任务]
    Task -->|结果或异常| Future[std::future]
    Future --> Caller

    Config[配置与生命周期 API] --> State[线程池共享状态]
    State --> Submit
    State --> Workers
    Workers -->|notFull 唤醒| Submit
    Reaper[Cached 线程回收] --> Workers
Loading

组件职责

组件 主要职责 对应实现
公共接口层 配置、启动、提交任务、查询状态和关闭线程池 include/threadpool.h
任务封装层 推导返回类型,将可调用对象封装为 std::packaged_task ThreadPool::submit()
调度与同步层 维护有界队列,协调生产者与 worker,处理提交超时 taskQue_notEmpty_notFull_
工作线程层 等待任务、执行任务,维护当前线程数和空闲线程数 threadFunc()
动态扩容层 根据任务积压创建额外 worker,并限制最大线程数 addThreadLocked()
线程回收层 标记已退出的 Cached worker,在锁外 join 并移除线程对象 reapFinishedThreadsLocked()
生命周期层 管理启动、优雅关闭、幂等关闭和关闭后重启 start()shutdown()

并发模型

  • taskQueMtx_ 保护任务队列、worker 状态、线程计数和运行状态;
  • notEmpty_ 在任务入队或线程池关闭时唤醒 worker;
  • notFull_ 在任务出队或线程池关闭时唤醒等待提交的线程;
  • join() 在释放 taskQueMtx_ 后执行,避免 worker 退出时形成锁等待;
  • Cached worker 先在锁内标记 finished,再由线程池统一回收对应的 std::thread
  • submit() 在条件变量等待和线程回收后都会重新检查线程池状态与队列容量。

更完整的设计说明见 Design NotesThread Pool Lifecycle

📁 项目结构

cpp_thread_pool/
├── .github/
│   └── workflows/
│       └── ci.yml                 # Windows、Linux、macOS 持续集成
├── docs/
│   ├── design.md                  # 设计取舍和实现笔记
│   └── threadpool_lifecycle.md    # 提交、执行、回收和关闭流程
├── examples/
│   └── example.cpp                # 使用示例
├── include/
│   └── threadpool.h               # 公共接口和 submit() 模板实现
├── src/
│   └── threadpool.cpp             # 生命周期、worker 和线程回收逻辑
├── tests/
│   └── threadpool_tests.cpp       # 关键行为测试
├── CMakeLists.txt                 # 构建与测试配置
├── README.md                      # 简体中文文档
├── README_EN.md                   # English documentation
└── LICENSE

🧰 环境要求

工具 要求
C++ 标准 C++17 或更高
CMake 3.10 或更高
编译器 GCC、Clang 或 MSVC

推荐使用 GCC 9+、Clang 10+ 或 MSVC 2019+。

🚀 构建与测试

cmake -S . -B build -DCMAKE_BUILD_TYPE=Release
cmake --build build --config Release --parallel
ctest --test-dir build --build-config Release --output-on-failure

运行示例:

# Visual Studio / MSVC
./build/Release/threadpool_demo.exe

# Ninja、MinGW、Linux 或 macOS
./build/threadpool_demo

GitHub Actions 会在向 main 分支推送代码或创建 Pull Request 时,分别在 Ubuntu、Windows 和 macOS 上执行配置、构建和测试。

💡 快速开始

#include "threadpool.h"

#include <iostream>

int main() {
  ThreadPool pool;
  pool.start(4);

  auto result = pool.submit([](int a, int b) {
    return a + b;
  }, 10, 20);

  std::cout << result.get() << '\n';
  pool.shutdown();
}

submit() 会自动推导返回类型。上例中 result 的类型为 std::future<int>

其他常见任务形式:

auto value = pool.submit([] { return 42; });       // std::future<int>
auto output = pool.submit([] {});                  // std::future<void>
auto text = pool.submit([] { return std::string("ok"); });

需要引用语义时请显式使用 std::refstd::cref,否则函数对象和参数会被拷贝或移动到任务内部。

任务异常

任务抛出的异常由 std::packaged_task 保存,并在调用 future.get() 时重新抛出:

auto result = pool.submit([]() -> int {
  throw std::runtime_error("task failed");
});

try {
  result.get();
} catch (const std::exception &e) {
  std::cerr << e.what() << '\n';
}

⚙️ 运行模式

Fixed 模式

默认使用 MODE_FIXED。线程池启动后创建固定数量的 worker,运行期间不会自动增减:

ThreadPool pool;
pool.start(4);

Cached 模式

Cached 模式会根据任务积压动态创建 worker。额外线程空闲超过设定时间后自动退出,随后由线程池完成 join 和容器清理:

ThreadPool pool;
pool.setMode(ThreadPoolMode::MODE_CACHED);
pool.setThreadSizeThreshold(8);
pool.setThreadMaxIdleTime(std::chrono::seconds(1));
pool.start(2);
配置 含义
start(2) 初始创建 2 个 worker,也是回收后的保底线程数
setThreadSizeThreshold(8) 最大允许 8 个 worker
setThreadMaxIdleTime(std::chrono::seconds(1)) 额外 worker 空闲超过 1 秒后退出

⏱️ 提交与拒绝策略

任务队列默认容量为 1024。当队列已满时,submit() 等待空位的默认超时时间为 1s,可在启动前修改:

ThreadPool pool;
pool.setTaskQueMaxThreshold(64);
pool.setSubmitTimeout(std::chrono::milliseconds(500));
pool.start(4);

submit() 可能抛出以下异常:

场景 异常信息
线程池未启动 Thread pool is not running
队列在超时时间内始终没有空位 Task queue is full; submit timed out
等待期间线程池进入关闭流程 Thread pool has stopped or is shutting down

提交超时只限制“队列满时等待空位”的时间,不限制任务执行时间,也不是任务取消机制。

🛑 关闭与重启

pool.shutdown();

shutdown() 会停止接收新任务,让已入队任务继续执行,等待所有 worker 退出并回收线程资源。它具有以下语义:

  • 可以重复调用,也支持多个外部线程并发调用;
  • 函数返回时,本轮关闭流程已经完成;
  • 关闭完成后,同一个对象可以重新配置并再次 start()
  • 未显式调用时,析构函数会自动执行关闭流程。

禁止在线程池自己的任务内部调用 shutdown(),也不要在 worker 任务内部销毁 ThreadPool 对象。若已入队任务永久阻塞,shutdown() 也会一直等待。

📚 API 速查

所有配置接口都应在 start() 前调用;shutdown() 完成后可重新配置。

API 说明
setMode(ThreadPoolMode mode) 设置 Fixed 或 Cached 模式
setTaskQueMaxThreshold(std::size_t threshold) 设置任务队列容量
setSubmitTimeout(std::chrono::milliseconds timeout) 设置队列满时的提交等待上限
setThreadSizeThreshold(std::size_t threshold) 设置 Cached 模式最大线程数
setThreadMaxIdleTime(std::chrono::seconds idleTime) 设置额外 Cached worker 的最大空闲时间
start(std::size_t initThreadSize = 4) 启动线程池
submit(F&& func, Args&&... args) 提交任务并返回 std::future<T>
currentThreadSize() const 查询当前 worker 总数
idleThreadSize() const 查询当前空闲 worker 数
shutdown() 停止接收任务,等待已入队任务完成并回收资源

🧩 设计文档

  • Design Notes:记录 std::future、任务参数所有权、Cached 线程回收、关闭语义和提交超时等设计取舍;
  • Thread Pool Lifecycle:使用 Mermaid 图展示任务提交、worker 执行、Cached 回收和关闭流程。

⚠️ 已知边界

  • 当前不支持任务取消、优先级队列、work stealing 和暂停/恢复;
  • 不支持为每次 submit() 单独设置超时时间;
  • std::future::get() 只能调用一次;
  • 状态查询接口主要用于学习和观察,不代表性能监控系统;
  • 当前测试以功能行为为主,尚未覆盖长时间压力测试和性能基准。

🗺️ 后续计划

  • 增加多生产者并发提交和长时间压力测试;
  • 增加 AddressSanitizer、UndefinedBehaviorSanitizer 和 ThreadSanitizer 构建;
  • 增加性能基准与不同线程数下的吞吐量对比;
  • 探索任务取消、优先级队列和每次提交独立超时。

📄 License

本项目基于 MIT License 开源。

About

A lightweight C++17 thread pool with fixed/cached modes, futures, bounded queue, configurable submit timeout, graceful shutdown, tests, and CI.

Topics

Resources

Stars

Watchers

Forks

Contributors

Languages