mold 内置 TBB Flow Graph 边(Edges)机制详解:make_edge 建边、remove_edge 拆边与消息传递协议
2026/9/15 17:30:19 网站建设 项目流程

mold 内置 TBB Flow Graph 边(Edges)机制详解:make_edge 建边、remove_edge 拆边与消息传递协议

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

本篇文章以 mold 仓库内置的 oneAPI TBB(Threading Building Blocks)官方用户指南中Edges(边)主题为核心,深入讲解 Flow Graph 中"边"的本质、make_edge/remove_edge的用法与底层实现,并结合并发限制、单后继/广播语义以及 push/pull 消息传递协议,帮助你从零搭建节点相连的数据流图并准确控制消息在节点间的流动方向。读完本文,你将能够独立使用边构建多节点流水线,理解边上的并发行为,并掌握拓扑变更的推荐做法。

mold 仓库在 third-party/tbb 目录下完整内置了 TBB 的源码与官方用户指南(RST 文档),其中 Flow Graph 是 TBB 面向数据流与依赖图并行编程的核心模型,而"边(Edges)"正是把各个节点串成一张图的关键粘合剂。

边的本质:连接节点的有向消息通道

大多数应用都包含多个节点,节点之间通过边(edge)彼此相连。在 Flow Graph 接口中,边是有向通道(directed channel),消息(message)沿着边从一个节点传递到另一个节点。边由调用函数make_edge(p, s)创建,其中:

  • p前驱节点(predecessor),即消息的发出方;
  • s后继节点(successor),即消息的接收方。

边是有向的:消息只会从p流向s,不会反向流动。这一点与 Nodes(节点) 主题中介绍的单节点模型形成了递进关系——单节点只能通过try_put手动喂入消息,而一旦引入边,消息便可以在运行时由库自动转发,无需开发者逐条手动投递。

需要特别注意的是,边是类型化的。从源码中make_edge的函数签名可以看到,它要求前驱必须是sender<T>、后继必须是receiver<T>,且两者消息类型T必须一致(见 flow_graph.h)。也就是说,只有输出类型与输入类型匹配的节点才能被边连接,编译器会在类型不匹配时直接报错。

双节点示例:从 n 到 m 的一条边

沿用 Nodes 主题 中的示例,我们新增一个节点m,它先对接收到的值做平方运算再打印,然后通过make_edge(n, m)把它与节点n相连:

graph g; function_node< int, int > n( g, unlimited, []( int v ) -> int { cout << v; spin_for( v ); cout << v; return v; } ); function_node< int, int > m( g, 1, []( int v ) -> int { v *= v; cout << v; spin_for( v ); cout << v; return v; } ); make_edge( n, m ); n.try_put( 1 ); n.try_put( 2 ); n.try_put( 3 ); g.wait_for_all();

其中spin_for是文档示例中用于模拟耗时计算的辅助函数(其实现未在指南中给出,可自行实现为一个忙等循环)。

这个例子包含了 Flow Graph 编程的四个基本动作:

  1. 构造图对象graph g;——所有节点与边都隶属于这张图;
  2. 构造节点——function_node<int, int>表示"一个输入、一个输出"的函数节点,模板两个参数分别是输入消息类型与输出消息类型;
  3. 建边——make_edge(n, m)建立从nm的有向边;
  4. 驱动与等待——n.try_put(1/2/3)向图注入三个消息,g.wait_for_all()阻塞直到所有消息处理完毕。

现在图中有两个function_nodenmmake_edge调用创建了从nm的一条边。节点nunlimited(无限)并发度创建,而m的并发度被限制为1。这意味着:

  • n的多次调用可以并行推进;
  • m的多次调用会被串行化(同一时刻最多只有一次调用在执行);
  • 因为存在从nm的边,n每次返回的值v都会被运行时库自动传递给节点m

这里unlimited指的是oneapi::tbb::flow::unlimited,即 TBB 预定义的"无限并发"常量(也可写作oneapi::tbb::flow::unlimited的完整限定形式,或直接使用 1、4、8 等具体数值)。

make_edge 的底层实现:注册后继者与运行时追踪

边并非某种独立于节点的"实体对象",其本质是节点之间的注册关系。查看 flow_graph.h 中make_edge的实现:

template< typename T > inline void internal_make_edge( sender<T> &p, receiver<T> &s ) { register_successor(p, s); fgt_make_edge( &p, &s ); } //! Makes an edge between a single predecessor and a single successor template< typename T > inline void make_edge( sender<T> &p, receiver<T> &s ) { internal_make_edge( p, s ); }

可以看到make_edge实际做了两件事:

  1. 调用register_successor(p, s):把后继节点s注册到前驱节点p的后继者列表里。在sender 基类中,register_successor是一个纯虚函数,由各具体节点类型实现;例如input_node在注册新后继后会立即调用spawn_put()尝试把已就绪的消息推给新后继(见 flow_graph.h)。这正是"建边之后消息立即开始流动"的底层原因。
  2. 调用fgt_make_edge(&p, &s):更新 Flow Graph 的运行时追踪/可视化信息。从源码结构看,fgt_前缀函数族(如fgt_make_edgefgt_remove_edge)用于记录图的拓扑变化,供性能剖析与调试工具使用。

值得注意的是,make_edge还提供了针对多输出端口/多输入端口节点的重载版本(见 flow_graph.h):

  • make_edge(T& output, V& input):连接多输出节点的端口 0与多输入节点的端口 0
  • make_edge(T& output, receiver<R>& input):连接多输出节点的端口 0 与普通接收节点;
  • make_edge(sender<S>& output, V& input):连接普通发送节点与多输入节点的端口 0。

也就是说,开发者不需要关心节点是否有多个端口,make_edge会自动选择"端口 0"完成最常见的一对一连接;更复杂的端口连接需求则通过节点自身的output_ports()/input_ports()接口配合std::get<N>手工指定。

并发控制与消息缓冲:边上的执行语义

function_node的构造签名是(见 Nodes 主题):

template< typename Body> function_node(graph &g, size_t concurrency, Body body)
参数说明
g节点所属的图对象
concurrency节点的并发度上限:从 1(完全串行)到unlimited(无限)之间的任意值
body用户定义的函数对象或 lambda 表达式,用于把输入消息转换为输出消息

并发度参数直接决定了边上消息被处理的节奏。以并发度 1 的节点为例,当它依次收到消息 1、2、3 时:

  1. 节点派生一个任务(task)处理第一个输入 1;
  2. 该任务完成后,再派生下一个任务处理 2;
  3. 同理,处理完 2 后才派生处理 3 的任务。

try_put调用本身不会阻塞等待任务派生:如果节点因并发度限制暂时无法立即派生任务处理消息,消息会被缓冲在节点内部;一旦并发条件允许,节点就会派生任务处理下一个被缓冲的消息。

如果把function_node的并发度改为unlimited,则库会在消息到达时立即派生任务,不管此前已经派生了多少任务。如果系统有足够多的线程,三个 body 调用将真正并行执行;如果系统只有一个线程,它们仍然串行执行。这一点引出了 Flow Graph 中一个重要的认知:

派生任务(spawn task)≠ 创建线程(create thread)。一张图可能派生大量任务,但真正执行这些任务的只有库线程池中可用数量的线程。

把这个语义套回双节点示例:n是 unlimited,三条消息的 body 调用可并行;m是并发度 1,从边上传来的所有消息被串行处理。于是"上游并行计算、下游串行汇聚"的流水线结构就通过一条边加两个并发度参数自然形成了。

删除边与拓扑操作约定

make_edge对称,TBB 提供remove_edge(p, s)用于拆除两个节点之间的边。其实现同样分两步:调用remove_successor(p, s)从前驱的后继者列表中移除该后继,并调用fgt_remove_edge更新追踪信息;多端口节点同样有对应的重载版本(见 flow_graph.h)。

关于边的增删,官方指南给出了明确的使用约定(见 Use make_edge and remove_edge):

  • 使用make_edgeremove_edge来表达图的拓扑;
  • 避免直接调用register_successorregister_predecessor
  • 避免直接调用remove_successorremove_predecessor

原因是:运行时库本身会直接调用这些节点级函数来实现拓扑的动态优化(例如前文input_node注册新后继后立即激活消息推送的逻辑)。如果应用代码也直接调用它们,就可能干扰库对拓扑的管理,产生难以排查的并发问题。因此,无论建边还是拆边,统一走flow::make_edge/flow::remove_edge这两个公开接口即可。

单后继还是广播:建边前必须了解的节点行为

边只是通道,消息沿边送出去之后的行为由节点类型决定。这是建边时最容易踩坑的地方。根据 Sending to One or Multiple Successors,预定义节点分两大类:

推送给单个后继(single successor)的节点:

  • buffer_node
  • queue_node
  • priority_queue_node
  • sequencer_node

这些缓冲类节点会把每条消息只推送给一个后继。它们存在的意义是临时保存消息,等待下游消费。考虑下面的例子:一个priority_queue_node同时连接两个function_node

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

这里有两个关键点:

  1. function_node默认会在输入侧缓冲消息;为了让priority_queue_node正确工作,示例把两个function_node的缓冲策略设为rejecting(拒绝型),使它们不内部缓冲,而是完全依赖上游priority_queue_node的缓冲;
  2. 每条被priority_queue_node缓冲的消息,会被送往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();

因此,把一个节点连到多个后继之前,务必确认它的输出语义是广播给所有后继,还是只推送给其中一个后继,否则图的行为会与直觉相悖。

边上的动态 push/pull 消息传递协议

边不只是"建好就固定"的静态通道,它在运行时会根据节点状态在push(推)pull(拉)两种模式之间动态切换(见 Flow Graph Basics: Message Passing Protocol)。

Flow Graph 通过节点间传递消息来运作,但后继节点可能暂时无法接收并处理来自前驱的消息。为了让图高效运行,当出现这种情况时,边的状态会从 push 反转为 pull:等后继有能力处理消息时,它可以主动向前驱查询是否有消息可用。如果没有这种反转机制,前驱节点就只能反复重试推送,白白浪费资源。

边处于 pull 模式时,后继一旦空闲就会尝试从前驱拉取消息:

  1. 如果前驱有消息,后继处理它,边保持 pull 模式
  2. 如果前驱没有消息,边从 pull 切回 push 模式

这个协议意味着:边本身携带了"前驱/后继当前是否就绪"的协商状态,消息的传递并非简单的"生产者死等消费者",而是推拉结合的自适应机制。理解这一点,有助于解释为什么make_edge之后节点无需任何手动同步——运行时库会沿着边自动完成消息的调度与背压(backpressure)处理。

小结:构建边时的检查清单

围绕"边"这一主题,将上述要点浓缩为一张可直接对照的清单:

关注点结论
建边/拆边接口只用flow::make_edge(p, s)flow::remove_edge(p, s),不要直接调用register_*/remove_*系列节点函数
边的方向与类型边是sender<T>receiver<T>的有向通道,两端消息类型必须一致;多端口节点默认连接端口 0
并发语义concurrency = 1串行处理、unlimited到达即处理、具体数值(如 4、8)限定并发上限;任务≠线程
消息缓冲try_put不阻塞;节点忙时会缓冲消息,具备处理能力后再派生任务
单后继 vs 广播buffer_node/queue_node/priority_queue_node/sequencer_node只推给单个后继;其余节点广播给所有可接受的后继
运行时行为边上运行 push/pull 动态协议:后继无法接收时边转为 pull,后继空闲后按需拉取,无消息时再切回 push

从单个节点的try_put到双节点的make_edge,再到多后继与广播场景,边始终是 Flow Graph 中"消息如何流动"的载体。掌握make_edge/remove_edge的用法与其背后的注册、追踪、并发与推拉协议机制,你就拥有了把任意预定义节点组合成高效并行流水线的核心能力。更进一步的内容,如节点类型总览、input_node的使用、数据竞争规避等,可继续阅读 Basic Flow Graph Concepts 目录 与 Flow Graph Tips 目录 下的其余主题。

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

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

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

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

立即咨询