Console Library 8.0.0
A header-only library that makes C++ simple
Loading...
Searching...
No Matches
pool.h
Go to the documentation of this file.
1
10
11/*
12Copyright (c) 2026 MrXie1109
13
14Permission is hereby granted, free of charge, to any person obtaining a copy
15of this software and associated documentation files (the "Software"), to deal
16in the Software without restriction, including without limitation the rights
17to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
18copies of the Software, and to permit persons to whom the Software is
19furnished to do so, subject to the following conditions:
20
21The above copyright notice and this permission notice shall be included in all
22copies or substantial portions of the Software.
23
24THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
25IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
26FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
27AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
28LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
29OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
30SOFTWARE.
31*/
32
33#pragma once
34#include <atomic>
35#include <condition_variable>
36#include <functional>
37#include <future>
38#include <memory>
39#include <mutex>
40#include <queue>
41#include <thread>
42#include <type_traits>
43#include <utility>
44#include <vector>
45
46#include "../core/csexc.h"
47
48namespace console {
57 class ThreadPool {
58 public:
59 enum class launch : bool {
62 };
63
64 private:
71 struct TaskBase {
72 virtual ~TaskBase() = default;
73 virtual void execute() = 0;
74 };
75
82 template <class T>
83 struct Task : public TaskBase {
84 std::packaged_task<T()> task;
85 Task(std::packaged_task<T()> &&t) : task(std::move(t)) {}
86 void execute() override { task(); }
87 };
88
89 std::vector<std::thread> workers;
90 std::queue<std::unique_ptr<TaskBase>> tasks;
91 mutable std::mutex mutex;
92 std::condition_variable cv;
93 std::atomic<bool> shutdown;
94 std::atomic<size_t> active_tasks;
96
97 public:
106 size_t num_threads = std::thread::hardware_concurrency()) :
107 shutdown(false), active_tasks(0), exit_launch_(exit_launch) {
108 if (num_threads == 0) num_threads = 2;
109 for (size_t i = 0; i < num_threads; ++i) {
110 workers.emplace_back([this] {
111 while (true) {
112 std::unique_ptr<TaskBase> task;
113 {
114 std::unique_lock<std::mutex> lock(mutex);
115 cv.wait(lock, [this] {
116 return shutdown.load(std::memory_order_acquire)
117 || !tasks.empty();
118 });
119 if (shutdown.load(std::memory_order_acquire)
120 && tasks.empty())
121 return;
122 task = std::move(tasks.front());
123 tasks.pop();
124 }
125 active_tasks.fetch_add(1, std::memory_order_release);
126 try {
127 task->execute();
128 active_tasks.fetch_sub(
129 1, std::memory_order_release);
130 if (active_tasks.load(std::memory_order_acquire)
131 == 0
132 && tasks.empty())
133 cv.notify_all();
134 } catch (...) {
135 active_tasks.fetch_sub(
136 1, std::memory_order_release);
137 if (active_tasks.load(std::memory_order_acquire)
138 == 0
139 && tasks.empty())
140 cv.notify_all();
141 }
142 }
143 });
144 }
145 }
146
154 ThreadPool(size_t num_threads, launch exit_launch = launch::close) :
155 ThreadPool(exit_launch, num_threads) {}
156
164 close();
165 else
166 abort();
167 }
168
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)...))> {
185 using return_type
186 = decltype(std::forward<F>(f)(std::forward<Args>(args)...));
187 auto bound_task
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();
191 {
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>(
196 new Task<return_type>(std::move(task))));
197 }
198 cv.notify_one();
199 return result;
200 }
201
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));
225 return futures;
226 }
227
233 void wait() {
234 std::unique_lock<std::mutex> lock(mutex);
235 cv.wait(lock, [this] {
236 return tasks.empty()
237 && active_tasks.load(std::memory_order_acquire) == 0;
238 });
239 }
240
246 size_t waiting_task_count() const {
247 std::lock_guard<std::mutex> lock(mutex);
248 return tasks.size();
249 }
250
256 size_t active_task_count() const {
257 return active_tasks.load(std::memory_order_acquire);
258 }
259
265 size_t active_worker_count() const { return workers.size(); }
266
273 void close() {
274 shutdown.store(true, std::memory_order_release);
275 cv.notify_all();
276 wait();
277 for (auto &worker : workers)
278 if (worker.joinable()) worker.join();
279 }
280
287 void abort() {
288 shutdown.store(true, std::memory_order_release);
289 {
290 std::lock_guard<std::mutex> lock(mutex);
291 std::queue<std::unique_ptr<TaskBase>>().swap(tasks);
292 }
293 cv.notify_all();
294 for (auto &worker : workers)
295 if (worker.joinable()) worker.detach();
296 }
297
299 ThreadPool(const ThreadPool &) = delete;
301 ThreadPool &operator=(const ThreadPool &) = delete;
303 ThreadPool(ThreadPool &&) = delete;
306 };
307
308 namespace pool {
314 static ThreadPool exe;
315 return exe;
316 }
317
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)...))> {
329 return instance().submit(
330 std::forward<F>(f), std::forward<Args>(args)...);
331 }
332
341 template <class F, class Container>
342 inline auto
343 map(F &&f, const Container &c) -> std::vector<std::future<decltype(f(
344 std::declval<typename Container::value_type>()))>> {
345 return instance().map(std::forward<F>(f), c);
346 }
347
351 inline void wait() {
352 instance().wait();
353 }
354
359 inline size_t waiting_task_count() {
360 return instance().waiting_task_count();
361 }
362
367 inline size_t active_task_count() {
368 return instance().active_task_count();
369 }
370
375 inline size_t active_worker_count() {
376 return instance().active_worker_count();
377 }
378
382 inline void close() {
383 return instance().close();
384 }
385
389 inline void abort() {
390 return instance().abort();
391 }
392 }
393
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)...);
408 }
409}
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 库使用的自定义异常类层次结构。
Definition pool.h:308
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
@ T
Definition kb.h:99
任务基类,提供多态接口。
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