oneTBB flow_graph 的 async_node:连接流图与外部异步活动的桥接节点 并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载导读oneapi::tbb::flow::async_node是 oneTBB 流图flow graph中一个专用于跨运行时协作的节点它把图内的消息投递给由用户或其他运行时线程池、事件循环、GPU 队列等管理的外部活动并借助 gateway 接口让外部活动把处理结果安全地送回图中。读完本文你将掌握async_node的完整 API、类型约束、并发限制与节点策略queueing / rejecting / lightweight的取舍以及如何写出可被 conformance_async_node.cpp 验证的异步桥接代码。async_node 的定位与核心语义async_node是 flow graph 中唯一允许 body 把工作外派出图执行的节点。官方参考文档 async_node_cls.rst 将其定义为A node that enables communication between a flow graph and an external activity managed by the user or another runtime.它仍然是一个标准的流图节点同时继承三类基类能力graph_node归属于某个graph实例receiverInput可以作为其他节点的后继接收Input类型的消息senderOutput可以作为其他节点的前驱向后继广播Output类型的消息。它执行用户提供的 body 来处理到达的消息。典型用法是 body 并不直接计算而是把消息提交给图外的异步服务例如一个线程池、一个 IO 完成端口或一个 SYCL 队列并保证消息能被传递到该外部活动。真正的回传工作由gateway_type接口完成——外部活动通过它把结果送回流图。从消息流属性看async_node具有discarding丢弃与broadcast-push广播推送属性这一特性组合在 forwarding_and_buffering.rst 的节点属性总表中有明确定义。类模板签名与构造 APIasync_node声明于头文件oneapi/tbb/flow_graph.h实现在 include/oneapi/tbb/flow_graph.h 中模板签名如下namespace oneapi { namespace tbb { namespace flow { template typename Input, typename Output, typename Policy /*implementation-defined*/ class async_node : public graph_node, public receiverInput, public senderOutput { public: templatetypename Body async_node( graph g, size_t concurrency, Body body, Policy /*unspecified*/ Policy(), node_priority_t priority no_priority ); templatetypename Body async_node( graph g, size_t concurrency, Body body, node_priority_t priority no_priority ); async_node( const async_node src ); ~async_node(); using gateway_type /*implementation-defined*/; gateway_type gateway(); bool try_put( const input_type v ); bool try_get( output_type v ); }; } // namespace flow } // namespace tbb } // namespace oneapi构造函数详解① 基础构造可指定优先级templatetypename Body async_node( graph g, size_t concurrency, Body body, node_priority_t priority no_priority );构造一个调用body副本的async_nodeconcurrency限制该节点同时执行body调用的数量。该重载可以指定节点优先级。② 带 Policy 的构造templatetypename Body async_node( graph g, size_t concurrency, Body body, Policy /*unspecified*/ Policy(), node_priority_t priority no_priority );同样的构造语义额外允许以Policy模板参数或构造实参指定节点策略详见下文策略一节。文档约定最多可以有concurrency次对body的调用并发执行。③ 拷贝构造async_node( const async_node src );新节点与src引用同一个graph对象拥有src初始 body 的副本并沿用相同的并发阈值但src的前驱与后继不会被复制。新 body 是从src构造时传入的原始 body 的副本拷贝构造而来——src的 body 在构造之后对成员变量的任何修改都不会影响新节点的 body。④ 预览特性构造Helper Functions for Expressing Graphstemplate typename Body async_node(decltype(follows(...)), std::size_t concurrency, Body body); template typename Body async_node(decltype(precedes(...)), std::size_t concurrency, Body body);这是 :ref:预览特性 preview_features中 Helper Functions for Expressing Graphs 提供的语法糖允许直接以follows(...)/precedes(...)声明节点与一组节点的前后继关系来构造async_node。成员函数语义成员行为gateway_type gateway()返回gateway_type接口的引用供外部活动与图通信bool try_put( const input_type v )若并发限制允许则在消息v上执行 body否则按节点策略排队或拒绝该消息。返回true表示输入被接受false表示被拒绝bool try_get( output_type v )恒返回false。这是 discarding 属性在 API 层面的体现——输出不被缓存无法通过try_get取回try_put的排队或拒绝行为直接受Policy模板参数控制与 functional_node_policies.rst 定义的语义一致。类型要求Requirements创建合法的async_node需要满足以下约束见 async_node_cls.rstInput类型必须满足 ISO C 标准的DefaultConstructible可默认构造与CopyConstructible可拷贝构造要求。Policy类型可显式指定为 lightweight、queueing、rejecting 策略或使用默认值不指定。Body类型必须满足 AsyncNodeBody 命名需求。从 C17 起Body还可以是指向Input类型中一个接受gateway_type参数的const 成员函数的指针——conformance 测试中正是用test_invoke::SmartIDsize_t::send_id_to_gatewaygateway_type这种成员函数指针形式验证std::invoke对 body 的调用路径见 conformance_async_node.cpp 中 async_node and std::invoke 用例。AsyncNodeBody 命名需求AsyncNodeBody 需求要求Body提供拷贝构造函数Body::Body( const Body )析构函数Body::~Body()调用运算符void Body::operator()( const Input v, GatewayType gateway )其中Input必须与async_node模板参数一致GatewayType必须与该节点的gateway_type成员类型一致。语义上输入值v由流图提交给外部活动而 gateway 接口则允许外部活动与所在流图通信。GatewayType 命名需求GatewayType 需求定义了外部活动侧回传结果所需的三方法接口bool T::try_put( const Output v ); // 把 v 广播给对应 async_node 的全部后继 void T::reserve_wait(); // 通知流图已向外部活动提交了工作 void T::release_wait(); // 通知流图提交给外部活动的工作已完成其中Output必须与async_node模板参数Output一致。reserve_wait()/release_wait()成对使用是流图判断外部异步工作是否仍在途的依赖记账机制直接关系到graph::wait_for_all()能否正确等待所有外派工作收尾。并发限制从 serial 到 unlimitedasync_node的concurrency参数是用户可设置的并发上限既可以使用预定义常量也可以传任意std::size_t值限制在1到unlimited之间namespace oneapi { namespace tbb { namespace flow { std::size_t unlimited /*implementation-defined*/; // 允许无限数量的 body 并发调用 std::size_t serial /*implementation-defined*/; // 只允许单个 body 调用并发执行 }}}unlimited不限制 body 的并发调用数适合每来一条消息就立即外派的场景serial任意时刻只有一个 body 调用在执行适合需要保证提交顺序的外部服务中间值例如concurrency 4表示最多 4 个 body 调用可同时进行超出部分由策略决定排队或拒绝。conformance 测试在 async_node constructors 用例中同时覆盖了unlimited、serial与自定义node_priority_t的组合async_nodeint, int fn1(g, unlimited, fun); async_nodeint, int fn2(g, unlimited, fun, oneapi::tbb::flow::node_priority_t(1)); async_nodeint, int, lightweight lw_node1(g, serial, fun, lightweight()); async_nodeint, int, lightweight lw_node2(g, serial, fun, lightweight(), oneapi::tbb::flow::node_priority_t(1));见 conformance_async_node.cpp。节点策略Policyqueueing、rejecting 与 lightweightPolicy模板参数同时适用于function_node、multifunction_node、async_node与continue_node见 functional_node_policies.rst以 tag 类集合形式呈现class queueing { /*unspecified*/ }; class rejecting { /*unspecified*/ }; class lightweight { /*unspecified*/ }; class queueing_lightweight { /*unspecified*/ }; class rejecting_lightweight { /*unspecified*/ };queueing —— 排队接受输入消息若无法立即处理会被保留下来待节点空闲时再处理。对async_node而言当并发额度用满时后续try_put的消息进入内部队列不丢失。rejecting —— 拒绝接受输入消息若无法立即处理节点不接受该消息try_put返回false由前驱节点决定如何处理例如前驱重试或放弃。conformance 测试中的 async_node with rejecting policy 用例conformance_async_node.cpp专门验证了async_nodeint, int, oneapi::tbb::flow::rejecting的拒绝行为。lightweight —— 轻量执行提示lightweight是非绑定性提示声明节点 body 处理耗时很短供实现降低节点执行开销如绕过任务调度。任何优化都不得对节点与图的执行产生可观察的副作用。需要注意lightweight可与 queueing / rejecting 组合等价于queueing_lightweight/rejecting_lightweight当节点默认 Policy 为 queueing 时仅指定lightweight即等价于queueing_lightweightbody 的operator()必须声明为noexceptlightweight 优化才生效。官方在 lightweight_policy.cpp 中给出了流水线示例对耗时极短的第二、三节点乘 2、立方应用lightweight使它们无需任务调度开销即可执行第一节点不指定以允许图的并发调用function_node int, int add( g, unlimited, [](const int v) { return v1; } ); function_node int, int, lightweight multiply( g, unlimited, [](const int v) noexcept { return v*2; } ); function_node int, int, lightweight cube( g, unlimited, [](const int v) noexcept { return v*v*v; } ); make_edge(add, multiply); make_edge(multiply, cube); for(int i 1; i 10; i) add.try_put(i); g.wait_for_all();节点优先级node_priority_tasync_node与其他功能节点一样支持在构造时设置节点优先级引导执行图的线程优先选取高优先级节点的任务typedef unsigned int node_priority_t; const node_priority_t no_priority node_priority_t(0);参数值越大优先级越高默认值no_priority0表示关闭该节点的优先级在同一个图中带优先级的节点任务优先于低优先级或无优先级节点任务线程选取可执行任务时选择最高优先级者。典型收益当图中存在关键路径节点时如 critical_path_in_graph.png 展示的依赖流图给关键路径节点f2设置更高优先级可引导线程尽早执行它避免先执行短任务导致整图总耗时被拉长。构造async_node时可用node_priority_t(1)等方式传参。Body 的拷贝语义与 copy_body传入async_node的 body 对象会被拷贝构造后对原 body 成员变量的修改不会影响节点内部的执行副本。若需要在节点外部检查 body 内部状态应使用 copy_body 函数模板适用于continue_node、function_node、multifunction_node、input_node与async_nodenamespace oneapi { namespace tbb { namespace flow { // Defined in header oneapi/tbb/flow_graph.h template typename Body, typename Node Body copy_body( Node n ); }}}conformance 测试 async_node body copying 用例通过conformance::copy_counting_object验证了 body 被拷贝且copy_body能取回更新后的副本conformance_async_node.cpp。一个完整的异步桥接模式综合以上 API典型的async_node用法是在 body 中把消息提交给外部活动并保存 gateway 引用用于回传#include oneapi/tbb/flow_graph.h using namespace oneapi::tbb::flow; struct AsyncBody { // 外部活动例如线程池在完成处理后通过 gateway 把结果放回图中 void operator()( const int v, async_nodeint, int::gateway_type gateway ) { gateway.reserve_wait(); // 通知流图外部工作已提交 submit_to_external_service(v, [gateway, v] { gateway.try_put(v * 2); // 外部活动完成把结果广播给后继 gateway.release_wait(); // 通知流图外部工作已完成 }); } }; graph g; async_nodeint, int a(g, serial, AsyncBody{}); // serial保证提交顺序 function_nodeint, int print(g, 1, [](const int v) { /* consume result */ return v; }); make_edge(a, print); a.try_put(42); g.wait_for_all(); // 等待包含外部活动在内的全部工作完成其中reserve_wait()/release_wait()的配对保证wait_for_all()不会在外部工作仍在途时提前返回——这正是async_node相比普通function_node的核心差异。源码与测试佐证类模板实现位于 include/oneapi/tbb/flow_graph.hclass async_node定义于该文件flow命名空间内与文档声明一一对应官方 conformance 测试 conformance_async_node.cpp 覆盖了构造、拷贝构造、继承关系、broadcast 转发、并发限制、body 拷贝、rejecting 策略、输出/输入类型以及std::invoke调用 body 等全部行为策略与并发的更完整语义可进一步参考 functional_node_policies.rst、predefined_concurrency_limits.rst 与 forwarding_and_buffering.rst。小结async_node是 oneTBB flow graph 中面向图外世界的桥接节点通过try_put接收消息、body 把消息外派给用户或第三方运行时、gateway_typetry_putreserve_wait/release_wait把结果回传并正确参与图的等待与完成判定。合理组合concurrencyserial/unlimited/ 自定义值、Policyqueueing / rejecting / lightweight与node_priority_t优先级即可在不阻塞图内线程的前提下将 IO、GPU 或其他异步工作流无缝接入 oneTBB 流图。赞分享并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载相关推荐oneTBB 流图 GatewayType 概念解析async_node 与外部活动的通信契约oneTBB 流图 GatewayType 概念解析async_node 与外部活动的通信契约 GatewayType 是 oneTBBoneAPI Thr并发编程高性能计算oneAPI TBB flow_graph 的 async_msg 模板类打通异步任务与流图的数据桥接oneAPI TBB flow_graph 的 async_msg 模板类打通异步任务与流图的数据桥接 导读 async_msg 是 oneAPI TBBo开发工具构建工具系统编程ToastFish 摸鱼背单词指南3 分钟把通知栏变单词卡ToastFish 摸鱼背单词指南3 分钟把通知栏变单词卡 等编译的三分钟右下角弹出个单词拼写、音标、词性、中文释义、例句全都在通知里。这就是 Toas并发编程高性能计算创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考