You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.
Dismiss alert
<li><ahref="#TaskflowPipelineDefineThePipes">Define the Pipes</a></li>
<li><ahref="#TaskflowPipelineDefineTheTaskGraph">Define the Task Graph</a></li>
<li><ahref="#TaskflowPipelineSubmitTheTaskGraph">Submit the Task Graph</a></li>
</ul>
</li>
</ul>
</nav>
<p>We study a taskflow processing pipeline that propagates a sequence of tokens through linearly dependent taskflows. The pipeline embeds a taskflow in each pipe to run a parallel algorithm using task graph parallelism.</p><sectionid="FormulateTheTaskflowProcessingPipelineProblem"><h2><ahref="#FormulateTheTaskflowProcessingPipelineProblem">Formulate the Taskflow Processing Pipeline Problem</a></h2><p>Many complex and irregular pipeline applications require each pipe to run a parallel algorithm using task graph parallelism. We can formulate such applications as scheduling a sequence of tokens through linearly dependent taskflows. The following example illustrates the pipeline propagation of three scheduling tokens through three linearly dependent taskflows:</p><divclass="m-graph"><svgstyle="width: 36.500rem; height: 18.800rem;" viewBox="0.00 0.00 364.75 188.00">
</div><p>Each pipe (stage) in the pipeline embeds a taskflow to perform a stage-specific parallel algorithm on an input scheduling token. Parallelism exhibits both inside and outside the three taskflows, combining both <em>task graph parallelism</em> and <em>pipeline parallelism</em>.</p></section><sectionid="CreateATaskflowProcessingPipeline"><h2><ahref="#CreateATaskflowProcessingPipeline">Create a Taskflow Processing Pipeline</a></h2><p>Using the example from the previous section, we create a pipeline of three <em>serial</em> pipes each running a taskflow on a sequence of five scheduling tokens. The overall implementation is shown below:</p><preclass="m-code"><spanclass="cp">#include</span><spanclass="w"></span><spanclass="cpf"><taskflow/taskflow.hpp></span>
<spanclass="p">}</span></pre><sectionid="TaskflowPipelineDefineTaskflows"><h3><ahref="#TaskflowPipelineDefineTaskflows">Define Taskflows</a></h3><p>First, we define three taskflows for the three pipes in the pipeline:</p><preclass="m-code"><spanclass="c1">// taskflow on the first pipe</span>
<spanclass="p">}</span></pre><p>As each taskflow corresponds to a pipe in the pipeline, we create a linear array to store the three taskflows:</p><preclass="m-code"><spanclass="n">std</span><spanclass="o">::</span><spanclass="n">array</span><spanclass="o"><</span><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">Taskflow</span><spanclass="p">,</span><spanclass="w"></span><spanclass="n">num_pipes</span><spanclass="o">></span><spanclass="w"></span><spanclass="n">taskflows</span><spanclass="p">;</span>
<spanclass="n">make_taskflow3</span><spanclass="p">(</span><spanclass="n">taskflows</span><spanclass="p">[</span><spanclass="mi">2</span><spanclass="p">]);</span></pre><p>Since the three taskflows are linearly dependent, at most one taskflow will run at a pipe. We can store the three taskflows in a linear array of dimension equal to the number of pipes. If there is a parallel pipe, we need to use two-dimensional array, as multiple taskflows at a stage can run simultaneously across parallel lines.</p></section><sectionid="TaskflowPipelineDefineThePipes"><h3><ahref="#TaskflowPipelineDefineThePipes">Define the Pipes</a></h3><p>The pipe definition is straightforward. Each pipe runs the corresponding taskflow, which can be indexed at <code>taskflows</code> with the pipe's identifier, <ahref="classtf_1_1Pipeflow.html#a4914c1f381a3016e98285b019cf60d6d" class="m-doc">tf::<wbr/>Pipeflow::<wbr/>pipe()</a>. The first pipe will cease the pipeline scheduling when it has processed five scheduling tokens:</p><preclass="m-code"><spanclass="c1">// first pipe runs taskflow1</span>
<spanclass="p">}}</span></pre><p>At each pipe, we use <ahref="classtf_1_1Executor.html#a8fcd9e0557922bb8194999f0cd433ea8" class="m-doc">tf::<wbr/>Executor::<wbr/>corun</a> to execute the corresponding taskflow and wait until the execution completes. This is important because we want the caller thread, which is the worker that invokes the pipe callable, to not block (i.e., <code>executor.run(taskflows[pf.pipe()]).wait()</code>) but participate in the work-stealing loop of the scheduler to avoid deadlock.</p></section><sectionid="TaskflowPipelineDefineTheTaskGraph"><h3><ahref="#TaskflowPipelineDefineTheTaskGraph">Define the Task Graph</a></h3><p>To build up the taskflow for the pipeline, we create a module task with the defined pipeline structure and connect it with two tasks that output helper messages before and after the pipeline:</p><preclass="m-code"><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">Task</span><spanclass="w"></span><spanclass="n">init</span><spanclass="w"></span><spanclass="o">=</span><spanclass="w"></span><spanclass="n">taskflow</span><spanclass="p">.</span><spanclass="n">emplace</span><spanclass="p">([](){</span><spanclass="w"></span><spanclass="n">std</span><spanclass="o">::</span><spanclass="n">cout</span><spanclass="w"></span><spanclass="o"><<</span><spanclass="w"></span><spanclass="s">"ready</span><spanclass="se">\n</span><spanclass="s">"</span><spanclass="p">;</span><spanclass="w"></span><spanclass="p">})</span>
</div></section><sectionid="TaskflowPipelineSubmitTheTaskGraph"><h3><ahref="#TaskflowPipelineSubmitTheTaskGraph">Submit the Task Graph</a></h3><p>Finally, we submit the taskflow to the execution and run it once:</p><preclass="m-code"><spanclass="n">executor</span><spanclass="p">.</span><spanclass="n">run</span><spanclass="p">(</span><spanclass="n">taskflow</span><spanclass="p">).</span><spanclass="n">wait</span><spanclass="p">();</span></pre><p>One possible output is shown below:</p><preclass="m-console"><spanclass="go">ready</span>