mold 项目中 oneTBB flow_graph 的 InputNodeBody 命名要求:input_node 数据源函数体的接口契约与实现解析
【免费下载链接】moldmold: A Modern Linker 🦠项目地址: https://gitcode.com/GitHub_Trending/mo/mold
导读:
InputNodeBody是 oneTBB(本仓库以third-party/tbb形式内置)flow_graph 中input_node节点对"数据源函数体(Body)"提出的形式化命名要求(Named Requirement),它规定了一个无输入、仅靠oneapi::tbb::flow_control&参数驱动、逐次产出Output类型消息的函数对象必须具备的接口与语义。本文以 InputNodeBody 规范文档 为主线,结合 flow_graph.h 中的input_node类、input_node_body概念约束以及 examples/graph 下的真实用例,系统讲解 Body 的接口要求、fc.stop()终止语义、源码中的调用链,并给出可直接运行的实战示例,帮助你正确编写并安全使用 input_node 的生成函数。
一、InputNodeBody 是什么:input_node 的"数据源契约"
在 oneTBB 的 flow_graph 中,input_node是一个特殊节点:它没有前驱(predecessor),是整张图的数据源头,负责按需生成消息并推送给后继节点。它的工作方式不是"被动等待输入",而是由用户提供一个生成函数体(Body),每当图需要更多数据时,节点就会调用一次该 Body 来产生下一项消息。
InputNodeBody就是 oneTBB 规范对这一函数体类型提出的形式化要求(对应规范编号[req.input_node_body],见 input_node_body.rst)。凡是要作为input_node的 Body 传入的类型,都必须满足这一组接口与语义约束,否则编译无法通过(或在使用 C++20 概念约束时直接报概念错误)。
其核心定义可以概括为三句话:
| 要求 | 伪签名 | 语义 |
|---|---|---|
| 拷贝构造 | Body::Body( const Body& ) | Body 必须可拷贝构造 |
| 析构 | Body::~Body() | Body 必须可析构 |
| 核心调用 | Output Body::operator()( oneapi::tbb::flow_control& fc ) | 生成下一项数据;无法继续生成时调用fc.stop() |
命名要求(Named Requirements)是 oneTBB 规范中描述"类型必须满足的接口与语义"的惯用形式,与 C++20 concept 一一对应。本仓库源码在 flow_graph.h 第 112-116 行 给出了等价的概念约束:
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>; };可以看到,规范文档中的三条要求与源码概念完全对齐:std::copy_constructible<Body>对应拷贝构造要求;body(fc)的可调用性对应operator()要求;返回值adaptive_same_as<Output>对应"返回类型必须与节点模板参数Output一致"的要求。
二、逐条解读三项接口要求
1. 拷贝构造函数:Body::Body( const Body& )
input_node在构造时会拷贝你传入的 Body(内部会创建一份"初始副本"作为重置时的模板),因此 Body 必须支持拷贝构造。从 flow_graph.h 第 661-669 行 可以看到,构造函数将传入的body复制成两份:
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_body:当前正在使用的 Body(每次调用operator()产生数据);my_init_body:初始状态的 Body 快照,用于图重置(reset())时恢复现场。
底层封装类定义在 _flow_graph_body_impl.h 第 76-91 行:
template< typename Output > class input_body : no_assign { public: virtual ~input_body() {} virtual Output operator()(d1::flow_control& fc) = 0; virtual input_body* clone() = 0; }; 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); } };input_node的拷贝构造函数同样依赖my_init_body->clone()来复制 Body(flow_graph.h 第 682-690 行)。因此一个实用建议是:Body 内部的可变状态(如计数器、迭代器、文件句柄)应当可以随拷贝独立工作,且拷贝不应产生共享的、会互相干扰的副作用;若状态不可拷贝,则应使用std::shared_ptr等间接手段(此时拷贝的是指针,仍需保证语义正确)。
2. 析构函数:Body::~Body()
要求很简单:Body 必须可析构。input_node的析构函数会delete my_body; delete my_init_body;(flow_graph.h 第 693 行),因此如果 Body 持有资源(如打开的文件、堆内存),请确保析构函数正确释放。
3. 核心调用运算符:Output Body::operator()( oneapi::tbb::flow_control& fc )
这是整个契约的核心,规范原文强调两点:
- 返回类型要求:
Output必须与构造该input_node时使用的模板类型参数Output完全一致(对应源码概念中的adaptive_same_as<Output>)。 - 生成与终止语义:每次调用时,Body 负责"生成下一项数据";当无法再生成新元素时,必须调用
fc.stop()通知节点数据流结束。
规范还特别指出一个容易被忽略的细节:由于Output必须被返回,即使停止生成,Body 也要返回一个合法的Output值——这个值会被节点立即丢弃,不会进入图中。也就是说,fc.stop()分支里的return只是语法上必须的占位,语义上不影响输出。
oneapi::tbb::flow_control本身是一个极简的哨兵类,定义在 _pipeline_filters.h 第 134-143 行:
class flow_control { bool is_pipeline_stopped = false; flow_control() = default; // ... 仅允许 pipeline 与 input_node 内部构造 public: void stop() { is_pipeline_stopped = true; } };注意flow_control的构造函数是私有的,用户代码无法自行构造,只能以oneapi::tbb::flow_control&引用形式从operator()的参数中接收;stop()是它唯一的公开接口。这种设计保证了"停止信号"只能由图运行时产生并传递给 Body,用户无法伪造。
三、源码级调用链:Body 是如何被驱动的
理解调用链能帮你更准确地设计 Body 的内部状态。以 flow_graph.h 第 821-843 行 的try_reserve_apply_body为例,核心逻辑如下:
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; // 每次调用 Body 前构造一个全新的 control fgt_begin_body( my_body ); my_cached_item = (*my_body)(control); // 调用 Body,传入 control 引用 my_has_cached_item = !control.is_pipeline_stopped; // 依据 stop 标志决定是否缓存结果 fgt_end_body( my_body ); } if ( my_has_cached_item ) { v = my_cached_item; my_reserved = true; return true; } else { return false; // 已 stop:不产生数据,不再入队任务 } }整个数据驱动流程可以归纳为:
input_node被activate()激活(或图运行时需要数据)时,通过spawn_put()(flow_graph.h 第 853-857 行)在图的 arena 中投递一个input_node_task_bypass任务;- 任务执行
apply_body_bypass()(flow_graph.h 第 861-872 行):先调用try_reserve_apply_body获取一项数据,再通过my_successors.try_put_task(v)尝试推送给后继节点; - 若推送成功则
try_consume(),否则try_release()释放预留项; - 若 Body 内调用了
fc.stop(),is_pipeline_stopped置真,节点不再缓存结果,apply_body_bypass返回nullptr,数据流自然终止。
从这组实现可以推断出两个重要的工程结论:
flow_control是"一次性"的:每次调用 Body 都会创建全新的control对象(d1::flow_control control;),因此在一个 Body 调用中调用stop()只影响本次生成结果,不会"粘滞"影响后续调用——但因为节点拿到is_pipeline_stopped == true后就不再请求下一次数据,实际效果就是整条流的终止。stop()后返回的值确实会被丢弃:my_has_cached_item为 false 时,v = my_cached_item这行不会执行,返回值根本没有进入缓存,与规范描述完全吻合。
四、实战写法:满足 InputNodeBody 的四种典型实现
下面给出满足契约的典型实现,均以oneapi::tbb::flow::input_node<int>(Output为int)为例。注意oneapi::tbb::flow::input_node要求Output满足std::copyable(见 flow_graph.h 第 645-647 行 的__TBB_requires(std::copyable<Output>))。
写法一:lambda 表达式(最简单,官方示例的默认形态)
oneapi::tbb::flow::graph g; int counter = 0; // 注意:lambda 按引用捕获时,若以值传参给 input_node,会拷贝引用本身 oneapi::tbb::flow::input_node<int> src(g, & -> int { if (counter < 10) { return counter++; // 生成下一项 } fc.stop(); // 数据耗尽,停止 return -1; // 占位返回值,立即被丢弃 }); src.activate(); g.wait_for_all();写法二:函数对象类(状态封装更清晰,对应规范中的"类型 Body")
本仓库示例 examples/graph/binpack/binpack.cpp 第 202-217 行 给出了教科书式实现:
class item_generator { size_type counter; public: item_generator() : counter(0) {} value_type operator()(oneapi::tbb::flow_control& fc) { if (counter < elements_num) { value_type result = input_array[counter]; ++counter; return result; } fc.stop(); return value_type{}; // 停止时必须返回合法值 } };写法三:读取外部数据源(文件/IO)并终止
examples/graph/fgbzip2/fgbzip2.cpp 第 246-255 行 展示了读取文件直到数据耗尽的模式:
oneapi::tbb::flow::input_node<BufferMsg> file_reader( g, &io -> BufferMsg { if (io.hasDataToRead()) { BufferMsg bufferMsg = BufferMsg::createBufferMsg(io.chunksRead(), io.chunkSize()); io.readChunk(bufferMsg.inputBuffer); return bufferMsg; } fc.stop(); return BufferMsg{}; }); file_reader.activate();写法四:const 成员函数形式
examples/graph/logic_sim/basics.hpp 第 249 行 与 examples/parallel_pipeline/square/square.cpp 第 99-109 行 还展示了将operator()声明为const的写法——此时 Body 内部状态需用mutable成员维护,契约本身并不要求非 const。
通用骨架模板(可直接套用):
struct my_input_body { // 内部状态:计数器 / 迭代器 / IO 上下文等 int next_item{0}; int total{10}; int operator()(oneapi::tbb::flow_control& fc) { if (next_item < total) { return next_item++; // 1) 能生成:返回下一项 } fc.stop(); // 2) 不能生成:通知停止 return int{}; // 3) 停止时仍返回任意合法 Output(被丢弃) } };五、易错点与最佳实践
- 忘记调用
fc.stop()→ 无限循环:若 Body 在数据耗尽后仍不断返回合法值,input_node会持续向图中注入数据,图可能永不结束(如果后继节点消费能力跟得上)。这是最经典的 bug。 stop()分支返回类型必须正确:规范要求返回值与Output相同。虽然该值会被丢弃,但类型不匹配会导致编译错误(C++20 概念下报adaptive_same_as<Output>不满足)。Output必须可拷贝(std::copyable):input_node通过值传递消息,消息类型需要满足拷贝语义;如果数据很大(如 fgbzip2 中的BufferMsg),应设计轻量句柄或共享内存方案,避免每次拷贝大块数据。- Body 的拷贝构造与状态:
input_node构造时会拷贝 Body(my_body与my_init_body两份),节点重置(reset(),见 flow_graph.h 第 797-808 行)时会用初始副本恢复状态。若 Body 持有文件指针等"不可安全复制"的资源,请用智能指针并在operator()内正确管理生命周期。 activate()与wait_for_all()配对:input_node默认处于非激活状态,需要显式调用activate()(flow_graph.h 第 781-786 行)后才会开始投递任务;结束时在主线程调用g.wait_for_all()等待整张图完成,参见 fgbzip2.cpp 第 275 行。- 不要把
fc存起来跨调用使用:flow_control是每次调用临时构造的局部对象,规范与实现都假定它只在单次operator()调用内有效;保存引用或延长其生命周期属于未定义行为。
六、小结
InputNodeBody用三条接口要求(拷贝构造、析构、Output operator()(flow_control&))定义了 flow_graph 数据源节点的完整契约:它既规定了"如何生成下一项"(通过返回值),也规定了"如何宣告终止"(通过fc.stop())。在 flow_graph.h 中,这一契约同时以 C++20 概念input_node_body形式落地为编译期约束;在 examples/graph/binpack/binpack.cpp 与 examples/graph/fgbzip2/fgbzip2.cpp 中,则可以找到直接可用的权威范例。掌握这套契约,你就能在任何 oneTBB 流图应用中安全、正确地编写数据源节点,实现"按需生成 + 明确终止"的稳定数据流。
【免费下载链接】moldmold: A Modern Linker 🦠项目地址: https://gitcode.com/GitHub_Trending/mo/mold
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考