ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

oneTBB 并行化复杂循环:parallel_for_each 与 parallel_pipeline 实战指南

oneTBB 并行化复杂循环:parallel_for_each 与 parallel_pipeline 实战指南 并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载复杂循环的并行化并不总是像parallel_for那样简单直接。当迭代空间在运行前无法预知、循环体可能动态增加新任务或数据需要像流水线一样经过多级处理时oneTBB 提供了两个专门的高层模板oneapi::tbb::parallel_for_each和oneapi::tbb::parallel_pipeline。本文以官方用户指南 Parallelizing_Complex_Loops.rst 为骨架结合仓库中的头文件实现与可运行示例完整讲解这两种模式的使用方法、函数对象约束、管线过滤器模式与性能调优原则帮助读者在简单的parallel_for模式之外掌握处理链表遍历、树形递归任务和分阶段数据流等场景的工程化方案。阅读本文后你将能够用parallel_for_each并行遍历未知长度的迭代空间并通过 feeder 动态补充任务用parallel_pipeline搭建读入—变换—写出的多级并行管线并正确选择 filter mode 与 token 数量理解管线的吞吐瓶颈、窗口大小与缓冲区复用的取舍。为什么需要复杂循环模式oneTBB 的 Parallelizing_Simple_Loops 一节介绍了最典型的可扩展并行形态——迭代之间彼此独立、可同时运行。但实际应用中有些循环并不满足这一前提迭代空间未知例如遍历一个链表循环结束时才知道总共要处理多少个元素循环体会追加工作例如处理树结构时处理完一个节点后才知道需要继续处理哪些后代节点数据需要按级流动例如读文件 → 变换 → 写文件的流水线不同阶段对数据的不同处理方式适合不同的并行度。上述三类场景分别由parallel_for_each和parallel_pipeline覆盖。官方文档对二者的定位是在简单的parallel_for之外为其他并行模式提供的备选方案。这一节的 toctree 还列出了配套的 Summary_of_Loops_and_Pipelines用于汇总高层模板的定位它们让开发者以任务模式task-pattern级别设计软件而无需关心底层线程的调度细节。编译提示本节文档沿用了旧版文档的编译命令表Windows 下icl /MD example.cpp tbb_debug.lib、Linux 下icc example.cpp -ltbb_debug去掉_debug即为生产版链接。在实际使用现代工具链时更推荐通过仓库根目录的 CMakeLists.txt 配置 CMake 构建并在链接阶段加入 oneTBB 库以避免未定义引用。模式一parallel_for_each——煮到熟为止的并行遍历适用场景与约束parallel_for_each解决的问题非常具体循环结束时迭代空间仍未完全确定或者循环体会不断往任务列表里添加新工作。官方文档给出的典型例子是链表链表的长度在遍历前不可知访问链表元素天然是串行的但如果每个元素的处理Foo至少需要几千条指令那么并行处理这些元素仍然值得。文档强调在并行编程中动态数组通常优于链表因为链表访问本质上串行仅当被迫使用链表且单元素处理开销足够大时才用parallel_for_each换取并行收益。这个每个 item 至少几千条指令的门槛是为了让并行调度的开销能够被实际计算摊薄。从串行到并行的三步改造文档以对链表每个元素执行Foo为例给出串行版本void SerialApplyFooToList( const std::listItem list ) { for( std::listItem::const_iterator ilist.begin(); i!list.end(); i ) Foo(*i); }第一步定义函数对象。parallel_for_each要求函数对象具备const限定的operator()这一点与标准库functional中的函数对象类似但约束更严格class ApplyFoo { public: void operator()( Item item ) const { Foo(item); } };第二步调用算法。把循环体替换为对parallel_for_each的一次调用void ParallelApplyFooToList( const std::listItem list ) { parallel_for_each( list.begin(), list.end(), ApplyFoo() ); }第三步可选理解可扩展性。文档明确点出一个关键事实parallel_for_each的一次调用绝不会让两个线程同时对同一个输入迭代器取值。这意味着顺序程序中的输入迭代器可以原样使用但工作的获取fetching of work本身是串行的因此算法并不天然可扩展。尽管如此在多数场景下仍能获得比纯串行更有用的加速。两种可扩展的取工作方式文档进一步给出两条让parallel_for_each获得可扩展性的途径随机访问迭代器当传入的迭代器是 random-access 类型时oneTBB 可以直接把区间切分给多个线程。这一行为在源码中有直接对应for_each_root_task针对 std::random_access_iterator_tag 有专门特化其实现是把区间委托给parallel_for配合parallel_for_body_wrapper处理而 forward/input 迭代器则走 input_block_handling_task 与 forward_block_handling_task按每块最多 4 个迭代任务max_block_size 4分批派发。从源码结构看这就是random-access 可扩展、其余迭代器逐块串行取值的实现依据。feeder 动态补充工作parallel_for_each的 body 如果接受第二个参数feeder类型为parallel_for_eachItem就可以在处理过程中调用feeder.add(item)追加新任务。文档以树遍历为例处理某个节点是处理其所有后代节点的前提那么处理完节点后用feeder.add把后代节点加进去即可。这次parallel_for_each调用不会提前结束直到所有已添加的 item 都被处理完。feeder 的实现细节同样可以在仓库中验证feeder 类 提供void add(const Item)与void add(Item)两个重载分别走internal_add_copy与internal_add_move底层feeder_impl会把新 item 包装成feeder_item_task并通过small_object_allocator分配、spawn到当前执行上下文见 feeder_impl是否创建 feeder 由编译期探测决定feeder_holder仅在 body 的operator()能接受feeder参数时才真正构造 feeder 对象见 feeder_holder因此带 feeder 与不带 feeder 的 body 可以混用同一算法接口函数对象的签名探测通过 parallel_for_each_operator_selector 的重载解析完成tbb::detail::invoke负责最终调用。仓库中的一致性测试 test_parallel_for_each.cpp 用static_assert验证了这些约束带 feeder 的 bodyWithFeederint合法而operator()非const的 bodyOperatorRoundBracketsNonConst会被编译期拒绝。使用注意点body 的operator()必须是const因为它可能被并发调用body 可以接受Item可修改元素或按值/常量引用处理元素具体取决于迭代器类型与 Item 的可移动性——源码中input_iteration_task_iterator_helper会对支持右值调用的 body 自动采用std::move_iterator传递元素见 include/oneapi/tbb/parallel_for_each.h#L240-L256若 Item 既不可拷贝也不可移动feeder.add会在断言中失败编译期特性探测会尽量提前发现该问题。模式二parallel_pipeline——装配线上的并行处理模式思想与三种过滤器模式流水线pipelining模仿传统制造业的装配线数据以 token 形式流经一系列过滤器filter每个过滤器对数据做某种处理。由于流式数据到达的顺序约束有些过滤器可以并行有些不能。例如视频处理中部分帧操作彼此独立可以同时进行而另一些帧操作必须先处理前面的帧。oneTBB 的parallel_pipeline与filter类实现了这一模式。filter_mode在 include/oneapi/tbb/parallel_pipeline.h#L38-L46 中定义共有三种取值其含义可结合 base_filter 的位标志 理解filter_mode位标志组合语义parallelfilter_is_out_of_order可并发处理多个 item且不保证顺序serial_in_orderfilter_is_serial一次处理一个 item所有该模式的过滤器按同一顺序处理serial_out_of_orderfilter_is_serial \| filter_is_out_of_order一次处理一个 item但不保证顺序官方示例文本格式化管线文档用读文本文件 → 把每个十进制数字平方 → 写回新文件的示例演示完整用法。原始文件 I/O 是顺序的因此输入输出过滤器必须是serial_in_order中间的平方变换只依赖局部数据可以指定为parallel让多个 chunk 同时变换再按顺序写出。为了摊薄并行调度的开销过滤器处理的是大约4000 字符的文本块chunk。每个 chunk 由TextSlice对象表示——注意该类在文档和 示例 square.cpp 中完全一致TextSlice的 C 声明只是内存中更大对象的头部实例必须通过类内allocate/free方法分配与释放即通过oneapi::tbb::tbb_allocatorchar申请sizeof(TextSlice)max_size1字节多出的 1 字节用于存放终止符。构建并运行管线的顶层代码来自文档与示例实现一致void RunPipeline( int ntoken, FILE* input_file, FILE* output_file ) { oneapi::tbb::parallel_pipeline( ntoken, oneapi::tbb::make_filtervoid,TextSlice*( oneapi::tbb::filter_mode::serial_in_order, MyInputFunc(input_file) ) oneapi::tbb::make_filterTextSlice*,TextSlice*( oneapi::tbb::filter_mode::parallel, MyTransformFunc() ) oneapi::tbb::make_filterTextSlice*,void( oneapi::tbb::filter_mode::serial_in_order, MyOutputFunc(output_file) ) ); }参数与接口约定继承自文档并补充实现细节ntoken第一个参数控制并行度。概念上 token 流经整个管线serial_in_order过滤器必须按序串行处理每个 tokenparallel过滤器可以让多个 token 并行处理。如果 token 无上限中间的无序过滤器可能因为输出过滤器跟不上而不断积累 token导致资源被中间过滤器过度消耗。ntoken指定同时在途in flight的 token 上限一旦达到该上限输入过滤器不再产生新 token直到输出过滤器销毁一个 token。该参数在 parallel_pipeline 声明 中对应max_number_of_live_tokens。第二个参数是过滤器链。每个过滤器由make_filterinputType, outputType(mode, functor)构造其中inputType过滤器输入值的类型输入过滤器源为voidoutputType过滤器输出值的类型输出过滤器汇为voidmode上述三种过滤器模式之一functor定义如何从输入值产生输出值的函数对象。过滤器用operator串联串联时前一过滤器的outputType必须等于后一过滤器的inputType。operator在 include/oneapi/tbb/parallel_pipeline.h#L104-L109 中实现它把两个过滤器节点组合成一个新的filter_node形成一棵解析树parallel_pipeline最终把这棵树扁平化为底层运行时的过滤器链见 include/oneapi/tbb/detail/_pipeline_filters.h#L335-L378 的filter_node与引用计数机制。先构造再运行的形式文档同样给出便于复用过滤器链void RunPipeline( int ntoken, FILE* input_file, FILE* output_file ) { oneapi::tbb::filtervoid,TextSlice* f1( oneapi::tbb::filter_mode::serial_in_order, MyInputFunc(input_file) ); oneapi::tbb::filterTextSlice*,TextSlice* f2(oneapi::tbb::filter_mode::parallel, MyTransformFunc() ); oneapi::tbb::filterTextSlice*,void f3(oneapi::tbb::filter_mode::serial_in_order, MyOutputFunc(output_file) ); oneapi::tbb::filtervoid,void f f1 f2 f3; oneapi::tbb::parallel_pipeline(ntoken,f); }注意filter类还有移动/拷贝构造与赋值、clear()方法见 filter 类定义filter 节点通过引用计数add_ref/remove_ref管理生命周期多个filter对象可以共享同一棵过滤器树因此先构造再运行是安全的。三个 functor 的职责与实现输出 functor最简单把TextSlice写到文件并释放。由于它位于管线末端operator()返回voidvoid MyOutputFunc::operator()( TextSlice* out ) const { size_t n fwrite( out-begin(), 1, out-size(), my_output_file ); if( n!out-size() ) { fprintf(stderr,Cant write into file %s\n, OutputFileName); exit(1); } out-free(); }变换 functor中间过滤器返回它产出的新TextSlice*。实现要点是先追加终止符让strtol在数字恰好在 slice 末尾时也能正确工作再把非数字字符原样复制、数字用strtol解析并平方后以sprintf写回输出 buffer 分配两倍输入长度文档与示例注释指出非负整数 n 的平方位数不可能超过 n 位数的两倍因此无需溢出检查TextSlice* MyTransformFunc::operator()( TextSlice* input ) const { *input-end() \0; char* p input-begin(); TextSlice* out TextSlice::allocate( 2*MAX_CHAR_PER_INPUT_SLICE ); char* q out-begin(); for(;;) { while( pinput-end() !isdigit(*p) ) *q *p; if( pinput-end() ) break; long x strtol( p, p, 10 ); long y x*x; sprintf(q,%ld,y); q strchr(q,0); } out-set_end(q); input-free(); return out; }输入 functor最复杂必须保证没有数字跨越 chunk 边界。当它发现某个可能跨入下一个 slice 的不完整数字时会把数字的残余部分复制到下一个 slice同时必须指示输入何时结束——通过调用特殊参数flow_control的stop()方法。所有用于管线第一个过滤器的 functor 都必须采用这一约定TextSlice* MyInputFunc::operator()( oneapi::tbb::flow_control fc ) const { if( !next_slice ) next_slice TextSlice::allocate( MAX_CHAR_PER_INPUT_SLICE ); size_t m next_slice-avail(); size_t n fread( next_slice-end(), 1, m, input_file ); if( !n next_slice-size()0 ) { // No more characters to process fc.stop(); return NULL; } else { TextSlice* t next_slice; next_slice TextSlice::allocate( MAX_CHAR_PER_INPUT_SLICE ); char* p t-end()n; if( nm ) { // Might have read partial number. // If so, transfer characters of partial number to next slice. while( pt-begin() isdigit(p[-1]) ) --p; assert(pt-begin(),Number too large to fit in buffer.\n); next_slice-append( p, t-end()n ); } t-set_end(p); return t; } }输入 functor 还必须提供拷贝构造函数因为 functor 在从filter_t即filter构造时会被拷贝一次管线运行时还会再拷贝一次。这一点可以从MyInputFunc的拷贝构造仅复制input_file看出。底层实现中concrete_filter保存的是const Body my_body引用见 include/oneapi/tbb/detail/_pipeline_filters.h#L232-L254而filter_node_leaf持有按值拷贝的const Body my_bodyinclude/oneapi/tbb/detail/_pipeline_filters.h#L438-L447这印证了构造时拷贝、运行时再拷贝的说法。重要警示body 的 operator() 必须是 const文档对此专门给出 CAUTION提供给管线过滤器的 body 对象可能被拷贝因此其operator()不应修改 body 自身。否则修改是否对调用parallel_pipeline的线程可见取决于operator()作用于原件还是拷贝件行为不确定。作为提醒parallel_pipeline要求 body 的operator()声明为const——这一点在MyInputFunc、MyTransformFunc、MyOutputFunc中均有体现。三种过滤器模式的顺序语义输入过滤器在本例中必须为serial_in_order因为它从顺序文件中读取 chunk输出过滤器必须按输入顺序写出 chunk因此同为serial_in_order所有serial_in_order过滤器按同一顺序处理 item如果某个 item 到达MyOutputFunc时晚于MyInputFunc建立起的顺序管线会自动推迟对该 item 的operator()调用直到其所有前驱处理完毕另一类串行过滤器serial_out_of_order不保持顺序中间过滤器只操作局部数据其 functor 的任意次调用都可以并发运行因此指定为parallel。运行时实现上base_filter用位标志区分串行/并行与有序/无序include/oneapi/tbb/detail/_pipeline_filters.h#L58-L64串行过滤器会配备输入缓冲区my_input_buffer用于按序暂存到达的 token。深入token 数量、吞吐瓶颈与窗口大小ntoken 的选择管线吞吐量受两个约束限制详见 Throughput_of_pipelinetoken 数量上限用N个 token 运行时并行执行的操作不可能超过N个。N选得太低限制并行度选得太高则可能占用过多资源例如更多缓冲区。官方示例 square.cpp 给出的经验值是nthreads * 4注释说明每个线程需要多于一个在途 token 才能让所有线程保持忙碌2–4 倍即可。最慢串行过滤器是瓶颈即使管线没有任何并行过滤器吞吐量也被最慢的串行过滤器限制——其他过滤器再快也无济于事。因此应尽量让串行过滤器保持轻快并把工作尽可能移到并行过滤器上。文本示例的现实速度上限文档明确提醒文本处理示例的加速比相对有限因为串行过滤器受限于系统 I/O 速度——即使文件在本地磁盘上加速比也不太可能明显超过 2。要让管线真正受益并行过滤器必须相对于串行过滤器承担较重的计算。窗口chunk大小的权衡窗口大小每个 token 的子问题规模同样限制吞吐窗口太小调度与传递开销可能压过有效工作窗口太大可能溢出 cache好的经验法则是在仍能装进 cache 的前提下尽量取大的窗口大小通常需要少量实验来确定。环形缓冲区复用进阶优化Using_Circular_Buffers 介绍了一种降低管线过滤器间分配/释放开销的技巧当第一个创建 item 的过滤器与最后一个消费 item 的过滤器都是serial_in_order时可以通过大小至少为ntoken的环形缓冲区来分配与释放 item。原因在于最多只有ntoken个 item 同时在途且 item 会按分配顺序被释放因此环形缓冲区回绕到复用某个 item 时该 item 必然已经结束上一次在管线中的使用。反之如果首尾过滤器不是serial_in_order则必须自己跟踪哪些缓冲区正在使用中因为缓冲区不会按分配顺序退役。线性管线的局限Non-Linear_Pipelines 说明parallel_pipeline只支持线性管线不直接处理分支、汇合等复杂拓扑。若业务需要分叉/合并结构可以把过滤器按拓扑排序压成线性顺序来近似。代价分析如下线性化损失的只是延迟latency而非吞吐量throughput。延迟是 token 从管线首端流到末端的时间例如一个原本 A、B 可并发、D、E 也可并发的五过滤器拓扑原始延迟相当于 3 级过滤器线性化后延迟变为 5 级吞吐量不变因为无论拓扑如何吞吐仍受最慢串行过滤器限制线性化后那些不需要处理、只需透传的过滤器A、B、D、E需要修改其 functor 以正确透传对象。文档的结论是parallel_pipeline不支持非线性管线是一个收益对代价的合理权衡——若支持非线性会显著增加编程复杂度且并不改善吞吐。从模式到实战结合示例运行与验证仓库中的 examples/parallel_pipeline/square 目录提供了与文档完全对应的可运行示例其 README.md 说明了构建与运行方式# 构建 cmake path_to_example cmake --build . # 运行也可使用预设 target make run_square # 预定义参数 make perf_run_square # 性能测试参数 make light_test_square # 快速验证参数命令行参数与 square.cpp 的 main 对应n-of-threads使用的线程数可用low[:high]区间或auto示例内部通过oneapi::tbb::global_control的max_allowed_parallelism控制input-file/output-file输入输出文件名max-slice-size每个 slice 的最大字符数对应文档中的MAX_CHAR_PER_INPUT_SLICE默认 4000silent除耗时外不输出其他信息。示例的main会先用 1 个线程做串行基准运行再用自动选择的线程数做并行运行并借助oneapi::tbb::tick_count计时方便读者直观对比加速效果——这正是文档所讲吞吐受串行过滤器限制这一结论的实测入口。小结两种复杂循环模式的选用准则场景首选模板关键注意点迭代空间未知链表等parallel_for_eachbody 的operator()必须const单元素处理需足够重几千条指令量级循环体动态追加任务树遍历等parallel_for_eachfeeder.addfeeder 追加的 item 会在本次调用内处理完迭代器为 random-access 时可扩展数据分阶段流动、各阶段并行度不同parallel_pipeline输入过滤器必须用flow_control::stop()标识结束body 的operator()必须const管线吞吐调优—合理选择ntoken经验值线程数×2~4窗口大小在能装进 cache 的前提下取大值把重活移入并行过滤器首尾均为 serial_in_order 的管线—可用大小≥ntoken 的环形缓冲区消除分配/释放开销需要分支/汇合拓扑—拓扑排序线性化接受延迟增加、吞吐不变这些高层模板的目标正如 Summary_of_Loops_and_Pipelines 所总结以任务模式级别设计软件、让 generic 模板按需定制从而在不必从头管理线程的前提下充分挖掘多核芯片的计算能力。需要进一步深入时可以阅读本文引用的源码文件 include/oneapi/tbb/parallel_for_each.h、include/oneapi/tbb/parallel_pipeline.h、include/oneapi/tbb/detail/_pipeline_filters.h以及示例 examples/parallel_pipeline/square/square.cpp 与测试 test/tbb/test_parallel_for_each.cpp。赞分享并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载相关推荐oneTBB 循环与流水线算法全景指南parallel_for、parallel_reduce、parallel_for_each 与 parallel_pipeline 实战详解oneTBB 循环与流水线算法全景指南parallel_for、parallel_reduce、parallel_for_each 与 parallel_pi并发编程高性能计算mistral.rs Python SDK 实战加载 Qwen3-Next 混合架构模型并做 Q4K 在位量化推理mistral.rs Python SDK 实战加载 Qwen3 Next 混合架构模型并做 Q4K 在位量化推理 本篇技术文章围绕 mistral.rs 的并发编程高性能计算如何使用oneTBB实现高效并行循环parallel_for实战教程 如何使用oneTBB实现高效并行循环parallel_for实战教程 在当今多核处理器普及的时代如何充分利用硬件资源提升程序性能成为开发者面临的重要挑并发编程高性能计算上一篇解决Ward服务器监控工具10大痛点从部署到告警全流程排错指南下一篇Orchis-theme 项目常见问题解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表