FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

updated parallel-for · baileyfu/taskflow-cpp@7bfd1ba · GitHub

Commit 7bfd1ba

Browse files
committed
updated parallel-for
1 parent 5da1379 commit 7bfd1ba

9 files changed

Lines changed: 263 additions & 338 deletions

File tree

‎CMakeLists.txt‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -614,6 +614,18 @@ add_test(prg.9threads ${TF_UTEST_ALGORITHM} -tc=prg.9threads)
614614
add_test(prg.10threads ${TF_UTEST_ALGORITHM} -tc=prg.10threads)
615615
add_test(prg.11threads ${TF_UTEST_ALGORITHM} -tc=prg.11threads)
616616
add_test(prg.12threads ${TF_UTEST_ALGORITHM} -tc=prg.12threads)
617+
add_test(prd.1thread ${TF_UTEST_ALGORITHM} -tc=prd.1thread)
618+
add_test(prd.2threads ${TF_UTEST_ALGORITHM} -tc=prd.2threads)
619+
add_test(prd.3threads ${TF_UTEST_ALGORITHM} -tc=prd.3threads)
620+
add_test(prd.4threads ${TF_UTEST_ALGORITHM} -tc=prd.4threads)
621+
add_test(prd.5threads ${TF_UTEST_ALGORITHM} -tc=prd.5threads)
622+
add_test(prd.6threads ${TF_UTEST_ALGORITHM} -tc=prd.6threads)
623+
add_test(prd.7threads ${TF_UTEST_ALGORITHM} -tc=prd.7threads)
624+
add_test(prd.8threads ${TF_UTEST_ALGORITHM} -tc=prd.8threads)
625+
add_test(prd.9threads ${TF_UTEST_ALGORITHM} -tc=prd.9threads)
626+
add_test(prd.10threads ${TF_UTEST_ALGORITHM} -tc=prd.10threads)
627+
add_test(prd.11threads ${TF_UTEST_ALGORITHM} -tc=prd.11threads)
628+
add_test(prd.12threads ${TF_UTEST_ALGORITHM} -tc=prd.12threads)
617629

618630
# unittest for traverse
619631
add_executable(traverse ${TF_UTEST_DIR}/traverse.cpp)

