news 2026/9/15 19:19:36

oneTBB Flow Graph 实用技巧全解:等待、建边、嵌套并行、资源限制与异常取消的完整实践指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
oneTBB Flow Graph 实用技巧全解:等待、建边、嵌套并行、资源限制与异常取消的完整实践指南

oneTBB Flow Graph 实用技巧全解:等待、建边、嵌套并行、资源限制与异常取消的完整实践指南

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

导读

本文系统梳理 oneTBB(Intel oneAPI Threading Building Blocks)Flow Graph 编程中最容易踩坑的五大类实战问题:图的等待与销毁、边的建立与拓扑约定、嵌套并行、资源消耗限制、异常处理与取消。针对每一类问题,本文给出明确的规则、参数语义与可直接复用的代码示例,并辅以本仓库(mold 链接器,内部大量使用 TBB 并行原语)源码作为并行运行时的背景印证。读完本文,你将掌握wait_for_allmake_edgelimiter_noderejecting缓冲策略、reserving令牌系统、graph::reset()等核心 API 的正确用法,能够写出稳定、无数据竞争、资源可控的 Flow Graph 程序。

一、背景:Flow Graph 与 TBB 并行运行时

oneTBB 的 Flow Graph 提供了一种消息驱动(message-driven)的并行编程模型:程序以“节点 + 边”的方式描述数据流/依赖关系,运行时负责将图映射为任务并调度执行。它特别适合表达粗粒度并行、流水线与动态数据流结构,相关概念参见 Basic_Flow_Graph_concepts.rst、Flow_Graph.rst 与 Data_Flow_Graph.rst。

本仓库(mold)作为一款现代 C++ 链接器,正是以 TBB 作为并行运行时:例如 gc-sections.cc 使用tbb::concurrent_unordered_maptbb::concurrent_vectortbb::parallel_for_each,cmdline.cc 使用tbb::global_control,arch-arm32.cc 使用tbb::parallel_fortbb::parallel_for_each。因此,掌握 Flow Graph(以及 TBB 任务、arena 等机制)的实战技巧,对阅读和开发这类重度并行应用具有直接价值。oneTBB 官方用户指南在 tbb_userguide 目录 中收录了 Flow Graph 的完整专题,其中 Flow_Graph_Tips.rst 汇总了五组“Tips and Tricks”,本文即围绕这五组主题展开:

  1. 等待与销毁 Flow Graph;
  2. 建立边(拓扑);
  3. 嵌套并行;
  4. 限制资源消耗;
  5. 异常处理与取消。

二、等待与销毁 Flow Graph 的技巧

对应文档:Flow-Graph-waiting-tips.rst,包含三条规则:始终调用wait_for_all、避免动态删除节点、在主线程之外销毁图。

2.1 始终调用 wait_for_all()

Flow Graph 编程中最常见的错误是忘记调用wait_for_all()graph::wait_for_all会阻塞,直到图派生出的所有任务全部完成。这不仅用于等待计算结束,在销毁图或其任何节点之前必须调用它。下面这个函数会直接导致程序失败:

void no_wait_for_all() { graph g; function_node< int, int > f( g, 1, []( int i ) -> int { return spin_for(i); } ); f.try_put(1); // program will fail when f and g are destroyed at the // end of the scope, since the body of f is not complete }

原因在于:函数作用域结束时g与其节点f被销毁,但执行f体的任务仍在运行(in flight)。任务完成后会回头查找其节点的后继(successor),而此时图和节点都已被删除,从而产生悬垂访问。在函数末尾加上g.wait_for_all()即可避免图的提前销毁:

void wait_for_all() { graph g; function_node< int, int > f( g, 1, []( int i ) -> int { return spin_for(i); } ); f.try_put(1); g.wait_for_all(); // 防止 f 与 g 被过早销毁 }

排查建议:如果使用 Flow Graph 时出现诡异行为,首先检查是否调用了wait_for_all

2.2 避免动态删除节点

关于节点与边的基本准则:

  • 避免动态删除节点(dynamic node removal);
  • 允许在运行中新增边和节点;
  • 允许在运行中删除边。

在消息仍在图中处理时动态删除节点是被明确劝阻的:销毁图或其节点可能导致内存被提前释放,而图中仍在运行的任务稍后会访问这些内存,造成程序失败。当图尚未空闲(idle)时删除节点可能引发间歇性、极难排查的故障。详细讨论见 avoid_dynamic_node_removal.rst。

2.3 在主线程之外运行与销毁图

有时你不想用wait_for_all()阻塞主应用线程,但销毁图之前最安全的做法仍是调用它。常见方案是把一个“构建图 + 等待图完成”的任务入队task_arena

class background_task { public: void operator()() { graph g; function_node< int, int > f( g, 1, []( int i ) -> int { return spin_for(i); } ); f.try_put(1); g.wait_for_all(); } }; void no_wait_for_all_enqueue() { task_arena a; a.enqueue(background_task()); // do other things without waiting… }

注意:入队任务何时执行是不确定的。如果需要使用该任务的执行结果,或需要确保它在程序结束前完成,就必须使用某种同步机制(如条件变量、future 等)从入队任务向外发出“图已完成”的信号。详见 destroy_graphs_outside_main_thread.rst。

三、建立边的技巧

对应文档:Flow_Graph_making_edges_tips.rst,涵盖边 API 约定、单后继/广播语义、图间通信、input_node激活时机与数据竞争防护。

3.1 使用 make_edge 与 remove_edge

创建与移除边的规范约定:

  • 使用flow::make_edgeflow::remove_edge
  • 避免直接调用register_successor/register_predecessor
  • 避免直接调用remove_successor/remove_predecessor

运行时库内部通过节点函数(如sender<T>::register_successor)实现这些边,但用户不应直接调用——这些底层函数被库用于在运行时对拓扑做优化。统一使用make_edge/remove_edge表达拓扑即可。详见 use_make_edge.rst。边的完整语义(单边/广播、缓冲、预留)可参考 Edges.rst 与 Flow_Graph_Single_Vs_Broadcast.rst。

3.2 理解单后继推送与广播:buffer 类节点的特殊语义

预定义节点的一个重要特性是推送到单个后继,还是广播给所有后继。以下节点将消息推送给单个后继

  • buffer_node
  • queue_node
  • priority_queue_node
  • sequencer_node

其他节点则把消息推送给所有愿意接收(accept)的后继。单后继推送的节点都属于 buffer 类节点,其作用是暂存消息直到下游消费。

考虑下面用priority_queue_node连接两个function_node的示例:每个被缓冲的消息只会被发送给f1f2之一(按优先级挑一个),而不会同时发给两者:

void use_buffer_and_two_nodes() { graph g; function_node< int, int, rejecting > f1( g, 1, []( int i ) -> int { spin_for(0.1); cout << "f1 consuming " << i << "\n"; return i; } ); function_node< int, int, rejecting > f2( g, 1, []( int i ) -> int { spin_for(0.2); cout << "f2 consuming " << i << "\n"; return i; } ); priority_queue_node< int > q(g); make_edge( q, f1 ); make_edge( q, f2 ); for ( int i = 10; i > 0; --i ) { q.try_put( i ); } g.wait_for_all(); }

要点:function_node默认在输入端排队缓冲消息。为了让priority_queue_node正常工作,这里把function_node的缓冲策略设为rejecting(第三个模板参数),使其内部不缓冲,而是依赖上游priority_queue_node的缓冲。

如果 buffer 类节点改为广播给所有后继,会引发一连串难题:某个后继接受而另一个拒绝时消息是否保留?保留的话按什么顺序派发(原优先级序还是接受时刻的优先级)?例如priority_queue_node中只有 "9",f1接受 "9" 而f2拒绝;随后 "100" 到达且f2可接受——f2该收到 "9" 还是优先级更高的 "100"?总之,保证所有后继都收到每条消息会带来垃圾回收难题并使推理复杂化。因此这些缓冲节点把每条消息推送给唯一一个后继。利用这一特性即可构建“每条消息按优先级顺序由 f1 或 f2 处理”的图结构。

如果确实希望f1f2按优先级收到全部值,可以各配一个priority_queue_node,再用broadcast_node把每个值同时推给两个队列:

graph g; function_node< int, int, rejecting > f1( g, 1, []( int i ) -> int { spin_for(0.1); cout << "f1 consuming " << i << "\n"; return i; } ); function_node< int, int, rejecting > f2( g, 1, []( int i ) -> int { spin_for(0.2); cout << "f2 consuming " << i << "\n"; return i; } ); priority_queue_node< int > q1(g); priority_queue_node< int > q2(g); broadcast_node< int > b(g); make_edge( b, q1 ); make_edge( b, q2 ); make_edge( q1, f1 ); make_edge( q2, f2 ); for ( int i = 10; i > 0; --i ) { b.try_put( i ); } g.wait_for_all();

所以,把节点连接到多个后继时,务必确认输出是广播到所有后继,还是只推送给单一后继。完整的消息传递协议可参考 Flow_Graph_Message_Passing_Protocol.rst。

3.3 图间通信:不要跨图建边

所有图节点的构造函数都需要一个graph对象引用。只在同一张图的节点之间建边是安全的——边向运行时表达的是拓扑关系;连接两张不同图的节点会让graph::wait_for_all、异常处理等整图操作难以推理。为优化性能,库可能在用户预期之外的时机调用节点的前驱/后继函数。

如果两张图必须通信,不要在它们之间建边,而应使用显式的try_put调用。这样运行时不会对两个节点的关系做任何假设,跨图边界的事件也更容易推理,但整图操作仍然棘手。例如:

graph g; function_node< int, int > n1( g, 1, [](int i) -> int { cout << "n1\n"; spin_for(i); return i; } ); function_node< int, int > n2( g, 1, [](int i) -> int { cout << "n2\n"; spin_for(i); return i; } ); make_edge( n1, n2 ); graph g2; function_node< int, int > m1( g2, 1, [](int i) -> int { cout << "m1\n"; spin_for(i); return i; } ); function_node< int, int > m2( g2, 1, & -> int { cout << "m2\n"; spin_for(i); n1.try_put(i); // 跨图通信:显式 try_put,而非边 return i; } ); make_edge( m1, m2 ); m1.try_put( 1 ); // The following call returns immediately: g.wait_for_all(); // The following call returns after m1 & m2 g2.wait_for_all(); // we reach here before n1 & n2 are finished // even though wait_for_all was called on both graphs

这里m1.try_put(1)触发m2m2又通过显式try_put把消息送入n1n1再发给n2。由于不存在边,运行时并不认为m2n1的前驱。要等待两张图的所有任务完成,必须对两张图都调用wait_for_all,且调用顺序很重要。上面(错误的)代码里,第一次g.wait_for_all()立即返回(此时g中尚无任务,任务都属于g2);g2.wait_for_all()只等m1m2完成,不会等n1n2。交换顺序即可得到正确结果:

g2.wait_for_all(); g.wait_for_all(); // all tasks are done

两张小图的交互尚且如此容易出错,两张更大(可能含环)的图就更难理解了。因此,跨图节点通信务必谨慎。详见 communicate_with_nodes.rst。

3.4 input_node:先建图、后激活

默认情况下,input_node非激活(inactive)状态构造:

template< typename Body > input_node( graph &g, Body body, bool is_active=true );

激活非激活的input_node需要调用其activate()

input_node< int > src( g, src_body(10), false ); // use it in calls to make_edge… src.activate();

所有input_node都以非激活状态构造,通常在整个 Flow Graph 建好之后再激活。例如:

make_edge( squarer, summer ); make_edge( cuber, summer ); input_node< int > src( g, src_body(10), false ); make_edge( src, squarer ); make_edge( src, cuber ); src.activate(); g.wait_for_all();

如果一开始就让input_node处于激活状态,它可能在src→squarer边连好后就立即向squarer发消息;等src→cuber边连上时,cuber只会收到后续消息,已错过先前的部分。

一般规则:让input_node保持非激活,待整图构建完毕再激活。缺点是构建与执行被串行化。

例外:如果图是 DAG 且每个input_node只有一个后继,可以构造完就激活,前提是按逆拓扑序建边——先建最深层的边,再逐步回退到最浅层。下面这个图即使src构造后立即激活也不会丢消息:

const int limit = 10; int count = 0; graph g; oneapi::tbb::flow::graph g; oneapi::tbb::flow::input_node<int> src( g, & -> int { if ( count < limit ) { return ++count; } fc.stop(); return {}; }); src.activate(); oneapi::tbb::flow::function_node<int,int> func1( g, 1, []( int i ) -> int { std::cout << i << "\n"; return i; } ); oneapi::tbb::flow::function_node<int,int> func2( g, 1, []( int i ) -> int { std::cout << i << "\n"; return i; } ); make_edge( func1, func2 ); make_edge( src, func1 ); g.wait_for_all();

安全的原因是func1→func2的边先于src→func1建立:若顺序颠倒,func1可能在func2挂接之前就产出消息并导致丢失。同时src只有一个后继;若src有多个后继,先挂接的后继可能收到后挂接后继收不到的消息。完整说明见 use_input_node.rst 与 Data_Flow_Graph.rst。

3.5 主动避免数据竞争

图中的边显式表达了希望库强制执行的数据依赖关系;function_nodemultifunction_node的并发限制(concurrency limit)限定了运行时允许的最大并发调用数。这些是库强制执行的边界,但库不会自动保护你免受数据竞争——必须自己显式地用这些机制防竞争。

例如下面的代码存在数据竞争:没有任何机制阻止对节点f所引用全局对象global_sum的并发访问:

graph g; int src_count = 1; int global_sum = 0; int limit = 100000; input_node< int > src( g, & -> int { if ( src_count <= limit ) { return src_count++; } else { fc.stop(); return int(); } } ); src.activate(); function_node< int, int > f( g, unlimited, & -> int { global_sum += i; // data race on global_sum return i; } ); make_edge( src, f ); g.wait_for_all(); cout << "global sum = " << global_sum << " and closed form = " << limit*(limit+1)/2 << "\n";

运行该示例时,由于数据竞争,global_sum大概率小于期望值。修复这个简单示例的办法之一是把f的并发度从unlimited改为1,强制f串行处理每个值。注意input_node也会更新全局变量src_count,但由于input_node总是串行执行,不存在竞争。详见 avoiding_data_races.rst。

四、嵌套并行技巧

对应文档:Flow_Graph_nested_parallelism_tips.rst。

4.1 在节点体内嵌套并行算法提升可扩展性

提升 Flow Graph 可扩展性的有力手段是在节点体内嵌套其他并行算法:把 Flow Graph 当作“协调语言”(coordination language),用图表达粗粒度并行,把更细粒度的并行嵌套在节点内部。

下面的示例创建了五个节点:input_nodematrix_source)从文件读取矩阵序列;两个function_noden1n2)对每个元素应用函数生成两个新矩阵;两个汇点节点(n1_sinkn2_sink)处理结果矩阵。n1n2的 lambda 内部用parallel_for并行地对矩阵元素应用函数:

graph g; input_node< double * > matrix_source( g, & -> double* { double *a = read_next_matrix(); if ( a ) { return a; } else { fc.stop(); return nullptr; } } ); function_node< double *, double * > n1( g, unlimited, & -> double * { double *b = new double[N]; parallel_for( 0, N, & { b[i] = f1(a[i]); } ); return b; } ); function_node< double *, double * > n2( g, unlimited, & -> double * { double *b = new double[N]; parallel_for( 0, N, & { b[i] = f2(a[i]); } ); return b; } ); function_node< double *, double * > n1_sink( g, unlimited, []( double *b ) -> double * { return consume_f1(b); } ); function_node< double *, double * > n2_sink( g, unlimited, []( double *b ) -> double * { return consume_f2(b); } ); make_edge( matrix_source, n1 ); make_edge( matrix_source, n2 ); make_edge( n1, n1_sink ); make_edge( n2, n2_sink ); matrix_source.activate(); g.wait_for_all();

其中read_next_matrixf1f2consume_f1consume_f2为外部提供的函数。详细说明见 use_nested_algorithms.rst。

4.2 嵌套 Flow Graph

除了在节点体内嵌套算法,还可以嵌套 Flow Graph。下面的外层图g有两个节点aba收到消息时构造并执行一个内层依赖图(dependence graph),b收到消息时构造并执行一个内层数据流图(data flow graph):

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();

优化建议:如果嵌套图在节点每次被调用时结构保持不变,每次都重新构造就是冗余开销。可让b复用一个跨调用持久存在的图:

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体内,只有在希望b阻塞至内层图完成时才需要在每次调用末尾调用h.wait_for_all()。第一个版本的b必须在每次调用末尾调用,因为图在作用域结束时被销毁;而在复用版本中,可以在m1.try_put(i)后直接返回而无需等待h空闲。详见 use_nested_flow_graphs.rst。

五、限制资源消耗的技巧

对应文档:Flow_Graph_resource_tips.rst。核心诉求:控制允许进入图中某部分的消息数量,或控制工作池中的最大任务数。oneTBB 提供多种机制:limiter_node、节点并发限制、基于令牌(token)的系统、把图绑定到指定task_arena

5.1 使用 limiter_node 限制流经某点的消息数

limiter_node的构造函数接收两个参数:

limiter_node( graph &g, size_t threshold );
  • 第一个参数:所属图对象的引用;
  • 第二个参数:允许通过的最大条数,超过后节点开始拒绝(reject)传入消息。

limiter_node内部维护一个“已放行消息”计数。当消息离开被控区域时,向limiter_nodedecrement 端口发送一条消息即可递减计数,从而允许更多消息通过。下面的示例中,input_node将生成M个大对象,但要求同时最多只有 3 个大对象到达function_node,避免input_node一次性生成全部M个对象:

graph g; int src_count = 0; int number_of_objects = 0; int max_objects = 3; input_node< big_object * > s( g, & -> big_object* { if ( src_count < M ) { big_object* v = new big_object(); ++src_count; return v; } else { fc.stop(); return nullptr; } } ); s.activate(); limiter_node< big_object * > l( g, max_objects ); function_node< big_object *, continue_msg > f( g, unlimited, []( big_object *v ) -> continue_msg { spin_for(1); delete v; return continue_msg(); } ); make_edge( l, f ); make_edge( f, l.decrement ); make_edge( s, l ); g.wait_for_all();

工作流程:limiter_node阈值为 3,内部计数达到 3 后开始拒绝传入消息;input_node看到消息被拒后停止调用其 body,并临时缓冲最后生成的值。function_node的输出(continue_msg)送回limiter_node的 decrement 端口,执行完毕后计数减一;计数低于阈值后消息重新从input_node流出。因此图中同时最多存在4个大对象:3 个已通过limiter_node,1 个缓冲在input_node中。详见 use_limiter_node.rst。

5.2 使用节点并发限制(rejecting 策略)

要控制单个节点的并发实例数,可用节点的并发限制;要让节点在达到并发上限后拒绝消息,需把它构造成rejecting节点。

function_node通过模板参数构造,第三个模板参数控制缓冲策略,默认是queueing

template < typename Input, typename Output = continue_msg, graph_buffer_policy = queueing > class function_node;
  • queueing 策略:达到并发上限的function_node仍接收传入消息,但内部缓冲;
  • rejecting 策略:达到并发上限后拒绝传入消息。

示例:用一个rejectingfunction_node放在input_node下游,限制图中同时存在的大对象数量:

graph g; int src_count = 0; int number_of_objects = 0; int max_objects = 3; input_node< big_object * > s( g, & -> big_object* { if ( src_count < M ) { big_object* v = new big_object(); ++src_count; return v; } else { fc.stop(); return nullptr; } } ); s.activate(); function_node< big_object *, continue_msg, rejecting > f( g, 3, []( big_object *v ) -> continue_msg { spin_for(1); delete v; return continue_msg(); } ); make_edge( s, f ); g.wait_for_all();

function_node同时最多处理 3 个大对象:并发运行 3 个实例时开始拒绝来自input_node的消息,input_node缓冲最后一个对象并暂时停止调用 body;当并发度下降后,function_nodeinput_node拉取新消息。因此同时最多存在 4 个大对象:3 个在function_node中,1 个缓冲在input_node中。详见 use_concurrency_limits.rst。

5.3 创建基于令牌(Token)的系统:reserving join_node

更灵活的限制方案是使用令牌:图中只有有限数量的令牌,消息必须与可用令牌配对才能进入图;消息被移出图时释放令牌,令牌再与新消息配对进入。oneTBB 的parallel_pipeline算法即依赖令牌系统。Flow Graph 接口没有显式令牌支持,但可以用join_node构造类似系统:

template<typename OutputTuple, graph_buffer_policy JP = queueing> class join_node;

join_node的缓冲策略(buffer policy)有以下三种:

  • queueing:输入按先入先出(FIFO)匹配——按到达顺序拼接成元组;
  • tag_matching:将具有匹配标签(tag)的输入拼接在一起;
  • reserving:内部不缓冲,仅在能从每个端口的上游源先“预留”(reserve)输入时才消费输入;若每个端口都能预留,则取得这些输入并拼成输出元组。

基于令牌的系统可用reservingjoin_node实现。下面的示例中,input_node生成M个大对象,buffer_node预填了 3 个令牌(token_t可以是任意类型,例如typedef int token_t;)。input_nodebuffer_node都连到reservingjoin_node:只有当join_node确认buffer_node端有可拉取的项时,才会从input_node拉取输入;input_node只有被join_node拉取时才生成输入:

graph g; int src_count = 0; int number_of_objects = 0; int max_objects = 3; input_node< big_object * > s( g, & -> big_object* { if ( src_count < M ) { big_object* v = new big_object(); ++src_count; return v; } else { fc.stop(); return nullptr; } } ); s.activate(); join_node< tuple_t, reserving > j(g); buffer_node< token_t > b(g); function_node< tuple_t, token_t > f( g, unlimited, []( const tuple_t &t ) -> token_t { spin_for(1); cout << get<1>(t) << "\n"; delete get<0>(t); return get<1>(t); } ); make_edge( s, input_port<0>(j) ); make_edge( b, input_port<1>(j) ); make_edge( j, f ); make_edge( f, b ); b.try_put( 1 ); b.try_put( 2 ); b.try_put( 3 ); g.wait_for_all();

代码中function_node把令牌返回给buffer_node:图中的这个让令牌被回收并重新与input_node的新输入配对。与前两节一样,图中最多存在 4 个大对象:3 个在function_node中,1 个缓冲在input_node中等待配对令牌。

令牌系统的灵活性体现在:

  • 没有预定义的token_t,可以使用任何类型做令牌,包括对象或数组指针;
  • 令牌不必是哑类型——可以把计算中必需的缓冲或对象本身当作令牌,例如直接用大对象自身作令牌,通过回到buffer_node的环构成大对象空闲链表(free list),免去反复分配/释放;
  • 除了用固定次数的try_put预填buffer_node,也可以挂一个生成令牌的input_node
  • 可以把function_node换成multifunction_node(每个输出端口可输出 0 个或多个消息),从而选择回收或不回收令牌、甚至产生更多令牌,动态增减图中允许的并发度
  • 由于令牌可在源处与输入配对,这种方式能跨整个图限制资源消耗。

详见 create_token_based_system.rst。join_node的预留协议细节可参考 Flow_Graph_Reservation.rst。

5.4 将 Flow Graph 绑定到指定的 task_arena

task_arena接口提供多种引导任务执行的机制:设置首选计算单元(core type)、限制计算单元子集、限制 arena 并发度等。某些场景下希望 Flow Graph 执行时也应用这些机制。

graph对象在构造时绑定到“构造线程占有一个槽位”的那个 arena。下面的示例把图绑定到“以性能最高的 core 类型为首选”的 arena:

std::vector<tbb::core_type_id> core_types = tbb::info::core_types(); tbb::task_arena arena( tbb::task_arena::constraints{}.set_core_type(core_types.back()) ); arena.execute( [&]() { graph g; function_node< int > f( g, unlimited, []( int ) { /*the most performant core type is defined as preferred.*/ } ); f.try_put(1); g.wait_for_all(); } );

graph也可以通过调用graph::reset()重新绑定到不同的task_arena:该函数会重新初始化并把图重新绑定到“graph::reset()方法执行所在的”arena。每当为图派生任务时,任务都会在图所绑定的 arena中生成,而不论发起线程位于哪个 arena:

graph g; function_node< int > f( g, unlimited, []( int ) { /*the most performant core type is defined as preferred.*/ } ); std::vector<tbb::core_type_id> core_types = tbb::info::core_types(); tbb::task_arena arena( tbb::task_arena::constraints{}.set_core_type(core_types.back()) ); arena.execute( [&]() { g.reset(); } ); f.try_put(1); g.wait_for_all();

完整可编译示例见 flow_graph_examples.cpp(含begin_attach_to_arena_1/begin_attach_to_arena_2两段标注)。进一步了解 arena 引导与工作隔离,可参考 Guiding_Task_Scheduler_Execution.rst 与 work_isolation.rst。

六、异常处理与取消技巧

对应文档:Flow-Graph-exception-tips.rst。图的执行可以被直接取消,也可以因异常传播出节点 body 而被取消;随后可选用graph::reset()重置图以便重新执行。

6.1 在抛异常的节点内部捕获异常

如果异常在节点 body 内部被捕获,执行如常继续;如果异常未在 body 内捕获、传播出节点 body,则图中所有节点的执行都会被取消,异常在graph::wait_for_all()的调用点被重新抛出。例如:

graph g; function_node< int, int > f1( g, 1, []( int i ) { return i; } ); function_node< int, int > f2( g, 1, []( const int i ) -> int { throw i; return i; } ); function_node< int, int > f3( g, 1, []( int i ) { return i; } ); make_edge( f1, f2 ); make_edge( f2, f3 ); f1.try_put(1); f1.try_put(2); g.wait_for_all();

第二个节点f2抛出未捕获异常,导致图执行被取消,异常在g.wait_for_all()处被重新抛出;此处若未处理,程序终止。可以在 body 内捕获并处理:

function_node< int, int > f2( g, 1, []( const int i ) -> int { try { throw i; } catch (int j) { cout << "Caught " << j << "\n"; } return i; } );

异常在 body 内被捕获处理后,对整个图的执行没有影响。也可以选择在wait_for_all调用点捕获:

try { g.wait_for_all(); } catch ( int j ) { cout << "Caught " << j << "\n"; }

此时图的执行被取消:对本例而言,输入 1 永远到不了f3,输入 2 永远到不了f2f3。详见 catching_exceptions.rst 与 Exceptions_and_Cancellation.rst。

6.2 显式取消图(不借助异常)

要在不抛异常的情况下取消图执行,可以为图使用显式的task_group_context,然后调用其cancel_group_execution()

task_group_context t; graph g(t); function_node< int, int > f1( g, 1, []( int i ) { return i; } ); function_node< int, int > f2( g, 1, []( const int i ) -> int { cout << "Begin " << i << "\n"; spin_for(0.2); cout << "End " << i << "\n"; return i; } ); function_node< int, int > f3( g, 1, []( int i ) { return i; } ); make_edge( f1, f2 ); make_edge( f2, f3 ); f1.try_put(1); f1.try_put(2); spin_for(0.1); t.cancel_group_execution(); g.wait_for_all();

取消语义:已经开始执行的节点会执行完毕,尚未开始的节点不会启动。本例中f2会对输入 1 打印 Begin 与 End,但不会收到输入 2。

也可以从节点 body 内获取所属的task_group_context并取消其图:

graph g; function_node< int, int > f1( g, 1, []( int i ) { return i; } ); function_node< int, int > f2( g, 1, []( const int i ) -> int { cout << "Begin " << i << "\n"; spin_for(0.2); cout << "End " << i << "\n"; task::self().group()->cancel_group_execution(); return i; } ); function_node< int, int > f3( g, 1, []( int i ) { return i; } ); make_edge( f1, f2 ); make_edge( f2, f3 ); f1.try_put(1); f1.try_put(2); g.wait_for_all();

即使构造时没有显式传入task_group_context,也能从节点 body 中拿到它。详见 cancel_a_graph.rst 与 Cancellation_Without_An_Exception.rst。

6.3 用 graph::reset() 重置被取消的图

当图因未处理异常或task_group_context被显式取消而停止执行时,图与节点可能处于不确定状态:缓冲中可能残留消息(如 6.2 节示例中的输入 2),此外执行期间的各类优化也可能让节点和边处于不确定状态。要重新执行/重启图,必须先reset()

try { g.wait_for_all(); } catch ( int j ) { cout << "Caught " << j << "\n"; // do something to fix the problem g.reset(); f1.try_put(1); f1.try_put(2); g.wait_for_all(); }

详见 use_graph_reset.rst。

6.4 嵌套并行的取消传播

嵌套并行是否被取消,取决于内层上下文是否绑定到外层上下文:若绑定,则会被级联取消;否则不会。

如果 Flow Graph 执行被取消(显式或异常导致),其节点内嵌套的并行算法或 Flow Graph 启动的任务可能被取消,也可能不被取消。与库中所有嵌套并行一样,取消关系由显式的task_group_context控制:若未给 Flow Graph 提供显式task_group_context,它默认创建**隔离(isolated)**的上下文。相关细节见 cancelling_nested_parallelism.rst、Cancellation_and_Nested_Parallelism.rst 与 task_group_thread_safety.rst。

七、速查:核心规则一览

主题核心规则
等待与销毁销毁图/节点前必调wait_for_all();禁止动态删除节点;主线程外运行的图用入队任务包裹构建与等待
建立边只用make_edge/remove_edge;buffer 类节点单后继推送;跨图通信用try_put不用边;input_node先建图后activate();库不防数据竞争,需自行用并发度/同步控制
嵌套并行节点体内嵌套parallel_for等算法提升可扩展性;嵌套图结构不变时复用以省去重建开销
资源限制limiter_node(阈值 + decrement 端口);function_noderejecting并发上限;reservingjoin_node构造令牌系统;graph绑定/reset()重绑定task_arena
异常与取消未捕获异常取消全图并在wait_for_all处重抛;task_group_context::cancel_group_execution()显式取消;重启前调用graph::reset();嵌套并行取消取决于上下文绑定

上述规则对应的原始文档均位于 tbb_userguide 目录,入口为 Flow_Graph_Tips.rst;预定义节点类型总览见 Predefined_Node_Types.rst,节点如何映射到任务可参考 Mapping_Nodes2Tasks.rst。在 mold 这类以 TBB 为并行运行时的项目中,遵循这些约定即可写出正确、可扩展、资源可控的并行程序。

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

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/15 19:19:22

如何快速修改AI短剧的镜头和替换包装?

如何快速修改AI短剧的镜头和替换包装&#xff1f;特种猫的做法是正片和包装分层管理&#xff1a;镜头只替换有问题的单个分镜&#xff0c;包装用独立模板轨道一键覆盖&#xff0c;不重跑整集。截至 2026 年&#xff0c;创作者普遍踩3个坑&#xff1a;改1个镜头要整集重新生成、…

作者头像 李华
网站建设 2026/9/15 19:17:47

Unity角色口型同步与眼神模拟:SALSA With RandomEyes实战调优指南

学过几年Unity动画&#xff0c;接手过不少数字人、对话NPC的项目&#xff0c;我敢说在“让角色开口说话”这件事上&#xff0c;最让我省心的方案就是SALSA With RandomEyes。这个插件从名字就能看出来&#xff0c;它干两件事&#xff1a;SALSA负责说话时的口型同步&#xff0c;…

作者头像 李华