mold 仓库第三方依赖 TBB Flow Graph 的 try_put_and_wait:单消息级同步等待接口详解
2026/9/14 19:13:59 网站建设 项目流程

mold 仓库第三方依赖 TBB Flow Graph 的 try_put_and_wait:单消息级同步等待接口详解

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

本文基于 mold 仓库third-party/tbb子目录随仓库分发的 Intel oneAPI Threading Building Blocks(TBB)Flow Graph 参考文档编写。try_put_and_wait是 Flow Graph 接收节点上新增的预览(Preview)接口:它把一个消息送入 Flow Graph,并阻塞等待与该消息相关的全部工作完成,从而把传统graph::wait_for_all的"整图级同步"细化为"单消息级同步"。读完本文,你将掌握该接口的启用宏、覆盖的节点类型、每个节点的语义细节、底层 wait-tree 实现原理以及实战示例与注意事项。

一、特性概述:为什么需要单消息级等待

在 TBB Flow Graph 中,节点之间通过try_put传递消息,任务由任务调度器异步执行。以往要让生产者线程确认某个输入消息已被图内所有节点处理完毕,只能调用graph::wait_for_all(),它会等待图中所有工作完成——包括与该输入消息完全无关的其他节点、其他并行提交的任务。

try_put_and_wait改变了这一局面:对任意支持该接口的接收节点调用node.try_put_and_wait(msg),等价于先执行node.try_put(msg),再等待仅与msg相关的工作全部结束。由于它不再等待无关任务,相比graph::wait_for_all可以显著降低延迟,特别适合"提交一批独立输入、逐个处理并立即回收结果"的高吞吐流水线场景。

调用返回后,以下条件必然成立(参考文档):

  • 图中任何由msg或其任意中间结果触发的任务都已执行完毕;
  • msg计算出的任何中间结果,不再残留在图中任何节点的缓冲区中。

注意这里的语义边界:由其他调用提交的消息所触发的任务,不保证完成。这正是它与wait_for_all的本质差异——同步范围是"消息级"而非"全图级"。

二、两类必须规避的使用限制

参考文档明确给出了两条 caution,使用不当会导致死锁或提前返回:

1. 不要在图的末端使用缓冲类节点。缓冲节点(buffer_nodequeue_node等)本身不会消费最终结果,如果它是图中最后一个接收方,try_put_and_wait会一直等待"永远不会发生"的消费动作,从而无限等待。因此请确保图末端是能够真正消费消息的function_node等执行类节点,或自行显式取出缓冲中的结果。

2.multifunction_nodeasync_node暂不支持。图中只要包含这两类节点,try_put_and_wait可能提前返回——即使初始输入消息的计算仍在进行中。设计图结构时应避免在等待路径上引入这两类节点。

三、启用方式:两个预览宏

该接口属于 TBB 预览特性,需要先定义宏再包含头文件(参考文档 API 章节):

#define TBB_PREVIEW_FLOW_GRAPH_FEATURES // 宏选项 1:一次性开启全部 Flow Graph 预览特性 #define TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT // 宏选项 2:仅开启本特性 #include <oneapi/tbb/flow_graph.h>

两个宏的关系可以在源码配置头中找到印证:third-party/tbb/include/oneapi/tbb/detail/_config.h中定义:

#ifndef __TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT #define __TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT (TBB_PREVIEW_FLOW_GRAPH_FEATURES \ || TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT) #endif

也就是说,即使只定义TBB_PREVIEW_FLOW_GRAPH_FEATURES__TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT也会被置 1;__TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT是内部头文件(flow_graph.h_flow_graph_impl.h等)中真正驱动#if编译分支的底层开关。整个接口在flow_graph.h中同样由#if __TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT包裹,未定义宏时该成员函数根本不存在。

四、API 一览:覆盖 11 种接收节点

参考文档的 Synopsis 给出了全部支持该接口的节点声明(namespace oneapi::tbb内,flow别名即指此命名空间):

template <typename Output, typename Policy = /*default-policy*/> class continue_node { public: bool try_put_and_wait(const continue_msg& input); }; template <typename Input, typename Output = continue_msg, typename Policy = /*default-policy*/> class function_node { public: bool try_put_and_wait(const Input& input); }; template <typename T> class overwrite_node { public: bool try_put_and_wait(const T& input); }; template <typename T> class write_once_node { public: bool try_put_and_wait(const T& input); }; template <typename T> class buffer_node { public: bool try_put_and_wait(const T& input); }; template <typename T> class queue_node { public: bool try_put_and_wait(const T& input); }; template <typename T, typename Compare = std::less<T>> class priority_queue_node { public: bool try_put_and_wait(const T& input); }; template <typename T> class sequencer_node { public: bool try_put_and_wait(const T& input); }; template <typename T, typename DecrementType = continue_msg> class limiter_node { public: bool try_put_and_wait(const T& input); }; template <typename T> class broadcast_node { public: bool try_put_and_wait(const T& input); }; template <typename TupleType> class split_node { public: bool try_put_and_wait(const TupleType& input); };

