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

added parallel_for · tornadocean/cpp-taskflow@548ad3c · GitHub

Commit 548ad3c

Browse files
added parallel_for
1 parent 1846710 commit 548ad3c

4 files changed

Lines changed: 126 additions & 28 deletions

File tree

‎CMakeLists.txt‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,10 @@ add_executable(matrix example/matrix.cpp)
6464
target_compile_options(matrix PRIVATE ${EXAMPLE_CXX_FLAGS})
6565
target_link_libraries(matrix ${EXAMPLE_EXE_LINKER_FLAGS})
6666

67+
add_executable(parallel_for example/parallel_for.cpp)
68+
target_compile_options(parallel_for PRIVATE ${EXAMPLE_CXX_FLAGS})
69+
target_link_libraries(parallel_for ${EXAMPLE_EXE_LINKER_FLAGS})
70+
6771

6872
# -----------------------------------------------------------------------------
6973
# Unittest

‎example/matrix.cpp‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -245,12 +245,12 @@ void taskflow_parallel(const std::vector<size_t>& D) {
245245
int main(int argc, char* argv[]) {
246246

247247
if(argc != 3) {
248-
std::cerr << "usage: ./matrix N [baseline|openmp|cppthread|taskflow]\n";
248+
std::cerr << "usage: ./matrix [baseline|openmp|cppthread|taskflow] N\n";
249249
std::exit(EXIT_FAILURE);
250250
}
251251

252252
// Create a unbalanced dimension for vector products.
253-
const auto N = std::stoul(argv[1]);
253+
const auto N = std::stoul(argv[2]);
254254

255255
std::vector<size_t> dimensions(N);
256256

@@ -266,7 +266,7 @@ int main(int argc, char* argv[]) {
266266
std::cout << "]\n";
267267

268268
// Run methods
269-
if(std::string_view method(argv[2]); method == "baseline") {
269+
if(std::string_view method(argv[1]); method == "baseline") {
270270
baseline(dimensions);
271271
}
272272
else if(method == "openmp") {

‎example/parallel_for.cpp‎

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
#include "taskflow.hpp"
2+
3+
// Function: fib
4+
int fib(int n) {
5+
if(n <= 2) return n;
6+
return (fib(n-1) + fib(n-2))%1024;
7+
}
8+
9+
// Procedure: sequential
10+
void sequential(int N) {
11+
auto tbeg = std::chrono::steady_clock::now();
12+
for(int i=0; i<N; ++i) {
13+
printf("fib[%d]=%d\n", i, fib(i));
14+
}
15+
auto tend = std::chrono::steady_clock::now();
16+
std::cout << "sequential version takes "
17+
<< std::chrono::duration_cast<std::chrono::milliseconds>(tend-tbeg).count()
18+
<< " ms\n";
19+
}
20+
21+
// Procedure: taskflow
22+
void taskflow(int N) {
23+
auto tbeg = std::chrono::steady_clock::now();
24+
tf::Taskflow tf;
25+
std::vector<int> range(N);
26+
for(int n=0; n<N; ++n) {
27+
range[n] = n;
28+
}
29+
tf.parallel_for(range.begin(), range.end(), [&] (int i) {
30+
printf("fib[%d]=%d\n", i, fib(i));
31+
});
32+
tf.wait_for_all();
33+
auto tend = std::chrono::steady_clock::now();
34+
std::cout << "taskflow version takes "
35+
<< std::chrono::duration_cast<std::chrono::milliseconds>(tend-tbeg).count()
36+
<< " ms\n";
37+
}
38+
39+
// Procedure: openmp
40+
void openmp(int N) {
41+
auto tbeg = std::chrono::steady_clock::now();
42+
#pragma omp parallel for
43+
for(int i=0; i<N; ++i) {
44+
printf("fib[%d]=%d\n", i, fib(i));
45+
}
46+
auto tend = std::chrono::steady_clock::now();
47+
std::cout << "taskflow version takes "
48+
<< std::chrono::duration_cast<std::chrono::milliseconds>(tend-tbeg).count()
49+
<< " ms\n";
50+
}
51+
52+
// ------------------------------------------------------------------------------------------------
53+
54+
// Function: main
55+
int main(int argc, char* argv[]) {
56+
57+
if(argc != 3) {
58+
std::cerr << "usage: ./parallel_for [baseline|openmp|taskflow] N\n";
59+
std::exit(EXIT_FAILURE);
60+
}
61+
62+
// Run methods
63+
if(std::string_view method(argv[1]); method == "baseline") {
64+
sequential(std::atoi(argv[2]));
65+
}
66+
else if(method == "openmp") {
67+
openmp(std::atoi(argv[2]));
68+
}
69+
else if(method == "taskflow") {
70+
taskflow(std::atoi(argv[2]));
71+
}
72+
else {
73+
std::cerr << "wrong method, shoud be [baseline|openmp|taskflow]\n";
74+
}
75+
76+
return 0;
77+
}

‎taskflow.hpp‎

Lines changed: 42 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -370,6 +370,7 @@ class BasicTaskflow {
370370
};
371371

372372

373+
BasicTaskflow();
373374
BasicTaskflow(unsigned);
374375
~BasicTaskflow();
375376

@@ -388,15 +389,17 @@ class BasicTaskflow {
388389
auto placeholder();
389390
auto dispatch();
390391
auto silent_dispatch();
391-
392-
BasicTaskflow& precede(Task, Task);
393-
BasicTaskflow& linearize(std::vector<Task>&);
394-
BasicTaskflow& linearize(std::initializer_list<Task>);
395-
BasicTaskflow& broadcast(Task, std::vector<Task>&);
396-
BasicTaskflow& broadcast(Task, std::initializer_list<Task>);
397-
BasicTaskflow& gather(std::vector<Task>&, Task);
398-
BasicTaskflow& gather(std::initializer_list<Task>, Task);
399-
BasicTaskflow& wait_for_all();
392+
auto precede(Task, Task);
393+
auto linearize(std::vector<Task>&);
394+
auto linearize(std::initializer_list<Task>);
395+
auto broadcast(Task, std::vector<Task>&);
396+
auto broadcast(Task, std::initializer_list<Task>);
397+
auto gather(std::vector<Task>&, Task);
398+
auto gather(std::initializer_list<Task>, Task);
399+
auto wait_for_all();
400+
401+
template<typename I, class C>
402+
auto parallel_for(I, I, C&&);
400403

401404
size_t num_nodes() const;
402405
size_t num_workers() const;
@@ -608,6 +611,11 @@ BasicTaskflow<F>::Task::Task(Node* t) : _node {t} {
608611
// BasicTaskflow
609612
//---------------------------------------------------------
610613

614+
// Constructor
615+
template <typename F>
616+
BasicTaskflow<F>::BasicTaskflow() : _threadpool{std::thread::hardware_concurrency()} {
617+
}
618+
611619
// Constructor
612620
template <typename F>
613621
BasicTaskflow<F>::BasicTaskflow(unsigned N) : _threadpool{N} {
@@ -639,9 +647,8 @@ size_t BasicTaskflow<F>::num_topologies() const {
639647

640648
// Procedure: precede
641649
template <typename F>
642-
BasicTaskflow<F>& BasicTaskflow<F>::precede(Task from, Task to) {
650+
auto BasicTaskflow<F>::precede(Task from, Task to) {
643651
from._node->precede(*(to._node));
644-
return *this;
645652
}
646653

647654
// Procedure: _linearize
@@ -659,44 +666,38 @@ void BasicTaskflow<F>::_linearize(L& keys) {
659666

660667
// Procedure: linearize
661668
template <typename F>
662-
BasicTaskflow<F>& BasicTaskflow<F>::linearize(std::vector<Task>& keys) {
669+
auto BasicTaskflow<F>::linearize(std::vector<Task>& keys) {
663670
_linearize(keys);
664-
return *this;
665671
}
666672

667673
// Procedure: linearize
668674
template <typename F>
669-
BasicTaskflow<F>& BasicTaskflow<F>::linearize(std::initializer_list<Task> keys) {
675+
auto BasicTaskflow<F>::linearize(std::initializer_list<Task> keys) {
670676
_linearize(keys);
671-
return *this;
672677
}
673678

674679
// Procedure: broadcast
675680
template <typename F>
676-
BasicTaskflow<F>& BasicTaskflow<F>::broadcast(Task from, std::vector<Task>& keys) {
681+
auto BasicTaskflow<F>::broadcast(Task from, std::vector<Task>& keys) {
677682
from.broadcast(keys);
678-
return *this;
679683
}
680684

681685
// Procedure: broadcast
682686
template <typename F>
683-
BasicTaskflow<F>& BasicTaskflow<F>::broadcast(Task from, std::initializer_list<Task> keys) {
687+
auto BasicTaskflow<F>::broadcast(Task from, std::initializer_list<Task> keys) {
684688
from.broadcast(keys);
685-
return *this;
686689
}
687690

688691
// Function: gather
689692
template <typename F>
690-
BasicTaskflow<F>& BasicTaskflow<F>::gather(std::vector<Task>& keys, Task to) {
693+
auto BasicTaskflow<F>::gather(std::vector<Task>& keys, Task to) {
691694
to.gather(keys);
692-
return *this;
693695
}
694696

695697
// Function: gather
696698
template <typename F>
697-
BasicTaskflow<F>& BasicTaskflow<F>::gather(std::initializer_list<Task> keys, Task to) {
699+
auto BasicTaskflow<F>::gather(std::initializer_list<Task> keys, Task to) {
698700
to.gather(keys);
699-
return *this;
700701
}
701702

702703
// Procedure: silent_dispatch
@@ -737,12 +738,11 @@ auto BasicTaskflow<F>::dispatch() {
737738

738739
// Procedure: wait_for_all
739740
template <typename F>
740-
BasicTaskflow<F>& BasicTaskflow<F>::wait_for_all() {
741+
auto BasicTaskflow<F>::wait_for_all() {
741742
if(!_nodes.empty()) {
742743
silent_dispatch();
743744
}
744745
_wait_for_topologies();
745-
return *this;
746746
}
747747

748748
// Procedure: _wait_for_topologies
@@ -806,6 +806,23 @@ auto BasicTaskflow<F>::emplace(C&&... cs) {
806806
return std::make_tuple(emplace(std::forward<C>(cs))...);
807807
}
808808

809+
// Function: parallel_for
810+
template <typename F>
811+
template <typename I, class C>
812+
auto BasicTaskflow<F>::parallel_for(I beg, I end, C&& c) {
813+
814+
auto source = placeholder();
815+
auto target = placeholder();
816+
817+
for(; beg != end; ++beg) {
818+
auto task = silent_emplace([&, itr=beg](){ c(*itr); });
819+
source.precede(task);
820+
task.precede(target);
821+
}
822+
823+
return std::make_pair(source, target);
824+
}
825+
809826
// Procedure: _schedule
810827
template <typename F>
811828
void BasicTaskflow<F>::_schedule(Node& task) {

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL