<li><ahref="#ParallelDataPipelineIncludeHeaderFile">Include the Header</a></li>
<li><ahref="#CreateADataPipelineModuleTask">Create a Data Pipeline Module Task</a></li>
<li><ahref="#UnderstandInternalDataStorage">Understand Internal Data Storage</a></li>
<li><ahref="#DataParallelPipelineLearnMore">Learn More about Taskflow Pipeline</a></li>
</ul>
</nav>
<p>Taskflow provides another variant, <ahref="classtf_1_1DataPipeline.html" class="m-doc">tf::<wbr/>DataPipeline</a>, on top of <ahref="classtf_1_1Pipeline.html" class="m-doc">tf::<wbr/>Pipeline</a> (see <ahref="TaskParallelPipeline.html" class="m-doc">Task-parallel Pipeline</a>) to help you implement data-parallel pipeline algorithms while leaving data management to Taskflow. We recommend you finishing reading TaskParallelPipeline first before learning <ahref="classtf_1_1DataPipeline.html" class="m-doc">tf::<wbr/>DataPipeline</a>.</p><sectionid="ParallelDataPipelineIncludeHeaderFile"><h2><ahref="#ParallelDataPipelineIncludeHeaderFile">Include the Header</a></h2><p>You need to include the header file, <code>taskflow/algorithm/data_pipeline.hpp</code>, for implementing data-parallel pipeline algorithms.</p><preclass="m-code"><spanclass="cp">#include</span><spanclass="w"></span><spanclass="cpf"><taskflow/algorithm/data_pipeline.hpp></span></pre></section><sectionid="CreateADataPipelineModuleTask"><h2><ahref="#CreateADataPipelineModuleTask">Create a Data Pipeline Module Task</a></h2><p>Similar to creating a task-parallel pipeline (<ahref="classtf_1_1Pipeline.html" class="m-doc">tf::<wbr/>Pipeline</a>), there are three steps to create a data-parallel pipeline application:</p><ol><li>Define the pipeline structure (e.g., pipe type, pipe callable, stopping rule, line count)</li><li>Define the data storage and layout, if needed for the application</li><li>Define the pipeline taskflow graph using composition</li></ol><p>The following example creates a data-parallel pipeline that generates a total of five dataflow tokens from <code>void</code> to <code>int</code> at the first stage, from <code>int</code> to <code>std::string</code> at the second stage, and <code>std::string</code> to <code>void</code> at the final stage. Data storage between stages is automatically managed by <ahref="classtf_1_1DataPipeline.html" class="m-doc">tf::<wbr/>DataPipeline</a>.</p><preclass="m-code"><spanclass="cp">#include</span><spanclass="w"></span><spanclass="cpf"><taskflow/taskflow.hpp></span>
<spanclass="w"></span><spanclass="n">printf</span><spanclass="p">(</span><spanclass="s">"second pipe returns a string of %d</span><spanclass="se">\n</span><spanclass="s">"</span><spanclass="p">,</span><spanclass="w"></span><spanclass="n">input</span><spanclass="w"></span><spanclass="o">+</span><spanclass="w"></span><spanclass="mi">100</span><spanclass="p">);</span>
<spanclass="p">}</span></pre><p>The interface of <ahref="classtf_1_1DataPipeline.html" class="m-doc">tf::<wbr/>DataPipeline</a> is very similar to <ahref="classtf_1_1Pipeline.html" class="m-doc">tf::<wbr/>Pipeline</a>, except that the library transparently manages the dataflow between pipes. To create a stage in a data-parallel pipeline, you should always use the helper function <ahref="namespacetf.html#a8975fa5762088789adb0b60f38208309" class="m-doc">tf::<wbr/>make_data_pipe</a>:</p><preclass="m-code"><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">make_data_pipe</span><spanclass="o"><</span><spanclass="kt">int</span><spanclass="p">,</span><spanclass="w"></span><spanclass="n">std</span><spanclass="o">::</span><spanclass="n">string</span><spanclass="o">></span><spanclass="p">(</span>
<spanclass="p">);</span></pre><p>The helper function starts with a pair of an input and an output types in its template arguments. Both types will always be decayed to their original form using <ahref="http://en.cppreference.com/w/cpp/types/decay.html" class="m-doc-external">std::<wbr/>decay</a> (e.g., <code>const int&</code> becomes <code>int</code>) for storage purpose. In terms of function arguments, the first argument specifies the direction of this data pipe, which can be either <ahref="namespacetf.html#abb7a11e41fd457f69e7ff45d4c769564a7b804a28d6154ab8007287532037f1d0" class="m-doc">tf::<wbr/>PipeType::<wbr/>SERIAL</a> or <ahref="namespacetf.html#abb7a11e41fd457f69e7ff45d4c769564adf13a99b035d6f0bce4f44ab18eec8eb" class="m-doc">tf::<wbr/>PipeType::<wbr/>PARALLEL</a>, and the second argument is a callable to invoke by the pipeline scheduler. The callable must take the input data type in its first argument and returns a value of the output data type. Additionally, the callable can take a <ahref="classtf_1_1Pipeflow.html" class="m-doc">tf::<wbr/>Pipeflow</a> reference in its second argument which allows you to query the runtime information of a stage task, such as its line number and token number.</p><preclass="m-code"><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">make_data_pipe</span><spanclass="o"><</span><spanclass="kt">int</span><spanclass="p">,</span><spanclass="w"></span><spanclass="n">std</span><spanclass="o">::</span><spanclass="n">string</span><spanclass="o">></span><spanclass="p">(</span>
<spanclass="p">)</span></pre><asideclass="m-note m-info"><h4>Note</h4><p>By default, <ahref="classtf_1_1DataPipeline.html" class="m-doc">tf::<wbr/>DataPipeline</a> passes the data in reference to your callable at which you can take it in copy or in reference depending on application needs.</p></aside><p>For the first pipe, the input type should always be <code>void</code> and the callable must take a <ahref="classtf_1_1Pipeflow.html" class="m-doc">tf::<wbr/>Pipeflow</a> reference in its argument. In this example, we will stop the pipeline when processing five tokens.</p><preclass="m-code"><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">make_data_pipe</span><spanclass="o"><</span><spanclass="kt">void</span><spanclass="p">,</span><spanclass="w"></span><spanclass="kt">int</span><spanclass="o">></span><spanclass="p">(</span><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">PipeType</span><spanclass="o">::</span><spanclass="n">SERIAL</span><spanclass="p">,</span><spanclass="w"></span><spanclass="p">[](</span><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">Pipeflow</span><spanclass="o">&</span><spanclass="w"></span><spanclass="n">pf</span><spanclass="p">)</span><spanclass="w"></span><spanclass="o">-></span><spanclass="w"></span><spanclass="kt">int</span><spanclass="p">{</span>
<spanclass="w"></span><spanclass="k">return</span><spanclass="w"></span><spanclass="mi">0</span><spanclass="p">;</span><spanclass="w"></span><spanclass="c1">// returns a dummy value</span>
<spanclass="p">}),</span></pre><p>Similarly, the output type of the last pipe should be <code>void</code> as no more data will go out of the final pipe.</p><preclass="m-code"><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">make_data_pipe</span><spanclass="o"><</span><spanclass="n">std</span><spanclass="o">::</span><spanclass="n">string</span><spanclass="p">,</span><spanclass="w"></span><spanclass="kt">void</span><spanclass="o">></span><spanclass="p">(</span><spanclass="n">tf</span><spanclass="o">::</span><spanclass="n">PipeType</span><spanclass="o">::</span><spanclass="n">SERIAL</span><spanclass="p">,</span><spanclass="w"></span><spanclass="p">[](</span><spanclass="n">std</span><spanclass="o">::</span><spanclass="n">string</span><spanclass="o">&</span><spanclass="w"></span><spanclass="n">input</span><spanclass="p">)</span><spanclass="w"></span><spanclass="p">{</span>
<spanclass="p">})</span></pre><p>Finally, you need to compose the pipeline graph by creating a module task (i.e., tf::Taskflow::compoased_of).</p><preclass="m-code"><spanclass="c1">// build the pipeline graph using composition</span>
</div></section><sectionid="UnderstandInternalDataStorage"><h2><ahref="#UnderstandInternalDataStorage">Understand Internal Data Storage</a></h2><p>By default, <ahref="classtf_1_1DataPipeline.html" class="m-doc">tf::<wbr/>DataPipeline</a> uses <ahref="https://en.cppreference.com/w/cpp/utility/variant">std::<wbr/>variant</a> to store a type-safe union of all input and output data types extracted from the given data pipes. To avoid false sharing, each line keeps a variant that is aligned with the cacheline size. When invoking a pipe callable, the input data is acquired in reference from the variant using <ahref="https://en.cppreference.com/w/cpp/utility/variant/get">std::<wbr/>get</a>. When returning from a pipe callable, the output data is stored back to the variant using assignment operator.</p></section><sectionid="DataParallelPipelineLearnMore"><h2><ahref="#DataParallelPipelineLearnMore">Learn More about Taskflow Pipeline</a></h2><p>Visit the following pages to learn more about pipeline:</p><ol><li><ahref="TaskParallelPipeline.html" class="m-doc">Task-parallel Pipeline</a></li><li><ahref="TaskParallelScalablePipeline.html" class="m-doc">Task-parallel Scalable Pipeline</a></li><li><ahref="TextProcessingPipeline.html" class="m-doc">Text Processing Pipeline</a></li><li><ahref="GraphProcessingPipeline.html" class="m-doc">Graph Processing Pipeline</a></li><li><ahref="TaskflowProcessingPipeline.html" class="m-doc">Taskflow Processing Pipeline</a></li></ol></section>
</div>
</div>
</div>
</article></main>
<divclass="m-doc-search" id="search">
<ahref="#!" onclick="return hideSearch()"></a>
<divclass="m-container">
<divclass="m-row">
<divclass="m-col-m-8 m-push-m-2">
<divclass="m-doc-search-header m-text m-small">
<div><spanclass="m-label m-default">Tab</span> / <spanclass="m-label m-default">T</span> to search, <spanclass="m-label m-default">Esc</span> to close</div>