可以看到,除join_nodeindexer_nodemultifunction_nodeasync_node外的绝大多数常用接收节点均已覆盖。每个节点的"等待"语义完全一致——等待input在图中的全部相关任务执行完毕、且相关中间结果不在任何缓冲区残留;区别只体现在各节点接收消息时的行为(Effects)与返回值(Returns)上。

五、各节点成员函数语义逐项解析

参考文档对每个节点的 Effects/Returns 有精确定义,整理如下:

1. continue_node(延续节点)

template <typename Output, typename Policy> bool continue_node<Output, Policy>::try_put_and_wait(const continue_msg& input)
  • Effects:递增收到的输入信号计数。若递增后的计数等于已知前驱节点数,则执行用户提供的body函数对象;随后等待该input在图中完成。
  • Returns:恒为true

2. function_node(函数节点)

template <typename Input, typename Output, typename Policy> bool function_node<Input, Output, Policy>::try_put_and_wait(const Input& input)
  • Effects:若并发限制允许,立即对输入消息input执行用户 body;否则根据节点的Policyserial/unlimited/queueing/rejecting等),将消息入队或拒绝。
  • Returns:输入被接受时返回true,否则返回false。这是少数可能返回false的节点之一,调用方需检查返回值。

3. overwrite_node(覆盖节点)

template <typename T> bool overwrite_node<T>::try_put_and_wait(const T& input)
  • Effects:将input存入内部单元素缓冲区,并向所有后继广播。
  • Returns:恒为true
  • ⚠️ 特别注意overwrite_node被后继取走后不会删除该元素,须显式调用clear()或写入新元素覆盖,否则try_put_and_wait可能无限等待。

4. write_once_node(只写一次节点)

template <typename T> bool write_once_node<T>::try_put_and_wait(const T& input)
  • Effects:若内部单元素缓冲区尚无有效值,则存入input并向后继广播;已有值则忽略本次输入。
  • Returns:构造后或调用clear()后的第一次调用返回true,其余返回false
  • ⚠️ 特别注意:同样不会自动清除元素,需显式调用clear()防止无限等待。

5. buffer_node / queue_node / priority_queue_node / sequencer_node(缓冲族节点)

bool buffer_node<T>::try_put_and_wait(const T& input); // 加入集合并尝试转发 bool queue_node<T>::try_put_and_wait(const T& input); // 加入并尝试转发最旧元素 bool priority_queue_node<T, Compare>::try_put_and_wait(const T& input); // 加入并尝试转发优先级最高且未转发的元素 bool sequencer_node<T>::try_put_and_wait(const T& input); // 加入并尝试按序转发下一个元素
  • Effects:四者的差别仅在"转发哪个元素"的策略上——buffer_node转发任意未转发元素,queue_node按 FIFO 转发最旧元素,priority_queue_node依据Compare(默认std::less<T>)转发最高优先级元素,sequencer_node按序号转发顺序中的下一个元素;等待语义相同。
  • Returns:均为true。结合第二节的 caution,这类节点若处于图末端会因无人消费而无限等待。

6. limiter_node(限流节点)

template <typename T, typename DecrementType = continue_msg> bool limiter_node<T, DecrementType>::try_put_and_wait(const T& input)
  • Effects:若当前广播计数低于阈值,向所有后继广播input
  • Returns:广播成功返回true,否则返回false

7. broadcast_node(广播节点)

template <typename T> bool broadcast_node<T>::try_put_and_wait(const T& input)
  • Effects:向所有后继广播input
  • Returns:恒为true——即使无法成功把消息转发给任何一个后继也返回true

8. split_node(拆分节点)

template <typename TupleType> bool split_node<TupleType>::try_put_and_wait(const TupleType& input)
  • Effects:把元组input的每个元素分别广播到对应输出端口——下标i的元素经端口i输出。
  • Returns:恒为true

六、底层实现原理:wait-tree 引用计数

从源码看,try_put_and_wait并非简单地在try_put之后调用wait_for_all,而是为单条消息建立独立的等待上下文。核心实现位于flow_graph.hreceiver<T>::try_put_and_wait

bool try_put_and_wait( const T& t ) { // Since try_put_and_wait is a blocking call, it is safe to create wait_context on stack d1::wait_context_vertex msg_wait_vertex{}; bool res = internal_try_put(t, message_metainfo{message_metainfo::waiters_type{&msg_wait_vertex}}); if (res) { __TBB_ASSERT(graph_reference().my_context != nullptr, "No wait_context associated with the Flow Graph"); d1::wait(msg_wait_vertex.get_context(), *graph_reference().my_context); } return res; }

其设计要点可以归纳为三块:

  1. 栈上等待顶点:为本次调用在栈上创建一个wait_context_vertex(wait-tree 顶点)。因为该调用是阻塞式的,等待期间栈帧必然存活,栈上创建是安全的。

  2. 元信息随消息传播message_metainfo携带等待者列表(std::forward_list<d1::wait_context_vertex*>,见flow_graph.h第 193-222 行),通过internal_try_put(t, metainfo)注入try_put_task。这条"消息→等待者"的关联随任务在图内逐级转发、缓冲、取回,从而把等待范围精确限定在该消息及其派生中间结果上。

  3. 引用计数生命周期管理d1::wait()阻塞当前线程,直到该消息在图中产生的所有任务完成。任务侧的支持来自_flow_graph_impl.h中的trackable_messages_graph_task:它持有消息对应的等待上下文列表与对应的引用顶点列表,任务终结(finalize)时按引用计数模式逐个release(1);当最后一个关联任务完成、所有计数归零后,栈上的wait_context_vertex被唤醒,调用返回。

从源码结构看,这种"消息级引用计数 + 等待上下文随消息传播"的机制,正是它能比wait_for_all更精准、更低延迟的原因:无关消息的任务完全不触碰本消息的等待顶点。

七、完整示例:parallel_for 流水线

参考文档自带一个可编译运行的示例(examples/try_put_and_wait_example.cpp):

#define TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT 1 #include <oneapi/tbb/flow_graph.h> #include <oneapi/tbb/parallel_for.h> #include <tuple> struct f1_body; struct f2_body; struct f3_body; struct f4_body; int main() { using namespace oneapi::tbb; flow::graph g; flow::broadcast_node<int> start_node(g); flow::function_node<int, int> f1(g, flow::unlimited, f1_body{}); flow::function_node<int, int> f2(g, flow::unlimited, f2_body{}); flow::function_node<int, int> f3(g, flow::unlimited, f3_body{}); flow::join_node<std::tuple<int, int>> join(g); flow::function_node<std::tuple<int, int>, int> f4(g, flow::serial, f4_body{}); flow::make_edge(start_node, f1); flow::make_edge(f1, f2); flow::make_edge(start_node, f3); flow::make_edge(f2, flow::input_port<0>(join)); flow::make_edge(f3, flow::input_port<1>(join)); flow::make_edge(join, f4); // Submit work into the graph parallel_for(0, 100, & { start_node.try_put_and_wait(input); // Post processing the result of input }); }

该示例的图拓扑为:start_nodebroadcast_node)把输入同时分发给两条路径——路径 A 经f1 → f2,路径 B 经f3,两条路径在join_node汇合后进入f4收尾。parallel_for的每次迭代调用一次start_node.try_put_and_wait(input)

  • 返回后即可保证与本次input相关的整条流水线(f1→f2、f3、join、f4)已全部执行完毕,可以在迭代内立即"后处理"该输入的结果;
  • 由其他迭代提交的输入的任务是否完成,则不作保证;
  • 100 次迭代在parallel_for中并行推进,每次迭代独立阻塞等待自己的消息,形成高吞吐的"多生产者、消息级确认"模式。

注意示例中try_put_and_wait作用于broadcast_node——图末端是真正消费消息的f4function_node),因而不会触发"缓冲节点在末端导致无限等待"的陷阱。

八、测试佐证:仓库内的验证覆盖

仓库third-party/tbb/test/tbb/目录下有一整套针对该特性的测试,覆盖了文档列出的几乎所有节点:test_function_node.cpptest_continue_node.cpptest_broadcast_node.cpptest_buffer_node.cpptest_queue_node.cpptest_priority_queue_node.cpptest_sequencer_node.cpptest_limiter_node.cpptest_overwrite_node.cpptest_write_once_node.cpptest_split_node.cpptest_indexer_node.cpptest_join_node_preview.cpp

其中test_buffering_try_put_and_wait.h是一个典型的消息级语义验证:它在单线程task_arena中构造缓冲节点 + function_node + 缓冲节点 + function_node链,当某个输入消息被处理时,body 内再向缓冲节点注入一批新工作项,随后验证try_put_and_wait返回后处理过的元素集合恰好只包含与等待消息相关的条目,从而确认等待粒度精确到单条消息。阅读这些测试可以帮助你理解边界情况(如处理中向图内重新投递消息时的计数行为)。

九、使用建议与要点回顾

  • 适用场景:需要对"每一条输入"做完成确认的流水线,如批量请求处理、在线服务中的请求级流水线;相比wait_for_all,可避免被无关流量拖慢。
  • 图末端约束:确保末端是消费型节点(如function_node);若必须使用overwrite_node/write_once_node,记得显式clear();不要用缓冲节点收尾。
  • 节点约束:图中存在multifunction_nodeasync_node时不要依赖本接口的完成语义。
  • 返回值检查function_nodelimiter_node可能返回false(拒绝/限流),应据此判断是否真的发生了等待。
  • 编译前提:必须在包含oneapi/tbb/flow_graph.h之前定义TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT(或TBB_PREVIEW_FLOW_GRAPH_FEATURES),否则该成员不存在;内部开关定义参见_config.h

上述内容以 mold 仓库内随附的 TBB 参考文档 try_put_and_wait.rst 为主体,结合头文件实现与测试用例整理而成。若需进一步深入,可继续阅读 flow_graph.h 中receivermessage_metainfo相关实现,以及_flow_graph_impl.htrackable_messages_graph_task的引用计数逻辑。

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

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询