Mold 仓库内嵌 Intel oneTBB Flow Graph 的 input_node 深度解析:接口语义、行为规则与源码级实现
2026/9/14 4:11:35 网站建设 项目流程

Mold 仓库内嵌 Intel oneTBB Flow Graph 的 input_node 深度解析:接口语义、行为规则与源码级实现

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

本文以 mold 仓库third-party/tbb子树中 TBB Flow Graph 规范文档input_node_cls.rst为主体,完整解读oneapi::tbb::flow::input_node这一“消息源”节点:它如何生成消息、如何广播到所有后继、单槽缓冲如何工作、fc.stop()如何终止数据流,并结合 flow_graph.h 的实际实现逐一印证规范中每一条成员函数语义,帮助读者在设计并行数据流图时正确使用与调试该类节点。

1. input_node 在 Flow Graph 中的定位

input_nodeoneapi::tbb::flow::graph数据流图中的源头节点(source node)。规范文档(input_node_cls.rst)对其一句话定义是:

一个通过调用用户提供的函数对象(functor)生成消息、并将结果广播(broadcast)到所有后继(successor)的节点。

它有三个结构性特征:

  1. 没有前驱(predecessor)input_node只能作为消息的生产者,不能作为接收者。从源码看,其input_type被显式定义为null_type(flow_graph.h):
// Input node has no input type typedef null_type input_type;
  1. 串行执行体:它从不会并发调用自身的body。规范明确写道:“It is a serial node and never calls itsbodyconcurrently.”(它是一个串行节点,绝不同步调用其body。)这一点与function_nodemultifunction_node等可按并发度并行调用 body 的节点形成对比。

  2. 单槽缓冲:规范说明 “This node can buffer a single item.” 如果某条消息没有后继接收,消息会被缓存,并在生成新消息之前先提供出去。源码中对应两个成员变量(flow_graph.h):

bool my_has_cached_item; output_type my_cached_item;

input_node同时继承graph_nodesender<Output>,即在图中既可被graph统一管理与重置,又可通过make_edge作为消息发送端连接后继。它的转发与缓冲策略为broadcast-push + buffering,这一点可以在 forwarding_and_buffering.rst 的策略汇总表(try_get? 列为 yes,Forwarding 列为 broadcast-push)中得到印证。

2. 类声明与模板约束

2.1 规范给出的类声明

规范文档给出的接口声明如下(定义于头文件<oneapi/tbb/flow_graph.h>,当前仓库对应 flow_graph.h):

// Defined in header <oneapi/tbb/flow_graph.h> namespace oneapi { namespace tbb { namespace flow { template < typename Output > class input_node : public graph_node, public sender<Output> { public: template< typename Body > input_node( graph &g, Body body ); input_node( const input_node &src ); ~input_node(); void activate(); bool try_get( Output &v ); }; } // namespace flow } // namespace tbb } // namespace oneapi

2.2 Output 类型要求

规范要求模板参数Output满足 ISO C++ 标准的三项类型要求:

  • DefaultConstructible(默认可构造)——因为节点内部需要一个output_type my_cached_item缓存槽,默认构造是前提;
  • CopyConstructible(可拷贝构造)——缓存消息、向try_get的形参拷贝都依赖拷贝构造;
  • CopyAssignable(可拷贝赋值)——my_cached_item = (*my_body)(control)等赋值操作依赖拷贝赋值。

当前实现直接用概念(concept)约束了这一点(flow_graph.h):

template < typename Output > __TBB_requires(std::copyable<Output>) class input_node : public graph_node, public sender< Output > {

2.3 Body 类型要求(InputNodeBody)

Body必须满足 InputNodeBody 命名要求,其核心是一条伪签名:

Output Body::operator()( oneapi::tbb::flow_control& fc )

语义要求有三:

  • Body必须可拷贝构造、有析构函数;
  • operator()的返回类型必须与input_node实例的模板参数Output相同;
  • 当无法再生成新元素时,必须调用fc.stop()通知流结束。由于返回值必须是合法的Output,在终止时 Body 可以返回任意合法值,该值会被立即丢弃(discard)。

flow_control的实现非常轻量(_pipeline_filters.h):

class flow_control { bool is_pipeline_stopped = false; // ... public: void stop() { is_pipeline_stopped = true; } };

fc.stop()只是把一个标志位置位;节点在调用完 body 后检查control.is_pipeline_stopped来决定是否继续生成(见第 5 节)。当前头文件把这一要求写成了 C++20 概念(flow_graph.h):

template <typename Body, typename Output> concept input_node_body = std::copy_constructible<Body> && requires( Body& body, tbb::detail::d1::flow_control& fc ) { { body(fc) } -> adaptive_same_as<Output>; };

构造函数上用__TBB_requires(input_node_body<Body, Output>)做编译期检查(flow_graph.h)。

3. 成员函数语义详解

3.1 构造函数input_node( graph &g, Body body )

规范:构造一个调用bodyinput_node节点默认处于非激活(inactive)状态,即在调用activate()之前不会生成任何消息。

源码实现(flow_graph.h):

template< typename Body > __TBB_requires(input_node_body<Body, Output>) __TBB_NOINLINE_SYM input_node( graph &g, Body body ) : graph_node(g), my_active(false) , my_body( new input_body_leaf< output_type, Body>(body) ) , my_init_body( new input_body_leaf< output_type, Body>(body) ) , my_successors(this), my_reserved(false), my_has_cached_item(false) { fgt_node_with_body(CODEPTR(), FLOW_INPUT_NODE, &this->my_graph, static_cast<sender<output_type> *>(this), this->my_body); }

这里有三个值得注意的实现细节:

  • my_active(false):印证了“默认非激活”的规范表述,激活只能靠activate()
  • body 被拷贝了两次my_body是运行时真正被调用的那份,my_init_body保存构造时的初始副本,用于后续graph::reset恢复初始状态(见 reset_node 实现);
  • body 被包装进input_body_leaf<Output, Body>(_flow_graph_body_impl.h),以虚函数多态形式支持clone()get_body()
template< typename Output, typename Body> class input_body_leaf : public input_body<Output> { public: input_body_leaf( const Body &_body ) : body(_body) { } Output operator()(d1::flow_control& fc) override { return body(fc); } input_body_leaf* clone() override { return new input_body_leaf< Output, Body >(body); } Body get_body() { return body; } private: Body body; };

3.2 拷贝构造函数input_node( const input_node &src )

规范逐条给出了拷贝语义:

  • 新节点与src构造时的初始状态相同
  • 新节点引用与src相同的graph对象;
  • 持有src所用初始 body 的拷贝,且拥有与src相同的初始激活状态
  • src的后继不会被拷贝
  • 新 body 是从src构造时提供的原始 body 的副本再拷贝构造而来,src构造之后其 body 成员变量的变化不会影响新节点。

源码实现(flow_graph.h)与之一致:

__TBB_NOINLINE_SYM input_node( const input_node& src ) : graph_node(src.my_graph), sender<Output>() , my_active(false) , my_body(src.my_init_body->clone()), my_init_body(src.my_init_body->clone()) , my_successors(this), my_reserved(false), my_has_cached_item(false) { ... }

注意它克隆的是src.my_init_body而非src.my_body——这从源码结构上印证了规范中“以构造时的初始 body 为基准”的语义;my_successors(this)则是全新的空后继集合。

3.3void activate()

规范:将节点置为激活状态,从而使能消息生成

实现(flow_graph.h):

//! Activates a node that was created in the inactive state void activate() { spin_mutex::scoped_lock lock(my_mutex); my_active = true; if (!my_successors.empty()) spawn_put(); }

激活时若已有后继,会立即派生一个生成任务(spawn_put),启动消息流水线。

3.4bool try_get( Output &v )

规范:若缓冲中有消息,则拷贝到v;否则(当节点处于激活状态时)调用body尝试生成一条新消息并拷贝到v返回true表示消息已拷贝到v,否则返回false

实现(flow_graph.h):

//! Request an item from the node bool try_get( output_type &v ) override { spin_mutex::scoped_lock lock(my_mutex); if ( my_reserved ) return false; if ( my_has_cached_item ) { v = my_cached_item; my_has_cached_item = false; return true; } // we've been asked to provide an item, but we have none. enqueue a task to // provide one. if ( my_active ) spawn_put(); return false; }

从源码结构看,try_getsender<Output>接口的一部分(虚函数override),它并不在本线程直接执行 body,而是通过spawn_put()把生成任务提交回图的任务系统(spawn_in_graph_arena),这保证了 body 始终在图的调度域内串行执行。此外节点还提供try_reserve/try_release/try_consume的预留协议(flow_graph.h),供后继以“预留-确认”方式取消息,my_reserved标志防止同一缓存项被重复取用。

4. 消息生成循环与 fc.stop() 终止条件

规范文档给出了input_node最核心的行为规则:

  1. input_node持续调用body并广播消息,直到body内部调用fc.stop(),或节点没有有效后继为止;
  2. 一条消息生成后可能被所有后继拒绝。此时消息被缓冲,在“后继被加入”或“try_get被调用”之后作为下一条消息发出;
  3. 只有当缓冲为空时才会再次调用body

源码中这条循环由try_reserve_apply_body承担“生成 + 入缓存”职责(flow_graph.h):

bool try_reserve_apply_body(output_type &v) { spin_mutex::scoped_lock lock(my_mutex); if ( my_reserved ) { return false; } if ( !my_has_cached_item ) { d1::flow_control control; fgt_begin_body( my_body ); my_cached_item = (*my_body)(control); my_has_cached_item = !control.is_pipeline_stopped; fgt_end_body( my_body ); } if ( my_has_cached_item ) { v = my_cached_item; my_reserved = true; return true; } else { return false; } }

关键逻辑是my_has_cached_item = !control.is_pipeline_stoppedfc.stop()之后缓存被判定为空,apply_body_bypass返回nullptr,任务链自然终止——这正对应规范中“直到 body 调用fc.stop()”的终止条件。

任务驱动侧由apply_body_bypassinput_node_task_bypass完成(flow_graph.h、_flow_graph_body_impl.h):

//! Applies the body. Returning SUCCESSFULLY_ENQUEUED okay; forward_task_bypass will handle it. graph_task* apply_body_bypass( ) { output_type v; if ( !try_reserve_apply_body(v) ) return nullptr; graph_task *last_task = my_successors.try_put_task(v); // 广播-推送给所有后继 if ( last_task ) try_consume(); // 至少一个后继接收:消耗缓存 else try_release(); // 无后继接收:保留缓存,等待后续再发 return last_task; }

my_successorsbroadcast_cache<output_type>,其try_put_task实现“推送给所有愿意接收的后继”的 broadcast-push 语义;若全部拒绝则返回空任务,缓存项保留(try_release),后续register_successor(有新后继加入时)或try_get会再次触发spawn_put重发——与规范第 2 条规则逐句对应。此外spawn_put中还有is_graph_active(this->my_graph)检查,图被取消时不再派生任务。

5. Body 拷贝语义与 copy_body 观测

规范特别提醒:传给input_node的 body 对象会被拷贝。对 body 内部成员变量的更新不会影响构造节点时使用的原始对象;如果需要从节点外部查看 body 的最新状态,可以使用copy_body函数取回一份更新的拷贝。

copy_body的接口定义见 copy_body_func.rst,实现见 flow_graph.h 与 input_node 的 copy_function_object:

template<typename Body> Body copy_function_object() { input_body<output_type> &body_ref = *this->my_body; return dynamic_cast< input_body_leaf<output_type, Body> & >(body_ref).get_body(); } template< typename Body, typename Node > Body copy_body( Node &n ) { return n.template copy_function_object<Body>(); }

copy_body<MyBody>(src_node)返回节点内部当前body 副本,而非构造时的初始副本。一个实用含义是:像src_body中维护的my_next_value这类生成进度,应该通过copy_body查询,而不是看外部那个已经“过期”的原始对象。

6. 类型推导指引(Deduction Guides)

规范给出了 C++17 类模板推导指引,允许省略显式模板实参:

template <typename Body> input_node(graph&, Body) -> input_node<std::decay_t<input_t<Body>>>;

其中input_t是指向Body输入参数类型(即 body 的operator()返回值)的别名。于是input_node src(g, src_body(10))会被推导为对应输出类型的input_node,无需写input_node<int> src(g, src_body(10))

7. 实战示例:正确组织激活时机

规范强调默认非激活的设计意图,配合 TBB 用户指南 use_input_node.rst 与 Data_Flow_Graph.rst,可以总结出两种工程实践。

7.1 基础用法:整图构建完成后统一激活

用户指南给出的经典范式是:先构造全部节点与边,最后再activate()

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是 broadcast-push 节点,若在只有第一条边(src -> squarer)时就激活,消息会立即发给squarer;之后才连上的cuber只能收到“之后”的消息,早期消息会缺失。因此最稳妥的做法是节点先处于非激活状态,待全图构建完毕再激活

body 的完整形态(来自 Data_Flow_Graph.rst 的最终示例)展示了“生成直至fc.stop()”的写法:

class src_body { const int my_limit; int my_next_value; public: src_body(int l) : my_limit(l), my_next_value(1) {} int operator()( oneapi::tbb::flow_control& fc ) { if ( my_next_value <= my_limit ) { return my_next_value++; } else { fc.stop(); return int(); } } }; int main() { int sum = 0; graph g; function_node< int, int > squarer( g, unlimited, [](const int &v) { return v*v; } ); function_node< int, int > cuber( g, unlimited, [](const int &v) { return v*v*v; } ); function_node< int, int > summer( g, 1, & -> int { return sum += v; } ); make_edge( squarer, summer ); make_edge( cuber, summer ); input_node< int > src( g, src_body(10) ); make_edge( src, squarer ); make_edge( src, cuber ); src.activate(); g.wait_for_all(); cout << "Sum is " << sum << "\n"; }

指南同时指出:这种“先建图后激活”的写法会把图的构建与执行串行化;input_node的优势在于它能响应下游行为(下游拒绝时消息留在单槽缓存里,而不是无限堆积),在复杂图中可以限制内存占用。

7.2 进阶用法:DAG 反拓扑序连边,允许边构建边执行

用户指南还给出了允许构建/执行重叠的受限场景:如果图是有向无环图(DAG)且每个input_node只有一个后继,可以按反拓扑序连边(先连深度最大的边,再向源端回溯),并在构造后立即激活而不丢消息:

const int limit = 10; int count = 0; 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生成的消息到达时下游已就位,不会被丢弃;其二,src只有单个后继,避免了“先接上的后继独占早期消息”的 broadcast 不公平问题。若这两个前提不满足,仍应回退到 7.1 的“全图就绪后激活”范式。

8. 生命周期补充:reset 与 body 恢复

规范正文未展开、但源码结构可以直接确认的一点是input_nodegraph::reset的支持。其reset_node实现(flow_graph.h):

//! resets the input_node to its initial state void reset_node( reset_flags f) override { my_active = false; my_reserved = false; my_has_cached_item = false; if(f & rf_clear_edges) my_successors.clear(); if(f & rf_reset_bodies) { input_body<output_type> *tmp = my_init_body->clone(); delete my_body; my_body = tmp; } }

从源码结构看:重置后节点回到非激活态、清空缓存与预留标志;rf_clear_edges会清除后继边,rf_reset_bodies则利用构造时保存的my_init_body克隆体,把被 body 执行过程修改过的成员状态恢复为构造初始值——这正是第 3.1 节中 body 被保存两份副本的意义所在。

9. 小结:input_node 的完整心智模型

综合规范文档与源码证据,可以把input_node归纳为一个状态机:

状态/事件行为依据
构造后非激活,my_active=false,缓存空input_node_cls.rst、flow_graph.h#L662
activate()置激活;已有后继则立即派生生成任务flow_graph.h#L781-L786
每次生成仅当缓存为空时调用一次 body(串行,不并发)规范正文、try_reserve_apply_body
消息广播推送给所有可接收的后继(broadcast-push)forwarding_and_buffering.rst
全部拒绝消息留在单槽缓存,待新后继加入或try_get再发apply_body_bypass
body 调用fc.stop()终止条件;此后不再调用 body_pipeline_filters.h#L134-L143
外部取消息try_get(Output&):命中缓存返回 true,否则(激活时)派生任务并返回 falseflow_graph.h#L712-L727
观测内部状态copy_body<Body>(node)返回当前 body 副本copy_body_func.rst、flow_graph.h#L2626-L2629

对使用 TBB Flow Graph 的开发者而言,掌握input_node的关键在于三点:理解单槽缓冲与 broadcast-push 策略决定了它在图中天然限流body 必须可拷贝且用fc.stop()优雅终止,其内部状态只能通过copy_body对外观测;激活时机要与建边顺序配合——要么全图就绪后统一激活,要么在满足 DAG + 单后继 + 反拓扑连边三个条件时才允许提前激活。

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

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

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

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

立即咨询