ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

mold 内置 TBB Flow Graph:join_node 多路消息聚合节点的规范与源码解析

mold 内置 TBB Flow Graph:join_node 多路消息聚合节点的规范与源码解析 mold 内置 TBB Flow Graphjoin_node 多路消息聚合节点的规范与源码解析【免费下载链接】moldmold: A Modern Linker 项目地址: https://gitcode.com/GitHub_Trending/mo/mold在构建大型并发计算流水线时如何让多个输入流在满足特定条件后汇合成一条输出流是数据流编程的核心问题。本文以 mold 仓库随附的 TBBthird-party/tbb中 flow_graph 的join_node规范文档为骨架完整讲解join_node的类定义、模板约束、三种缓冲策略reserving、queueing、key_matching以及各成员函数的语义并结合 flow_graph.h 与 _flow_graph_join_impl.h 的源码实现帮助读者掌握在 TBB 流图中构建多路聚合节点的实战方法。一、join_node 的定位多路输入聚合成元组join_node是一个从一组在输入端口上接收到的消息构造元组tuple并将该元组广播到所有后继节点的节点。它是graph_node与senderOutputTuple的组合体内部包含一个输入端口元组其中每个端口都是对应消息类型的receiverType因此天然支持多个不同类型的输入接收器并且要求所有输入端口使用相同的缓冲策略。规范文档 join_node_cls.rst 给出的完整类声明如下定义于头文件oneapi/tbb/flow_graph.h// Defined in header oneapi/tbb/flow_graph.h namespace oneapi { namespace tbb { namespace flow { using tag_value /*implementation-specific*/; templatetypename OutputTuple, class JoinPolicy /*implementation-defined*/ class join_node : public graph_node, public sender OutputTuple { public: using input_ports_type /*implementation-defined*/; explicit join_node( graph g ); join_node( const join_node src ); input_ports_type input_ports( ); bool try_get( OutputTuple v ); }; templatetypename OutputTuple, typename K, class KHashtbb_hash_compareK class join_node OutputTuple, key_matchingK,KHash : public graph_node, public sender OutputTuple { public: using input_ports_type /*implementation-defined*/; explicit join_node( graph g ); join_node( const join_node src ); template typename B0, typename... BN join_node( graph g, B0 b0, BN... bn ); input_ports_type input_ports( ); bool try_get( OutputTuple v ); }; } // namespace flow } // namespace tbb } // namespace oneapi可以看到规范区分了两个层次主模板join_nodeOutputTuple, JoinPolicy承载reserving与queueing两种策略偏特化join_nodeOutputTuple, key_matchingK,KHash额外提供一个接受键提取函数对象包b0, ..., bn的构造函数。二、模板参数与使用前提Requirements规范对模板参数有严格的前提约束使用join_node前必须逐条核对OutputTuple必须是std::tuple的实例化。元组中存储的每一种类型都必须同时满足 ISO C 标准的 [defaultconstructible]默认可构造、[copyconstructible]可拷贝构造与 [copyassignable]可拷贝赋值要求。这一点在 flow_graph.h 中同样体现——各策略偏特化内部都通过std::tuple_sizeOutputTuple::value推导端口数N因此非 tuple 类型无法使用。JoinPolicy必须是指定的缓冲策略之一见下一节的 join_node_policiesreserving、queueing或key_matching。KHash类型必须满足 HashCompare 要求默认值为tbb_hash_compareK。Bi类型键提取函数对象必须满足 JoinNodeFunctionObject 要求即函数对象可拷贝、可析构且operator()(const Input v)返回类型Key与join_node的模板参数K一致、Input与OutputTuple对应元素一致。此外自 C17 起Bi还可以是Input上返回Key的 const 成员函数指针或Input上类型为Key的数据成员指针。join_node具备buffering缓冲与broadcast-push广播推送两种 节点属性消息会在输入端口处缓冲直至能构造出元组随后以广播方式推给全部后继。三、三种缓冲策略Join Policiesjoin_node的行为完全由其缓冲策略决定。根据 join_node_policies.rst 的规范// Defined in header oneapi/tbb/flow_graph.h namespace oneapi { namespace tbb { namespace flow { struct reserving; struct queueing; templatetypename K, class KHashtbb_hash_compareK struct key_matching; using tag_matching key_matchingtag_value; } // namespace flow } // namespace tbb } // namespace oneapi3.1 queueing无界 FIFO 队列每个输入端口被put时消息被加入该端口的无界先进先出FIFO队列当每个端口队列中至少有一条消息时join_node将各队列头部消息组成元组并广播给所有后继若至少一个后继接受该元组则各端口队列头部被移除否则消息原样留在各队列中等待下一次尝试。这是默认策略从 flow_graph.h 的主模板前向声明templatetypename OutputTuple, typename JPqueueing class join_node;可以确认JoinPolicy缺省即为queueing。3.2 reserving预留式握手端口被put时join_node只标记该端口可能有消息可用并返回false表示消息未被消费当所有端口都被标记后join_node尝试从各端口的已知前驱预留reserve一条消息若某个端口预留失败则取消该端口的标记并释放所有已获得的预留全部回滚若全部端口预留成功则广播包含这些消息的元组只要有后继接受预留即被消费consume否则预留被释放release。3.3 key_matching按键匹配端口被put时用户提供的函数对象被应用于消息以计算其键key消息随后被加入该输入端口的哈希表当给定某个 key 在每个输入端口都有对应消息时join_node从各端口取出所有匹配消息、构造元组并尝试广播若没有后继接受该元组元组会被保存并在后续try_get时再转发。3.4 tag_matchingkey_matching 的特例tag_matching即key_matchingtag_value的别名源码 flow_graph.h 中注释明确写道 tag_matching join_node is a specialization of key_matching, and is source-compatible它接受tag_value类型的键。规范指出tag_value是一个无符号整型专用于定义tag_matching策略。四、成员类型与成员函数详解4.1 成员类型input_ports_type是输入端口元组的别名a tuple of input ports与 flow_graph.h 中各偏特化暴露的typedef typename unfolded_type::input_ports_type input_ports_type;一致。4.2 构造函数1空节点构造explicit join_node( graph g );构造一个属于流图g的空join_node。实现上以reserving偏特化为例flow_graph.h__TBB_NOINLINE_SYM explicit join_node(graph g) : unfolded_type(g) { fgt_multiinput_nodeN( CODEPTR(), FLOW_JOIN_NODE_RESERVING, this-my_graph, this-input_ports(), static_cast sender output_type *(this) ); }可见构造函数在完成基类初始化后还会调用fgt_multiinput_nodeN向 Graph 流工具GFT注册该多输入节点节点类型标识为FLOW_JOIN_NODE_RESERVING/FLOW_JOIN_NODE_QUEUEING/FLOW_JOIN_NODE_TAG_MATCHING这是后续可视化调试多输入节点的依据。2带键提取函数的构造仅 key_matching 偏特化template typename B0, typename... BN join_node( graph g, B0 b0, BN... bn );该构造函数仅存在于key_matching偏特化中使用函数对象b0, ..., bN为输入端口0到N分别确定标签key。约束只有当std::tuple_sizeOutputTuple::value 1 sizeof...(BN)时才参与重载决议——即函数对象数量必须恰好等于元组元素数源码中体现为std::enable_if1 sizeof...(Bodies) N的 SFINAE 条件见 flow_graph.hC20 concepts 环境下则由join_node_functionsconcept 校验每个函数对象均满足join_node_function_objectFunction, 元组元素, K。警告规范原文传入join_node构造函数的函数对象不得抛出异常。它们会被并行调用因此应当是纯函数、耗时极短且非阻塞的。3拷贝构造join_node( const join_node src );创建一个与src在构造时刻具有相同初始状态的join_node。注意前驱列表、输入端口中的消息以及后继都不会被拷贝。源码中该构造函数会重新注册 GFT 节点事件如fgt_multiinput_node即新节点在自己的图中重新建立调试跟踪身份。4.3 input_ports()input_ports_type input_ports( )返回一个std::tuple其每个元素都是对应输入消息类型的receiverT继承体可像任何receiverT一样与make_edge等机制连接。各端口的具体行为由所选join_node策略决定。规范还提到非成员函数模板 input_port 可以简化获取特定输入端口的语法input_port(n)直接返回元组中第n个端口的引用避免书写冗长的std::get...(my_join.input_ports())。4.4 try_get()bool try_get( output_type v )尝试依据join_node的缓冲策略构造元组成功则把元组拷贝到v并返回true否则返回false。从源码结构看_flow_graph_join_impl.htry_get并非直接访问端口数据而是构造一个join_node_base_operation(v, try__get)并提交给节点内部的aggregator聚合器串行执行bool try_get( output_type v) override { join_node_base_operation op_data(v, try__get); my_aggregator.execute(op_data); return op_data.status SUCCEEDED; }聚合器的handle_operations在处理try__get操作时先检查tuple_build_may_succeed()能否成功构造元组再调用try_to_make_tuple真正取出/预留各端口消息构造成功后执行tuple_accepted()消费队列头部或预留失败则标记FAILED。这种所有操作经聚合器串行化的设计保证了多线程下端口数据访问的一致性。五、非成员类型与类模板推导指南Deduction Guide规范定义了推导指南使得 C17 下可以用更简洁的写法声明key_matching节点template typename Body, typename... Bodies join_node(graph, Body, Bodies...) -join_nodestd::tuplestd::decay_tinput_tBody, std::decay_tinput_tBodies..., key_matchingoutput_tBody;其中input_t是传入函数对象输入参数类型的别名output_t是传入函数对象返回类型的别名。也就是说只需传入键提取函数对象编译器即可反推出OutputTuple各函数对象参数类型的去引用组合与key_matchingK以第一个函数对象的返回类型作为键类型。六、源码实现要点三个层次的展开理解join_node的实现关键在 flow_graph.h 与 _flow_graph_join_impl.h 中的分层结构端口类型层reserving_port、queueing_port、key_matching_port分别对应三种策略。每个端口持有一个指向join_node_FE前端的反向指针通过join_helperN::set_join_node_pointer逐一装配_flow_graph_join_impl.h。join_helperN是一组递归模板负责以 O(N) 方式对元组逐元素执行reserve/get_item/consume/release/reset等批量操作——例如reserve从第N-1个端口开始尝试预留任一口失败即回滚已预留的部分release_my_reservation精确实现了第三节所述全部预留成功才广播否则全部释放的语义。前端join_node_FE层负责策略相关的数据结构。以key_matching为例_flow_graph_join_impl.hFE 持有my_inputs端口元组、key_to_count_buffer键到各端口就绪计数的缓冲以及输出缓冲try_to_make_tuple在聚合器保护下按current_key组装元组tuple_accepted()负责在各端口消费消息后重置键计数。基类join_node_base层同时继承graph_node、join_node_FE与senderOutputTuple把收消息FE 的端口与发消息broadcast_cache后继集合粘合起来。其聚合器处理四类操作reg_succ注册后继若元组可构造且图处于激活状态则派生 bypass 转发任务并置位forwarder_busy、rem_succ移除后继、try__get手动拉取与do_fwrd_bypassbypass 模式下的连续自动转发循环try_to_make_tupletry_put_task直到构造失败或后继拒绝。unfolded_join_node则是把策略相关的端口元组类型展开给公共join_node偏特化使用的中间模板。七、实战示例dining_philosophers 中的 reserving joinTBB 自带示例 dining_philosophers.cpp 展示了join_node的典型用法typedef oneapi::tbb::flow::join_nodejoin_output, oneapi::tbb::flow::reserving join_node_type; ... std::vectorjoin_node_type join_vector(num_philosophers, join_node_type(g));该图让拿起左手与拿起右手两条输入流在reserving策略的join_node处汇合——只有当哲学家左右两只筷子都就绪时才触发进餐节点。选用reserving而非queueing的意图在于reserving端口在消息未被最终消费前不会真正吞掉消息put返回false即回退适合需要成对匹配失败则原路释放的场景避免了 FIFO 队列可能造成的左右手消息错位配对。八、小结join_node是 TBB flow graph 中实现多路聚合的核心节点以std::tuple描述多路异构输入以JoinPolicy决定聚合时机与缓冲方式queueing默认FIFO 对齐聚合、reserving预留—消费/回滚的握手聚合、key_matching哈希按键聚合含tag_matching特例三种策略覆盖了绝大多数汇流需求input_ports()返回的端口元组可像普通receiverT一样建边try_get()提供手动拉取通道C17 推导指南让key_matching节点声明更简洁实现上采用端口层 前端层 聚合器基类的三层结构所有并发操作经 aggregator 串行化兼顾了线程安全与 bypass 转发性能。在 mold 仓库中该规范文档位于 third-party/tbb/doc/main/specification/source/flow_graph/ 目录配套的策略说明、input_port辅助函数、JoinNodeFunctionObject 命名要求以及forwarding_and_buffering节点属性等文档均在同目录下可结合阅读以获得 flow_graph 节点体系的完整图景。【免费下载链接】moldmold: A Modern Linker 项目地址: https://gitcode.com/GitHub_Trending/mo/mold创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表