35#include <condition_variable>
85 Task(std::packaged_task<
T()> &&t) :
task(std::move(t)) {}
90 std::queue<std::unique_ptr<TaskBase>>
tasks;
92 std::condition_variable
cv;
106 size_t num_threads = std::thread::hardware_concurrency()) :
108 if (num_threads == 0) num_threads = 2;
109 for (
size_t i = 0; i < num_threads; ++i) {
112 std::unique_ptr<TaskBase> task;
114 std::unique_lock<std::mutex> lock(
mutex);
115 cv.wait(lock, [
this] {
116 return shutdown.load(std::memory_order_acquire)
119 if (
shutdown.load(std::memory_order_acquire)
122 task = std::move(
tasks.front());
129 1, std::memory_order_release);
136 1, std::memory_order_release);
182 template <
class F,
class... Args>
183 auto submit(F &&f, Args &&...args) -> std::future<
184 decltype(std::forward<F>(f)(std::forward<Args>(args)...))> {
186 =
decltype(std::forward<F>(f)(std::forward<Args>(args)...));
188 = std::bind(std::forward<F>(f), std::forward<Args>(args)...);
189 std::packaged_task<return_type()> task(bound_task);
190 std::future<return_type> result = task.get_future();
192 std::lock_guard<std::mutex> lock(
mutex);
193 if (
shutdown.load(std::memory_order_acquire))
194 throw AsyncError(
"ThreadPool is shutting down");
195 tasks.emplace(std::unique_ptr<TaskBase>(
215 template <
class F,
class Container>
216 auto map(F &&func,
const Container &items)
217 -> std::vector<std::future<
decltype(func(
218 std::declval<typename Container::value_type>()))>> {
219 using return_type =
decltype(func(
220 std::declval<typename Container::value_type>()));
221 std::vector<std::future<return_type>> futures;
222 futures.reserve(items.size());
223 for (
const auto &item : items)
224 futures.push_back(
submit(func, item));
234 std::unique_lock<std::mutex> lock(
mutex);
235 cv.wait(lock, [
this] {
247 std::lock_guard<std::mutex> lock(
mutex);
274 shutdown.store(
true, std::memory_order_release);
278 if (worker.joinable()) worker.join();
288 shutdown.store(
true, std::memory_order_release);
290 std::lock_guard<std::mutex> lock(
mutex);
291 std::queue<std::unique_ptr<TaskBase>>().swap(
tasks);
295 if (worker.joinable()) worker.detach();
326 template <
class F,
class... Args>
327 inline auto submit(F &&f, Args &&...args) -> std::future<
328 decltype(std::forward<F>(f)(std::forward<Args>(args)...))> {
330 std::forward<F>(f), std::forward<Args>(args)...);
341 template <
class F,
class Container>
343 map(F &&f,
const Container &c) -> std::vector<std::future<
decltype(f(
344 std::declval<typename Container::value_type>()))>> {
404 template <
class F,
class... Args>
405 inline auto async(F &&f, Args &&...args) -> std::future<
406 decltype(std::forward<F>(f)(std::forward<Args>(args)...))> {
407 return pool::submit(std::forward<F>(f), std::forward<Args>(args)...);
AsyncError(const std::string &msg)
构造 ConsoleError。
Definition csexc.h:193
线程池执行器类,用于管理和执行并发任务。
Definition pool.h:57
auto submit(F &&f, Args &&...args) -> std::future< decltype(std::forward< F >(f)(std::forward< Args >(args)...))>
提交一个可调用对象到线程池执行。
Definition pool.h:183
size_t active_worker_count() const
获取当前存活的工作线程数。
Definition pool.h:265
std::atomic< bool > shutdown
线程池关闭标志
Definition pool.h:93
std::atomic< size_t > active_tasks
当前正在执行的任务数
Definition pool.h:94
std::mutex mutex
保护任务队列的互斥锁
Definition pool.h:91
~ThreadPool()
析构函数,根据退出策略决定关闭方式。
Definition pool.h:162
std::vector< std::thread > workers
工作线程容器
Definition pool.h:89
launch
Definition pool.h:59
@ close
退出时调用 close
Definition pool.h:61
auto map(F &&func, const Container &items) -> std::vector< std::future< decltype(func(std::declval< typename Container::value_type >()))> >
批量提交任务,对容器中的每个元素应用函数。
Definition pool.h:216
void close()
优雅关闭线程池,等待所有任务完成。
Definition pool.h:273
size_t waiting_task_count() const
获取当前任务队列中的任务数量。
Definition pool.h:246
void abort()
粗暴地关闭线程池,用于需要立刻结束一切。
Definition pool.h:287
size_t active_task_count() const
获取当前正在执行的任务数量。
Definition pool.h:256
ThreadPool(const ThreadPool &)=delete
禁止拷贝构造。
std::queue< std::unique_ptr< TaskBase > > tasks
待执行任务队列
Definition pool.h:90
std::condition_variable cv
用于线程等待和唤醒的条件变量
Definition pool.h:92
ThreadPool(launch exit_launch=launch::close, size_t num_threads=std::thread::hardware_concurrency())
构造函数,创建指定数量的工作线程。
Definition pool.h:105
ThreadPool(size_t num_threads, launch exit_launch=launch::close)
构造函数,创建指定数量的工作线程。
Definition pool.h:154
ThreadPool(ThreadPool &&)=delete
禁止移动构造。
ThreadPool & operator=(ThreadPool &&)=delete
禁止移动赋值。
launch exit_launch_
退出策略
Definition pool.h:95
ThreadPool & operator=(const ThreadPool &)=delete
禁止拷贝赋值。
void wait()
等待所有已提交的任务完成。
Definition pool.h:233
定义 console 库使用的自定义异常类层次结构。
size_t active_worker_count()
获取当前正在执行的工作线程数量。
Definition pool.h:375
size_t waiting_task_count()
获取等待执行的任务数量。
Definition pool.h:359
void close()
关闭线程池。
Definition pool.h:382
ThreadPool & instance()
获取线程池单例。
Definition pool.h:313
auto submit(F &&f, Args &&...args) -> std::future< decltype(std::forward< F >(f)(std::forward< Args >(args)...))>
提交任务到线程池。
Definition pool.h:327
size_t active_task_count()
获取当前正在执行的任务数量。
Definition pool.h:367
void abort()
中止线程池。
Definition pool.h:389
void wait()
等待所有任务完成。
Definition pool.h:351
auto map(F &&f, const Container &c) -> std::vector< std::future< decltype(f(std::declval< typename Container::value_type >()))> >
对容器中的元素应用函数,并返回结果的 future 对象。
Definition pool.h:343
auto async(F &&f, Args &&...args) -> std::future< decltype(std::forward< F >(f)(std::forward< Args >(args)...))>
提交一个可调用对象到全局线程池执行。
Definition pool.h:405
任务基类,提供多态接口。
Definition pool.h:71
virtual ~TaskBase()=default
具体任务类模板,封装 std::packaged_task。
Definition pool.h:83
Task(std::packaged_task< T()> &&t)
Definition pool.h:85
void execute() override
Definition pool.h:86
std::packaged_task< T()> task
Definition pool.h:84