尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

mold 仓库 TBB 并发编程指南:何时不应使用队列,用 parallel_pipeline 与 parallel_for_each 取代显式队列

mold 仓库 TBB 并发编程指南:何时不应使用队列,用 parallel_pipeline 与 parallel_for_each 取代显式队列 mold 仓库 TBB 并发编程指南何时不应使用队列用 parallel_pipeline 与 parallel_for_each 取代显式队列【免费下载链接】moldmold: A Modern Linker 项目地址: https://gitcode.com/GitHub_Trending/mo/mold在并行程序中队列常被用来缓冲生产者与消费者之间的数据流但显式队列并非总是最优选择。本指南以 mold 仓库内置的第三方 oneAPI Threading Building BlocksTBB用户指南文档 When_Not_to_Use_Queues.rst 为骨架结合 concurrent_queue 类文档、parallel_pipeline 头文件 与 parallel_for_each 头文件 等源码实现系统讲解队列成为性能瓶颈的根本原因以及如何用parallel_pipeline和parallel_for_each替代显式队列让读者掌握何时该用队列、何时不该用的判别方法与可落地的替代实现。为什么先质疑队列这个默认选择TBB 提供了两类并发队列容器详见 Concurrent_Queue_Classes.rstconcurrent_queueT, Alloc无界、无阻塞操作核心操作是push与try_pop。try_pop在一个原子操作内完成检查是否有元素 弹出避免了先查 empty 再 pop这类组合操作固有的竞态问题如 Concurrent_Queue_Classes.rst 中std::queue反例所示empty()刚返回 true另一线程可能恰好取走最后一个元素。它保证若同一线程连续 push 两个值另一线程按序 pop 出这两个值时弹出顺序与压入顺序一致。concurrent_bounded_queueT, Alloc在concurrent_queue基础上增加阻塞操作与容量限制。pop(item)会等待直到成功push(item)会等待直到不超出容量try_push(item)仅在不会超限时压入size()返回有符号整数其定义为已开始的 push 次数减去已开始的 pop 次数若空队列上有 n 个挂起的 popsize()返回 -n生产者借此得知有多少消费者在等待empty()当且仅当size()非正时为真。默认无界可通过set_capacity设置容量但有界队列比无界队列慢若程序其他约束已能防止队列过大就不应设置容量见 Concurrent_Queue_Classes.rst。然而正如 When_Not_to_Use_Queues.rst 开篇所强调的在使用显式队列之前应当先考虑用parallel_for_each或parallel_pipeline代替。这两者在多数场景下比队列更高效原因要从队列的三个固有弱点说起。显式队列的三个固有瓶颈原文档从三个角度剖析了队列为何天生低效这些结论适用于一切以显式队列为核心的生产者—消费者结构队列本质上是瓶颈它必须维护先进先出FIFO顺序。无论底层用何种无锁算法实现为保住 FIFO 语义进出队列的操作都要在同一个临界点上竞争所有数据流都被迫穿过这一个咽喉要道。消费线程可能空等执行pop/try_pop的线程必须等到值被 push 之后才能继续。在队列为空的时间窗内消费者线程只能空转或阻塞CPU 资源被浪费。队列是被动的数据结构线程 push 一个值后该值可能要在队列里停留一段时间才被 pop。期间这个值以及它引用的所有对象在缓存中逐渐变冷cold更糟的是若由另一个 CPU 上的线程将其 pop该值及其引用对象还必须被搬运到另一颗处理器上缓存局部性被彻底破坏。换句话说显式队列把同步与数据传递耦合在了一个被动容器上既制造了等待又破坏了数据热度。替代方案一parallel_pipeline——隐式线程化的流水线parallel_pipeline正是针对上述瓶颈设计的。原文档指出它的线程化是隐式的implicit它优化工作线程的使用让线程在值尚未到来时去执行其他工作而不是空等同时它会尽力让被处理的数据项在缓存中保持热hot。这两点恰好一一对应地化解了队列的第 2、3 条弱点而第 1 条 FIFO 瓶颈则通过只对必须保序的环节保序来回避详见下文 filter_mode。算法签名与过滤器链TBB 规范文档 parallel_pipeline_func.rst 给出了两种重载形式头文件 parallel_pipeline.h 中实现了对应接口// 定义于 oneapi/tbb/parallel_pipeline.h namespace oneapi::tbb { void parallel_pipeline( size_t max_number_of_live_tokens, const filtervoid,void filter_chain ); void parallel_pipeline( size_t max_number_of_live_tokens, const filtervoid,void filter_chain, task_group_context context ); }流水线由一组filter依次串联而成。通用构造方式为parallel_pipeline( max_number_of_live_tokens, make_filtervoid,I1(mode0,g0) make_filterI1,I2(mode1,g1) make_filterI2,I3(mode2,g2) ... make_filterIn,void(moden,gn) );关键约束与语义来自 parallel_pipeline_func.rst每个 filter 用两个模板参数指定输入类型与输出类型第一个 filter 的输入类型必须是void最后一个 filter 的输出类型必须是void。各 filter 通过operator拼接拼接要求左侧 filter 的输出类型与右侧 filter 的输入类型一致最终合并为一个filtervoid,void链。max_number_of_live_tokens是在途 token 数的上限即同时处于处理流程中的数据项个数上限。它限定了整体并行度一旦达到上限输入 filter 不会创建新 token直到输出 filter 销毁一个 token详见 Working_on_the_Assembly_Line_pipeline.rst。这个机制防止了中间的无序 filter 因下游跟不上而无限累积 token导致的资源失控。可选参数context指定任务组上下文缺省时算法运行在自己绑定的上下文中。filter_mode三种执行模式filter_mode是定义在 parallel_pipeline.h 中的枚举类三种取值及其语义参见 filter_mode_enum.rst模式语义filter_mode::parallel可同时处理多个数据项且不要求特定顺序filter_is_out_of_orderfilter_mode::serial_in_order一次只处理一个数据项所有serial_in_orderfilter 按第一个此类 filter 确立的顺序依次处理自动保证整体顺序filter_mode::serial_out_of_order一次只处理一个数据项但不保持顺序filter_is_serial \| filter_is_out_of_order在 parallel_pipeline.h 的实现中模式通过base_filter的位标志组合表达parallel对应filter_is_out_of_orderserial_in_order对应filter_is_serialserial_out_of_order是两者的按位或。这一设计让流水线调度器能够区分可并发执行的环节与必须串行的环节从而只对必要的环节付出保序代价——这正是它比全量 FIFO 队列高效的结构性原因。flow_control通知输入结束第一个 filter 的 functor 需要额外的flow_control参数用于通知流水线输入流已结束参见 flow_control_cls.rstclass flow_control { public: void stop(); // 表示第一个 filter 已到达输入流的末尾 };functor 收到flow_control fc后若仍有下一个值则返回该值若到达输入末尾则调用fc.stop()并返回一个不会被传递给下一级 filter 的哑值通常是nullptr。完整可运行示例均方根计算规范文档 parallel_pipeline_func.rst 给出了一个完整的均方根Root-Mean-Square示例直观展示了三级 filter 的搭建方式float RootMeanSquare( float* first, float* last ) { float sum0; parallel_pipeline( /*max_number_of_live_token*/16, make_filtervoid,float*( filter_mode::serial_in_order, - float*{ if( firstlast ) { return first; } else { fc.stop(); return nullptr; } } ) make_filterfloat*,float( filter_mode::parallel, [](float* p){return (*p)*(*p);} ) make_filterfloat,void( filter_mode::serial_in_order, {sumx;} ) ); return sqrt(sum); }这个例子浓缩了全部要点输入 filter用serial_in_order顺序喂数用flow_control::stop()结束输入中间 filter是纯函数式的平方运算标记为parallel从而允许任意多个数据项并发计算输出 filter用serial_in_order累加保证求和顺序与输入一致。三个 filter 的类型首尾相接void→float*、float*→float、float→void。更复杂的实战案例文本格式化流水线用户指南 Working_on_the_Assembly_Line_pipeline.rst 给出了一个更贴近真实工程的例子读取文本文件、把其中的十进制数字替换为其平方值、再按序写出。其核心结构是顺序读 → 并行转换 → 顺序写的三级流水线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) ) ); }该示例中有几个值得注意的工程细节按块chunk处理为摊薄并行调度的开销数据按约 4000 字符的TextSlice分块流动filter 之间传递的是TextSlice*指针避免复制大块数据的开销。顺序语义输入 filter 必须serial_in_order顺序读文件输出 filter 也必须serial_in_order按原顺序写回当某个数据项到达输出 filter 时其前驱尚未处理完流水线会自动延迟调用输出 functor直至前驱完成。中间 filter 只操作纯局部数据故声明为parallel任意多个调用可并发运行。跨块边界处理输入 functor 需要保证数字不被切断在相邻块边界上——当读到疑似跨块的数字时把部分数字拷贝到下一块并通过flow_control fc参数在输入耗尽时调用fc.stop()终止流水线输入 functor 必须使用该惯用法。body 必须是 const由于传入 filter 的 body 对象可能被拷贝其operator()不得修改 body 自身且必须声明为constWorking_on_the_Assembly_Line_pipeline.rst 中的 CAUTION 明确要求。替代方案二parallel_for_each——无界数据流与动态追加工作当数据天然是一个可遍历的集合而非需要多阶段变换的流时parallel_for_each是比队列更直接的替代。其公开接口定义在 parallel_for_each.htemplatetypename Iterator, typename Body void parallel_for_each(Iterator first, Iterator last, const Body body); templatetypename Range, typename Body void parallel_for_each(Range rng, const Body body); // 另有接受 task_group_context 的重载版本parallel_for_each的独特之处在于feeder 机制parallel_for_each.hbody 的operator()除了接收数据项之外还可以额外接收一个feederItem并在处理过程中调用feeder.add(item)动态追加新的工作项新项会被立即调度到任务池中并行处理。从 parallel_for_each.h 的实现看追加项会生成feeder_item_task并通过spawn提交到执行上下文parallel_for_each.h而随机访问迭代器版本则直接复用parallel_for的分块调度parallel_for_each.h。parallel_for_each对 body 的要求见 parallel_for_each.h非常轻量B::operator()(item, feederitem_type) const或B::operator()(item) const外加数据项的可拷贝构造与析构。它同样遵循隐式线程化原则——工作项由调度器按需分配到工作线程线程不会像 pop 空队列那样空等追加的新工作项也会就地热执行避免数据项在被动容器中变冷。补充parallel_pipeline 的线性限制与取舍需要明确的是parallel_pipeline只支持线性流水线Non-Linear_Pipelines.rst。对于更复杂的拓扑如分叉、汇合文档给出的做法是先把 filter 拓扑排序成线性顺序再按序串联。其代价分析值得记住强制线性化只损失延迟latency不损失吞吐throughput。延迟指一个 token 从流水线头流到尾的时间原非线性拓扑中 A/B 可并发、D/E 可并发延迟为三级线性化后延迟变成五级。吞吐始终受限于最慢的串行 filter与拓扑无关。因此若parallel_pipeline支持非线性拓扑只会增加大量编程复杂度而不会提升吞吐——线性限制是收益与成本之间的合理取舍Non-Linear_Pipelines.rst。决策建议何时该用队列何时不该用综合原文档与上述源码分析可给出如下判别准则场景推荐做法数据流需经过多阶段变换各阶段串并行属性明确用parallel_pipeline靠隐式调度避免空等、保持缓存热度数据是一组已知集合处理过程可能动态产生新工作项用parallel_for_eachfeeder.add需要显式缓冲、跨模块解耦或必须由上层逻辑自行控制同步时机用concurrent_queue无界、无阻塞必须阻塞等待、且需要容量上限控制背压用concurrent_bounded_queue注意有界会变慢非必要不设容量原文档给出的理由在这里可以闭环队列必须维护 FIFO 而成为天然瓶颈→ 对应parallel_pipeline只在serial_in_order环节保序pop 线程可能空等→ 对应隐式调度让线程值未到先做别的事队列是被动容器导致缓存变冷/跨核搬运→ 对应流水线尽力保持数据项在缓存中热。当你的生产者—消费者结构本质上就是一条流水线或一次遍历时显式队列几乎总可以被这两种算法替代并获得更好的资源利用与缓存行为。本文全部依据均来自 mold 仓库内置的 TBB 文档与源码读者可沿 When_Not_to_Use_Queues.rst → Concurrent_Queue_Classes.rst → Working_on_the_Assembly_Line_pipeline.rst → Non-Linear_Pipelines.rst 的脉络继续深入并在 parallel_pipeline.h 与 parallel_for_each.h 中核对每个 API 的真实行为。【免费下载链接】moldmold: A Modern Linker 项目地址: https://gitcode.com/GitHub_Trending/mo/mold创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表