FazBrowse GitHub Viewer
|
Trending
|
URL:
|
Home
Tools:
[Download Repo ZIP]
[View Raw Code]
[Original HTTPS Page]
taskflow/benchmarks/thread_pool/ThreadPool.hpp at master · taskflow/taskflow · GitHub
taskflow
/
taskflow
Public
Uh oh!
There was an error while loading.
Please reload this page
.
Notifications
You must be signed in to change notification settings
Fork
1.4k
Star
12.2k
Code
Issues
20
Pull requests
16
Actions
Security and quality
0
Insights
Additional navigation options
Code
Issues
Pull requests
Actions
Security and quality
Insights
Expand file tree
Breadcrumbs
taskflow
/
benchmarks
/
thread_pool
/
ThreadPool.hpp
Copy path
More file actions
More file actions
Latest commit
History
History
History
141 lines (121 loc) · 2.91 KB
Breadcrumbs
taskflow
/
benchmarks
/
thread_pool
/
ThreadPool.hpp
Copy path
File metadata and controls
141 lines (121 loc) · 2.91 KB
Raw
Copy raw file
Download raw file
Open symbols panel
Edit and raw actions
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
#
pragma
once
#
include
<
vector
>
#
include
<
queue
>
#
include
<
memory
>
#
include
<
thread
>
#
include
<
mutex
>
#
include
<
condition_variable
>
#
include
<
future
>
#
include
<
functional
>
#
include
<
stdexcept
>
#
include
<
map
>
#
include
<
type_traits
>
#
include
<
iostream
>
class
ThreadPool
{
public:
ThreadPool
(
size_t
);
template
<
class
F
,
class
... Args>
auto
enqueue
(F&& f, Args&&... args)
-> std::future<typename std::invoke_result<F, Args...>::type>;
~ThreadPool
();
int
thread_number
(std::thread::id id)
{
if
(id_map.
find
(id) != id_map.
end
())
return
(
int
)id_map[id];
return
-
1
;
}
size_t
num_threads
()
{
return
num_threads_;
}
static
ThreadPool*
get
()
{
return
instance
(
0
);
}
static
ThreadPool*
instance
(
uint32_t
numthreads)
{
std::unique_lock<std::mutex>
lock
(singleton_mutex);
if
(!singleton) {
singleton =
new
ThreadPool
(numthreads ? numthreads :
hardware_concurrency
());
}
return
singleton;
}
static
void
release
()
{
std::unique_lock<std::mutex>
lock
(singleton_mutex);
delete
singleton;
singleton =
nullptr
;
}
static
uint32_t
hardware_concurrency
()
{
return
std::thread::hardware_concurrency
();
}
private:
std::vector<std::thread> workers;
std::queue<std::function<
void
()>> tasks;
std::mutex queue_mutex;
std::condition_variable condition;
bool
stop;
std::map<std::thread::id,
size_t
> id_map;
size_t
num_threads_;
static
ThreadPool* singleton;
static
std::mutex singleton_mutex;
};
inline
ThreadPool::ThreadPool
(
size_t
threads) : stop(
false
), num_threads_(threads)
{
if
(threads ==
1
)
return
;
for
(
size_t
i =
0
; i < threads; ++i)
workers.
emplace_back
([
this
] {
for
(;;)
{
std::function<
void
()> task;
{
std::unique_lock<std::mutex>
lock
(
this
->
queue_mutex
);
this
->
condition
.
wait
(lock,
[
this
] {
return
this
->
stop
|| !
this
->
tasks
.
empty
(); });
if
(
this
->
stop
&&
this
->
tasks
.
empty
())
return
;
task =
std::move
(
this
->
tasks
.
front
());
this
->
tasks
.
pop
();
}
task
();
}
});
size_t
thread_count =
0
;
for
(std::thread& worker : workers)
{
id_map[worker.
get_id
()] = thread_count;
thread_count++;
}
}
//
add new work item to the pool
template
<
class
F
,
class
... Args>
auto
ThreadPool::enqueue
(F&& f, Args&&... args)
-> std::future<typename std::invoke_result<F, Args...>::type>
{
assert
(num_threads_ >
1
);
using
return_type =
typename
std::invoke_result<F, Args...>::type;
auto
task = std::make_shared<std::packaged_task<
return_type
()>>(
std::bind
(std::forward<F>(f), std::forward<Args>(args)...));
std::future<return_type> res = task->
get_future
();
{
std::unique_lock<std::mutex>
lock
(queue_mutex);
if
(stop)
throw
std::runtime_error
(
"
enqueue on stopped ThreadPool
"
);
tasks.
emplace
([task]() { (*task)(); });
}
condition.
notify_one
();
return
res;
}
inline
ThreadPool::~ThreadPool
()
{
{
std::unique_lock<std::mutex>
lock
(queue_mutex);
stop =
true
;
}
condition.
notify_all
();
for
(std::thread& worker : workers)
worker.
join
();
}
Back
|
FazBrowse Home
|
New Git URL