| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 5da1379 commit 7bfd1ba
9 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -614,6 +614,18 @@ add_test(prg.9threads ${TF_UTEST_ALGORITHM} -tc=prg.9threads) | |||
| 614 | 614 | add_test(prg.10threads ${TF_UTEST_ALGORITHM} -tc=prg.10threads) | |
| 615 | 615 | add_test(prg.11threads ${TF_UTEST_ALGORITHM} -tc=prg.11threads) | |
| 616 | 616 | 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) | ||
| 617 | 629 | ||
| 618 | 630 | # unittest for traverse | |
| 619 | 631 | add_executable(traverse ${TF_UTEST_DIR}/traverse.cpp) | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -88,10 +88,10 @@ void transform_reduce() { | |||
| 88 | 88 | auto tbeg = std::chrono::steady_clock::now(); | |
| 89 | 89 | tf::Taskflow tf; | |
| 90 | 90 | 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 | + //); | ||
| 95 | 95 | tf::Executor().run(tf).get(); | |
| 96 | 96 | auto tend = std::chrono::steady_clock::now(); | |
| 97 | 97 | std::cout << "[taskflow] transform_reduce " | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -108,7 +108,7 @@ Task FlowBuilder::parallel_for_guided(B&& beg, E&& end, C&& c, H&& chunk_size){ | |||
| 108 | 108 | return; | |
| 109 | 109 | } | |
| 110 | 110 | ||
| 111 | - size_t chunk_size = (h == 0) ? 1 : chunk_size; | ||
| 111 | + size_t chunk_size = (h == 0) ? 1 : h; | ||
| 112 | 112 | size_t W = sf._executor.num_workers(); | |
| 113 | 113 | size_t N = std::distance(beg, end); | |
| 114 | 114 | ||
@@ -206,7 +206,7 @@ Task FlowBuilder::parallel_for_guided( | |||
| 206 | 206 | TF_THROW("invalid range [", beg, ", ", end, ") with step size ", inc); | |
| 207 | 207 | } | |
| 208 | 208 | ||
| 209 | - size_t chunk_size = (h == 0) ? 1 : chunk_size; | ||
| 209 | + size_t chunk_size = (h == 0) ? 1 : h; | ||
| 210 | 210 | size_t W = sf._executor.num_workers(); | |
| 211 | 211 | size_t N = distance(beg, end, inc); | |
| 212 | 212 | ||
@@ -358,7 +358,9 @@ Task FlowBuilder::parallel_for_factoring(B&& beg, E&& end, C&& c){ | |||
| 358 | 358 | ||
| 359 | 359 | // Function: parallel_for_factoring | |
| 360 | 360 | 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 | + ){ | ||
| 362 | 364 | ||
| 363 | 365 | using I = underlying_index_t<B, E, S>; | |
| 364 | 366 | using namespace std::string_literals; | |
@@ -436,7 +438,9 @@ Task FlowBuilder::parallel_for_factoring(B&& beg, E&& end, S&& inc, C&& c){ | |||
| 436 | 438 | ||
| 437 | 439 | // Function: parallel_for_dynamic | |
| 438 | 440 | 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 | + ) { | ||
| 440 | 444 | ||
| 441 | 445 | using I = underlying_iterator_t<B, E>; | |
| 442 | 446 | using namespace std::string_literals; | |
@@ -445,7 +449,7 @@ Task FlowBuilder::parallel_for_dynamic(B&& beg, E&& end, C&& c, H&& chunk_size){ | |||
| 445 | 449 | [b=std::forward<B>(beg), | |
| 446 | 450 | e=std::forward<E>(end), | |
| 447 | 451 | 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 { | ||
| 449 | 453 | ||
| 450 | 454 | I beg = b; | |
| 451 | 455 | I end = e; | |
@@ -524,7 +528,7 @@ Task FlowBuilder::parallel_for_dynamic( | |||
| 524 | 528 | TF_THROW("invalid range [", beg, ", ", end, ") with step size ", inc); | |
| 525 | 529 | } | |
| 526 | 530 | ||
| 527 | - size_t chunk_size = (h == 0) ? 1 : chunk_size; | ||
| 531 | + size_t chunk_size = (h == 0) ? 1 : h; | ||
| 528 | 532 | size_t W = sf._executor.num_workers(); | |
| 529 | 533 | size_t N = distance(beg, end, inc); | |
| 530 | 534 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -9,7 +9,12 @@ namespace tf { | |||
| 9 | 9 | // ---------------------------------------------------------------------------- | |
| 10 | 10 | ||
| 11 | 11 | 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 | + ) { | ||
| 13 | 18 | return parallel_reduce_guided( | |
| 14 | 19 | std::forward<B>(beg), | |
| 15 | 20 | std::forward<E>(end), | |
@@ -147,6 +152,108 @@ Task FlowBuilder::parallel_reduce_guided( | |||
| 147 | 152 | return task; | |
| 148 | 153 | } | |
| 149 | 154 | ||
| 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 | + //) | ||
| 150 | 257 | ||
| 151 | 258 | } // end of namespace tf ----------------------------------------------------- | |
| 152 | 259 | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments