☰
oneTBB 嵌套流图(Nested Flow Graphs)实战指南:在节点内部构建并执行子图
2026/10/8 1:31:22 网站建设 项目流程
  • 并发编程
  • 高性能计算

【免费下载链接】oneTBB

oneAPI Threading Building Blocks (oneTBB)

项目地址:https://gitcode.com/gh_mirrors/on/oneTBB
点击查看免费下载

导读

oneAPI Threading Building Blocks(oneTBB)的 Flow Graph 支持在算法嵌套之外的另一层组合能力:图的嵌套(graph nesting)。本文围绕 use_nested_flow_graphs.rst 的核心内容展开,讲解如何在一个function_node的 body 中临时构造并执行一个内层依赖图(dependence graph)或数据流图(data flow graph),以及当内层图结构在多轮调用间保持不变时,如何通过持久复用图来消除重复构造的开销。读完本文,你将掌握嵌套流图的两种实现模式、wait_for_all()的正确使用时机,以及嵌套图与task_arena、节点优先级组合使用的源码级细节。


一、嵌套流图:在节点内部再建一张图

oneTBB Flow Graph 的组合能力体现在两个层次:

  1. 算法嵌套:在一个节点内部调用parallel_for、parallel_reduce等并行算法;
  2. 图嵌套:在一个节点的 body 内直接构造并运行另一张独立的graph。

嵌套流图适用于把一个大的数据流问题拆分为多个内部子流程的场景,例如外层图负责整体调度与扇出/扇入,而每个节点内部运行一段独立的子流程(子图)。每个graph对象都拥有独立的task_group_context(见 flow_graph.h 的graph::graph()构造实现,其中my_context = new ... task_group_context(FLOW_TASKS)),因此内层图的取消与异常处理与外层图相互隔离,这是嵌套能够成立的结构基础。

1.1 在节点中构造并执行内层依赖图

以下示例来自 use_nested_flow_graphs.rst:外层图g有两个节点a和b,节点a收到消息后在 body 内新建一张依赖图h(四个continue_node风格的节点n1~n4,先扇出后扇入),节点b收到消息后则在 body 内新建一张数据流图:

graph g; function_node< int, int > a( g, unlimited, []( int i ) -> int { graph h; node_t n1( h, = { cout << "n1: " << i << "\n"; } ); node_t n2( h, = { cout << "n2: " << i << "\n"; } ); node_t n3( h, = { cout << "n3: " << i << "\n"; } ); node_t n4( h, = { cout << "n4: " << i << "\n"; } ); make_edge( n1, n2 ); make_edge( n1, n3 ); make_edge( n2, n4 ); make_edge( n3, n4 ); n1.try_put(continue_msg()); h.wait_for_all(); return i; } ); function_node< int, int > b( g, unlimited, []( int i ) -> int { graph h; function_node< int, int > m1( h, unlimited, []( int j ) -> int { cout << "m1: " << j << "\n"; return j; } ); function_node< int, int > m2( h, unlimited, []( int j ) -> int { cout << "m2: " << j << "\n"; return j; } ); function_node< int, int > m3( h, unlimited, []( int j ) -> int { cout << "m3: " << j << "\n"; return j; } ); function_node< int, int > m4( h, unlimited, []( int j ) -> int { cout << "m4: " << j << "\n"; return j; } ); make_edge( m1, m2 ); make_edge( m1, m3 ); make_edge( m2, m4 ); make_edge( m3, m4 ); m1.try_put(i); h.wait_for_all(); return i; } ); make_edge( a, b ); for ( int i = 0; i < 3; ++i ) { a.try_put(i); } g.wait_for_all();

关键执行流程:

  • 外层通过a.try_put(i)注入 3 个消息;make_edge(a, b)使a的输出作为b的输入。
  • 节点a的每次调用都会新建依赖图h、连边、以n1.try_put(continue_msg())启动,并调用h.wait_for_all()等待内层图空闲后返回。
  • 节点b的每次调用则新建一张由 4 个function_node< int, int >构成的数据流图,同样以m1.try_put(i)启动、h.wait_for_all()收尾。