‎examples/reduce.cpp‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,10 +88,10 @@ void transform_reduce() {
8888
auto tbeg = std::chrono::steady_clock::now();
8989
tf::Taskflow tf;
9090
auto tmin = std::numeric_limits<int>::max();
91-
tf.transform_reduce(data.begin(), data.end(), tmin,
92-
[] (int l, int r) { return std::min(l, r); },
93-
[] (const Data& d) { return d.transform(); }
94-
);
91+
//tf.transform_reduce(data.begin(), data.end(), tmin,
92+
// [] (int l, int r) { return std::min(l, r); },
93+
// [] (const Data& d) { return d.transform(); }
94+
//);
9595
tf::Executor().run(tf).get();
9696
auto tend = std::chrono::steady_clock::now();
9797
std::cout << "[taskflow] transform_reduce "

‎taskflow/algorithm/parallel_for.hpp‎

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -108,7 +108,7 @@ Task FlowBuilder::parallel_for_guided(B&& beg, E&& end, C&& c, H&& chunk_size){
108108
return;
109109
}
110110

111-
size_t chunk_size = (h == 0) ? 1 : chunk_size;
111+
size_t chunk_size = (h == 0) ? 1 : h;
112112
size_t W = sf._executor.num_workers();
113113
size_t N = std::distance(beg, end);
114114

@@ -206,7 +206,7 @@ Task FlowBuilder::parallel_for_guided(
206206
TF_THROW("invalid range [", beg, ", ", end, ") with step size ", inc);
207207
}
208208

209-
size_t chunk_size = (h == 0) ? 1 : chunk_size;
209+
size_t chunk_size = (h == 0) ? 1 : h;
210210
size_t W = sf._executor.num_workers();
211211
size_t N = distance(beg, end, inc);
212212

@@ -358,7 +358,9 @@ Task FlowBuilder::parallel_for_factoring(B&& beg, E&& end, C&& c){
358358

359359
// Function: parallel_for_factoring
360360
template <typename B, typename E, typename S, typename C>
361-
Task FlowBuilder::parallel_for_factoring(B&& beg, E&& end, S&& inc, C&& c){
361+
Task FlowBuilder::parallel_for_factoring(
362+
B&& beg, E&& end, S&& inc, C&& c
363+
){
362364

363365
using I = underlying_index_t<B, E, S>;
364366
using namespace std::string_literals;
@@ -436,7 +438,9 @@ Task FlowBuilder::parallel_for_factoring(B&& beg, E&& end, S&& inc, C&& c){
436438

437439
// Function: parallel_for_dynamic
438440
template <typename B, typename E, typename C, typename H>
439-
Task FlowBuilder::parallel_for_dynamic(B&& beg, E&& end, C&& c, H&& chunk_size){
441+
Task FlowBuilder::parallel_for_dynamic(
442+
B&& beg, E&& end, C&& c, H&& chunk_size
443+
) {
440444

441445
using I = underlying_iterator_t<B, E>;
442446
using namespace std::string_literals;
@@ -445,7 +449,7 @@ Task FlowBuilder::parallel_for_dynamic(B&& beg, E&& end, C&& c, H&& chunk_size){
445449
[b=std::forward<B>(beg),
446450
e=std::forward<E>(end),
447451
c=std::forward<C>(c),
448-
h=std::forward<H>(chunk_size)](Subflow& sf) mutable {
452+
h=std::forward<H>(chunk_size)] (Subflow& sf) mutable {
449453

450454
I beg = b;
451455
I end = e;
@@ -524,7 +528,7 @@ Task FlowBuilder::parallel_for_dynamic(
524528
TF_THROW("invalid range [", beg, ", ", end, ") with step size ", inc);
525529
}
526530

527-
size_t chunk_size = (h == 0) ? 1 : chunk_size;
531+
size_t chunk_size = (h == 0) ? 1 : h;
528532
size_t W = sf._executor.num_workers();
529533
size_t N = distance(beg, end, inc);
530534

‎taskflow/algorithm/reduce.hpp‎

Lines changed: 108 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,12 @@ namespace tf {
99
// ----------------------------------------------------------------------------
1010

1111
template <typename B, typename E, typename T, typename O>
12-
Task FlowBuilder::parallel_reduce(B&& beg, E&& end, T& init, O&& bop) {
12+
Task FlowBuilder::parallel_reduce(
13+
B&& beg,
14+
E&& end,
15+
T& init,
16+
O&& bop
17+
) {
1318
return parallel_reduce_guided(
1419
std::forward<B>(beg),
1520
std::forward<E>(end),
@@ -147,6 +152,108 @@ Task FlowBuilder::parallel_reduce_guided(
147152
return task;
148153
}
149154

155+
// ----------------------------------------------------------------------------
156+
// parallel_reduce_dynamic
157+
// ----------------------------------------------------------------------------
158+
159+
template <typename B, typename E, typename T, typename O, typename H>
160+
Task FlowBuilder::parallel_reduce_dynamic(
161+
B&& beg, E&& end, T& init, O&& bop, H&& chunk_size
162+
) {
163+
164+
using I = underlying_iterator_t<B, E>;
165+
using namespace std::string_literals;
166+
167+
Task task = emplace(
168+
[b=std::forward<B>(beg),
169+
e=std::forward<E>(end),
170+
&r=init,
171+
o=std::forward<O>(bop),
172+
c=std::forward<H>(chunk_size)
173+
] (Subflow& sf) mutable {
174+
175+
// fetch the iterator values
176+
I beg = b;
177+
I end = e;
178+
179+
if(beg == end) {
180+
return;
181+
}
182+
183+
size_t C = (c == 0) ? 1 : c;
184+
size_t W = sf._executor.num_workers();
185+
size_t N = std::distance(beg, end);
186+
187+
// only myself - no need to spawn another graph
188+
if(W <= 1 || N <= C) {
189+
for(; beg!=end; r = o(r, *beg++));
190+
return;
191+
}
192+
193+
if(N < W) {
194+
W = N;
195+
}
196+
197+
std::mutex mutex;
198+
std::atomic<size_t> next(0);
199+
200+
for(size_t w=0; w<W; w++) {
201+
202+
if(w*2 >= N) {
203+
break;
204+
}
205+
206+
sf.emplace([&mutex, &next, &r, beg, N, &o, C] () mutable {
207+
208+
size_t s0 = next.fetch_add(2, std::memory_order_relaxed);
209+
210+
if(s0 >= N) {
211+
return;
212+
}
213+
214+
std::advance(beg, s0);
215+
216+
if(N - s0 == 1) {
217+
std::lock_guard<std::mutex> lock(mutex);
218+
r = o(r, *beg);
219+
return;
220+
}
221+
222+
auto beg1 = beg++;
223+
auto beg2 = beg++;
224+
225+
T sum = o(*beg1, *beg2);
226+
227+
size_t z = s0 + 2;
228+
229+
while(1) {
230+
s0 = next.fetch_add(C, std::memory_order_relaxed);
231+
if(s0 >= N) {
232+
break;
233+
}
234+
size_t e0 = (C <= (N - s0)) ? s0 + C : N;
235+
std::advance(beg, s0-z);
236+
for(size_t x=s0; x<e0; x++, beg++) {
237+
sum = o(sum, *beg);
238+
}
239+
z = e0;
240+
}
241+
242+
std::lock_guard<std::mutex> lock(mutex);
243+
r = o(r, sum);
244+
245+
}).name("prd_"s + std::to_string(w));
246+
}
247+
248+
sf.join();
249+
});
250+
251+
return task;
252+
}
253+
254+
//TF_ALGORITHM_GUIDED_PARTITION(
255+
// sum = o(sum, *beg);
256+
//)
150257

151258
} // end of namespace tf -----------------------------------------------------
152259

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL