
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_node是oneapi::tbb::flow::graph数据流图中的源头节点source node。规范文档input_node_cls.rst对其一句话定义是一个通过调用用户提供的函数对象functor生成消息、并将结果广播broadcast到所有后继successor的节点。它有三个结构性特征没有前驱predecessorinput_node只能作为消息的生产者不能作为接收者。从源码看其input_type被显式定义为null_typeflow_graph.h// Input node has no input type typedef null_type input_type;串行执行体它从不会并发调用自身的body。规范明确写道“It is a serial node and never calls itsbodyconcurrently.”它是一个串行节点绝不同步调用其body。这一点与function_node、multifunction_node等可按并发度并行调用 body 的节点形成对比。单槽缓冲规范说明 “This node can buffer a single item.” 如果某条消息没有后继接收消息会被缓存并在生成新消息之前先提供出去。源码中对应两个成员变量flow_graph.hbool my_has_cached_item; output_type my_cached_item;input_node同时继承graph_node与senderOutput即在图中既可被graph统一管理与重置又可通过make_edge作为消息发送端连接后继。它的转发与缓冲策略为broadcast-push buffering这一点可以在 forwarding_and_buffering.rst 的策略汇总表try_get? 列为 yesForwarding 列为 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 senderOutput { 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 oneapi2.2 Output 类型要求规范要求模板参数Output满足 ISO C 标准的三项类型要求DefaultConstructible默认可构造——因为节点内部需要一个output_type my_cached_item缓存槽默认构造是前提CopyConstructible可拷贝构造——缓存消息、向try_get的形参拷贝都依赖拷贝构造CopyAssignable可拷贝赋值——my_cached_item (*my_body)(control)等赋值操作依赖拷贝赋值。当前实现直接用概念concept约束了这一点flow_graph.htemplate typename Output __TBB_requires(std::copyableOutput) class input_node : public graph_node, public sender Output {2.3 Body 类型要求InputNodeBodyBody必须满足 InputNodeBody 命名要求其核心是一条伪签名Output Body::operator()( oneapi::tbb::flow_control fc )语义要求有三Body必须可拷贝构造、有析构函数operator()的返回类型必须与input_node实例的模板参数Output相同当无法再生成新元素时必须调用fc.stop()通知流结束。由于返回值必须是合法的Output在终止时 Body 可以返回任意合法值该值会被立即丢弃discard。flow_control的实现非常轻量_pipeline_filters.hclass flow_control { bool is_pipeline_stopped false; // ... public: void stop() { is_pipeline_stopped true; } };fc.stop()只是把一个标志位置位节点在调用完 body 后检查control.is_pipeline_stopped来决定是否继续生成见第 5 节。当前头文件把这一要求写成了 C20 概念flow_graph.htemplate typename Body, typename Output concept input_node_body std::copy_constructibleBody requires( Body body, tbb::detail::d1::flow_control fc ) { { body(fc) } - adaptive_same_asOutput; };构造函数上用__TBB_requires(input_node_bodyBody, Output)做编译期检查flow_graph.h。3. 成员函数语义详解3.1 构造函数input_node( graph g, Body body )规范构造一个调用body的input_node。节点默认处于非激活inactive状态即在调用activate()之前不会生成任何消息。源码实现flow_graph.htemplate typename Body __TBB_requires(input_node_bodyBody, 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_castsenderoutput_type *(this), this-my_body); }这里有三个值得注意的实现细节my_active(false)印证了“默认非激活”的规范表述激活只能靠activate()body 被拷贝了两次my_body是运行时真正被调用的那份my_init_body保存构造时的初始副本用于后续graph::reset恢复初始状态见 reset_node 实现body 被包装进input_body_leafOutput, Body_flow_graph_body_impl.h以虚函数多态形式支持clone()与get_body()template typename Output, typename Body class input_body_leaf : public input_bodyOutput { 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), senderOutput() , 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; } // weve been asked to provide an item, but we have none. enqueue a task to // provide one. if ( my_active ) spawn_put(); return false; }从源码结构看try_get是senderOutput接口的一部分虚函数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最核心的行为规则input_node会持续调用body并广播消息直到body内部调用fc.stop()或节点没有有效后继为止一条消息生成后可能被所有后继拒绝。此时消息被缓冲在“后继被加入”或“try_get被调用”之后作为下一条消息发出只有当缓冲为空时才会再次调用body。源码中这条循环由try_reserve_apply_body承担“生成 入缓存”职责flow_graph.hbool 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_bypass和input_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_successors是broadcast_cacheoutput_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_objecttemplatetypename Body Body copy_function_object() { input_bodyoutput_type body_ref *this-my_body; return dynamic_cast input_body_leafoutput_type, Body (body_ref).get_body(); } template typename Body, typename Node Body copy_body( Node n ) { return n.template copy_function_objectBody(); }即copy_bodyMyBody(src_node)返回节点内部当前body 副本而非构造时的初始副本。一个实用含义是像src_body中维护的my_next_value这类生成进度应该通过copy_body查询而不是看外部那个已经“过期”的原始对象。6. 类型推导指引Deduction Guides规范给出了 C17 类模板推导指引允许省略显式模板实参template typename Body input_node(graph, Body) - input_nodestd::decay_tinput_tBody;其中input_t是指向Body输入参数类型即 body 的operator()返回值的别名。于是input_node src(g, src_body(10))会被推导为对应输出类型的input_node无需写input_nodeint 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_nodeint src( g, - int { if ( count limit ) { return count; } fc.stop(); return {}; } ); src.activate(); oneapi::tbb::flow::function_nodeint,int func1( g, 1, []( int i ) - int { std::cout i \n; return i; } ); oneapi::tbb::flow::function_nodeint,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_node对graph::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_bodyoutput_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_activefalse缓存空input_node_cls.rst、flow_graph.h#L662activate()置激活已有后继则立即派生生成任务flow_graph.h#L781-L786每次生成仅当缓存为空时调用一次 body串行不并发规范正文、try_reserve_apply_body消息广播推送给所有可接收的后继broadcast-pushforwarding_and_buffering.rst全部拒绝消息留在单槽缓存待新后继加入或try_get再发apply_body_bypassbody 调用fc.stop()终止条件此后不再调用 body_pipeline_filters.h#L134-L143外部取消息try_get(Output)命中缓存返回 true否则激活时派生任务并返回 falseflow_graph.h#L712-L727观测内部状态copy_bodyBody(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),仅供参考