mold 仓库 TBB 并发编程指南:何时不应使用队列,用 parallel_pipeline 与 parallel_for_each 取代显式队列
2026/9/15 16:12:13 网站建设 项目流程

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

【免费下载链接】moldmold: A Modern Linker 🦠项目地址: https://gitcode.com/GitHub_Trending/mo/mold

在并行程序中,队列常被用来缓冲生产者与消费者之间的数据流,但显式队列并非总是最优选择。本指南以 mold 仓库内置的第三方 oneAPI Threading Building Blocks(TBB)用户指南文档 When_Not_to_Use_Queues.rst 为骨架,结合 concurrent_queue 类文档、parallel_pipeline 头文件 与 parallel_for_each 头文件 等源码实现,系统讲解队列成为性能瓶颈的根本原因,以及如何用parallel_pipelineparallel_for_each替代显式队列,让读者掌握"何时该用队列、何时不该用"的判别方法与可落地的替代实现。

为什么先质疑"队列"这个默认选择

TBB 提供了两类并发队列容器(详见 Concurrent_Queue_Classes.rst):

  • concurrent_queue<T, Alloc>:无界、无阻塞操作,核心操作是pushtry_poptry_pop在一个原子操作内完成"检查是否有元素 + 弹出",避免了"先查 empty 再 pop"这类组合操作固有的竞态问题(如 Concurrent_Queue_Classes.rst 中std::queue反例所示:empty()刚返回 true,另一线程可能恰好取走最后一个元素)。它保证"若同一线程连续 push 两个值,另一线程按序 pop 出这两个值时,弹出顺序与压入顺序一致"。
  • concurrent_bounded_queue<T, Alloc>:在concurrent_queue基础上增加阻塞操作与容量限制。pop(item)会等待直到成功,push(item)会等待直到不超出容量,try_push(item)仅在不会超限时压入;size()返回有符号整数,其定义为"已开始的 push 次数减去已开始的 pop 次数",若空队列上有 n 个挂起的 pop,size()返回 -n,生产者借此得知有多少消费者在等待;empty()当且仅当size()非正时为真。默认无界,可通过set_capacity设置容量,但有界队列比无界队列慢,若程序其他约束已能防止队列过大,就不应设置容量(见 Concurrent_Queue_Classes.rst)。

然而,正如 When_Not_to_Use_Queues.rst 开篇所强调的:在使用显式队列之前,应当先考虑用parallel_for_eachparallel_pipeline代替。这两者在多数场景下比队列更高效,原因要从队列的三个固有弱点说起。

显式队列的三个固有瓶颈

原文档从三个角度剖析了队列为何天生低效,这些结论适用于一切"以显式队列为核心的生产者—消费者"结构:

  1. 队列本质上是瓶颈:它必须维护先进先出(FIFO)顺序。无论底层用何种无锁算法实现,为保住 FIFO 语义,进出队列的操作都要在同一个临界点上竞争,所有数据流都被迫穿过这一个"咽喉要道"。
  2. 消费线程可能空等:执行pop/try_pop的线程必须等到值被 push 之后才能继续。在队列为空的时间窗内,消费者线程只能空转或阻塞,CPU 资源被浪费。
  3. 队列是被动的数据结构:线程 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 filter<void,void>& filter_chain ); void parallel_pipeline( size_t max_number_of_live_tokens, const filter<void,void>& filter_chain, task_group_context& context ); }

流水线由一组filter依次串联而成。通用构造方式为:

parallel_pipeline( max_number_of_live_tokens, make_filter<void,I1>(mode0,g0) & make_filter<I1,I2>(mode1,g1) & make_filter<I2,I3>(mode2,g2) & ... make_filter<In,void>(moden,gn) );

关键约束与语义(来自 parallel_pipeline_func.rst):

  • 每个 filter 用两个模板参数指定输入类型输出类型;第一个 filter 的输入类型必须是void,最后一个 filter 的输出类型必须是void
  • 各 filter 通过operator&拼接,拼接要求左侧 filter 的输出类型与右侧 filter 的输入类型一致,最终合并为一个filter<void,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_order
filter_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.rst):

class 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 sum=0; parallel_pipeline( /*max_number_of_live_token=*/16, make_filter<void,float*>( filter_mode::serial_in_order, &-> float*{ if( first<last ) { return first++; } else { fc.stop(); return nullptr; } } ) & make_filter<float*,float>( filter_mode::parallel, [](float* p){return (*p)*(*p);} ) & make_filter<float,void>( filter_mode::serial_in_order, & {sum+=x;} ) ); return sqrt(sum); }

这个例子浓缩了全部要点:输入 filterserial_in_order顺序喂数,用flow_control::stop()结束输入;中间 filter是纯函数式的平方运算,标记为parallel从而允许任意多个数据项并发计算;输出 filterserial_in_order累加,保证求和顺序与输入一致。三个 filter 的类型首尾相接:void→float*float*→floatfloat→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_filter<void,TextSlice*>( oneapi::tbb::filter_mode::serial_in_order, MyInputFunc(input_file) ) & oneapi::tbb::make_filter<TextSlice*,TextSlice*>( oneapi::tbb::filter_mode::parallel, MyTransformFunc() ) & oneapi::tbb::make_filter<TextSlice*,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 自身,且必须声明为const(Working_on_the_Assembly_Line_pipeline.rst 中的 CAUTION 明确要求)。

替代方案二:parallel_for_each——无界数据流与动态追加工作

当数据天然是一个"可遍历的集合"而非"需要多阶段变换的流"时,parallel_for_each是比队列更直接的替代。其公开接口定义在 parallel_for_each.h:

template<typename Iterator, typename Body> void parallel_for_each(Iterator first, Iterator last, const Body& body); template<typename Range, typename Body> void parallel_for_each(Range& rng, const Body& body); // 另有接受 task_group_context& 的重载版本

parallel_for_each的独特之处在于feeder 机制(parallel_for_each.h):body 的operator()除了接收数据项之外,还可以额外接收一个feeder<Item>&,并在处理过程中调用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, feeder<item_type>&) constB::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_each+feeder.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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询