这里wait_for_all()是必须的:因为h是节点 body 的局部变量,body 返回时h的析构函数会先调用wait_for_all()再销毁内部节点与上下文(见 flow_graph.h 中graph::~graph()的实现:wait_for_all(); if (own_context) { ... })。若不等待就返回,内层图任务会在析构过程中被强制等待,语义上等价但会把等待延迟到析构阶段,且无法保证h已经空闲后再复用,因此显式调用wait_for_all()是更清晰、可控的写法。


二、结构不变时的优化:持久复用内层图

如果内层图的结构在多次调用之间保持不变,那么每次调用都重新构造一次图就是多余的——make_edge、节点注册、arena 重新 attach 都会带来额外开销。此时可以把内层图提升为外层节点可捕获的持久对象,只在节点 body 中注入消息。

2.1 复用依赖图结构的改进版

将b的内层图提升到外层作用域,b的 body 通过引用捕获并直接try_put:

graph h; function_node< int, int > m1( h, unlimited, []( int j ) -> int { cout << "m1: " << j << "\n"; return j; } ); function_node< int, int > m2( h, unlimited, []( int j ) -> int { cout << "m2: " << j << "\n"; return j; } ); function_node< int, int > m3( h, unlimited, []( int j ) -> int { cout << "m3: " << j << "\n"; return j; } ); function_node< int, int > m4( h, unlimited, []( int j ) -> int { cout << "m4: " << j << "\n"; return j; } ); make_edge( m1, m2 ); make_edge( m1, m3 ); make_edge( m2, m4 ); make_edge( m3, m4 ); graph g; function_node< int, int > a( g, unlimited, []( int i ) -> int { graph h; node_t n1( h, = { cout << "n1: " << i << "\n"; } ); node_t n2( h, = { cout << "n2: " << i << "\n"; } ); node_t n3( h, = { cout << "n3: " << i << "\n"; } ); node_t n4( h, = { cout << "n4: " << i << "\n"; } ); make_edge( n1, n2 ); make_edge( n1, n3 ); make_edge( n2, n4 ); make_edge( n3, n4 ); n1.try_put(continue_msg()); h.wait_for_all(); return i; } ); function_node< int, int > b( g, unlimited, & -> int { m1.try_put(i); h.wait_for_all(); // optional since h is not destroyed return i; } ); make_edge( a, b ); for ( int i = 0; i < 3; ++i ) { a.try_put(i); } g.wait_for_all();

对比第一版,注意两处差异:

  • b的 lambda 捕获方式由[=]变为[&],以便引用外层持久图h及其节点;
  • b不再在局部构造内层图,只执行m1.try_put(i)即可复用已建好的图。

2.2wait_for_all()何时可以省略

原文档明确指出:持久图场景下,h.wait_for_all()是可选的。原因如下:

  • 第一版中h是局部对象,离开作用域即析构,graph::~graph()内部会执行wait_for_all(),所以不等待是不行的;
  • 改进版中h是外层持久对象,b的 body 返回后h并不会被销毁。因此若业务允许b不阻塞等待内层图完成,可以直接m1.try_put(i); return i;,让内层图在外层图的调度下继续运行。

不过要注意权衡:省略等待意味着b返回时内层消息可能尚未处理完,若后续逻辑依赖内层结果,则仍需调用h.wait_for_all();反之,若希望最大化重叠执行(overlap),省略等待能让b的调用立即返回,内层图与外层后续节点并行推进。

从实现上看,graph::wait_for_all()的核心是my_task_arena->execute(... d1::wait(my_wait_context_vertex.get_context(), *my_context)),等待线程在阻塞期间会去窃取任务工作("The waiting thread will go off and steal work while it is blocked in the wait_for_all",见 _flow_graph_impl.h),因此无论是否等待,等待线程本身都不会空转。


三、嵌套图的并发与优先级:源码与测试中的佐证

3.1 内层图与 task_arena

每个graph构造时会prepare_task_arena():优先 attach 到当前活动的task_arena,失败时新建一个默认初始化的 arena(见 _flow_graph_impl.h)。因此内层图默认运行在创建它的线程当前所在的 arena上。若需要让内层图运行在特定 arena,可以在构造内层图前用task_arena::execute切换,或在图生命周期跨越多个execute调用时通过graph::reset()重新 attach(见 flow_graph.h 中prepare_task_arena(/*reinit=*/true)的注释说明)。

3.2 测试用例:NestedCase

仓库测试 test_flow_graph_priorities.cpp 中的NestedCase直接验证了嵌套图的正确性,其要点:

  • 外层图outer_graph上挂 10 个function_node<int,int>,彼此全连接(任意两节点之间都建边);
  • 每个外层节点 body(OuterBody::operator())内部新建graph inner_graph,挂 4 个continue_node<continue_msg>(start_node、mid_node1、mid_node2、end_node),并组成扇出扇入结构;
  • 内层节点可以携带优先级:mid_node1为node_priority_t(5)、end_node为node_priority_t(15),证明嵌套图与节点优先级机制兼容;
  • 测试通过task_arena控制内外层图是否运行在同一个 arena(same_arena与different_arena两种配置,见test_in_arena的INFO输出),并遍历不同线程数验证;
  • 每次运行前通过outer_graph.reset()、inner_graph.reset()重置图状态,说明嵌套图可安全地反复启动(reset 会依次重置 context、所有注册节点,并重新 attach arena,见 flow_graph.h)。

对应的测试注册在 test_flow_graph_priorities.cpp 的TEST_CASE("Nested test case"),覆盖"内层图在外层节点体内反复构造执行"与"内外层图分处不同 arena"两类场景,可作为嵌套流图实现的参考基准。


四、实践要点与注意事项

  1. 局部图必须等待:若内层graph是节点 body 内的局部对象,离开作用域时析构会强制执行wait_for_all(),请在返回前显式调用h.wait_for_all(),语义更清晰,也避免析构阶段的隐式等待。
  2. 持久图按需等待:若内层图持久存在于外层作用域,h.wait_for_all()可以省略,以实现内外层图的异步重叠执行;需要确定性结果时再等待。
  3. 捕获方式:持久复用要求节点 lambda 以引用捕获[&]外层图对象;临时构造则可用值捕获[=]。
  4. 上下文隔离:每个graph拥有独立task_group_context(FLOW_TASKS类型),嵌套图之间取消与异常状态互不干扰;wait_for_all()返回前会同步cancelled/caught_exception状态。
  5. arena 归属:内层图默认 attach 到当前线程所在 arena;需要指定 arena 时结合task_arena::execute构造内层图,或使用graph::reset()重新 attach。
  6. 性能取舍:仅当内层图结构在多次调用间不变时才值得持久复用;若每次调用结构都不同,临时构造是唯一选择,此时应控制单次构造的节点规模,避免构造开销淹没图本身的执行时间。

五、延伸阅读

  • 图对象与wait_for_all/reserve_wait/release_wait的完整语义:flow_graph.h
  • graph::reset与 arena 重 attach 的实现:flow_graph.h
  • 嵌套图与节点优先级组合的测试:test_flow_graph_priorities.cpp
  • 流图基础概念与图对象说明:Flow_Graph.rst、Graph_Object.rst
  • 将嵌套流图与算法嵌套(parallel_for等)结合的组合模式:Parallelizing_Flow_Graph.rst
  • 并发编程
  • 高性能计算

【免费下载链接】oneTBB

oneAPI Threading Building Blocks (oneTBB)

项目地址:https://gitcode.com/gh_mirrors/on/oneTBB
点击查看免费下载
上一篇:Android抽屉动画效果终极指南:MaterialDrawer转场动画实现详解
下一篇:NetSonar高级技巧:如何自定义ping服务与导出网络性能报告

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

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

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

立即咨询