ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

oneTBB Flow Graph 并发限制实战:用 rejecting function_node 构建背压式流图

oneTBB Flow Graph 并发限制实战:用 rejecting function_node 构建背压式流图 oneTBB Flow Graph 并发限制实战用 rejecting function_node 构建背压式流图【免费下载链接】moldmold: A Modern Linker 项目地址: https://gitcode.com/GitHub_Trending/mo/mold本文以 oneTBBThreading Building BlocksFlow Graph 的并发限制机制为核心讲解function_node的并发阈值与queueing/rejecting两种缓冲策略的本质区别并通过input_node rejecting function_node的完整示例演示如何把流图中同时在途的大对象数量钳制在固定上限。读完本文你将掌握如何在数据流图中实现背压backpressure、如何控制节点实例并发数以及如何在 oneTBB 源码层面理解这两种策略的底层实现。为什么需要并发限制在 Flow Graph 中一个节点可能被上游持续投喂消息。对于 CPU 密集型的function_node如果上游产出速度远快于下游消费速度消息会无限堆积内存占用与资源消耗随之失控。oneTBB 提供两类手段控制同时在途的消息数量节点级并发限制通过function_node的并发阈值concurrency limit约束同一时刻正在执行 body 的实例数图级限流借助limiter_node、join_nodetoken 机制等节点在整张图中设置消息数量上限。本文主角是第一种即在单个节点上设置并发限制当节点达到并发上限后它可以选择继续接收并缓冲消息也可以选择直接拒绝新消息——这由节点的缓冲策略graph_buffer_policy决定。function_node 的并发阈值与缓冲策略function_node的模板声明如下完整定义见 flow_graph.htemplate typename Input, typename Output continue_msg, graph_buffer_policy queueing class function_node;三个模板参数分别是参数含义默认值Input节点接收的消息类型必填Output节点产出的消息类型continue_msg仅用于触发后继不携带数据graph_buffer_policy输入缓冲策略queueing排队缓冲构造函数的第二个参数concurrency就是并发限制function_node( graph g, size_t concurrency, Body body, Policy Policy(), node_priority_t a_priority no_priority );concurrency表示该节点最多允许多少个 body 实例同时执行。它既可以是具体数值也可以使用unlimited不限制。queueing达标后继续收内部排队当策略为queueing默认时即使节点已经达到并发限制它仍然接受上游投递的消息只是不立即执行而是把消息暂存在节点内部的输入队列中等待某个并发槽位释放后再取出来处理。消息不会丢失但队列可能无限增长。rejecting达标后拒收形成背压当策略为rejecting时节点达到并发限制后try_put会返回失败拒绝该消息。被拒绝的消息不会进入任何缓冲而是弹回上游——由上游决定如何处理。这正是 Flow Graph 中实现背压的核心手段让慢的下游通过拒绝消息反向拖慢快的上游。从源码看两种策略的实现差异两种策略的底层差异在输入基类function_input_base的构造函数中体现得淋漓尽致见 _flow_graph_node_impl.hfunction_input_base( graph g, size_t max_concurrency, node_priority_t a_priority, bool is_no_throw ) : my_graph_ref(g), my_max_concurrency(max_concurrency) , my_concurrency(0), my_priority(a_priority), my_is_no_throw(is_no_throw) , my_queue(!has_policyrejecting, Policy::value ? new input_queue_type() : nullptr) , my_predecessors(this) , forwarder_busy(false) { my_aggregator.initialize_handler(handler_type(this)); }关键点有两处my_max_concurrency与my_concurrency前者保存构造时传入的并发上限后者记录当前正在执行的实例数。每次占用一个并发槽位时递增、body 执行完毕释放时递减节点据此判断是否已达上限。my_queue的条件分配!has_policyrejecting, Policy::value为真即不是 rejecting时才new input_queue_type()也就是说rejecting 节点根本不分配输入队列物理上就不具备收下再说的能力消息要么被立即执行有空闲槽位要么被拒绝。而 queueing 节点会为每个输入类型分配一个function_input_queue。策略标签本身定义在 _flow_graph_body_impl.h 的graph_policy_namespace命名空间中struct rejecting { }; struct reserving { }; struct queueing { }; struct lightweight { };同时 _flow_graph_node_impl.h 中还有一条静态断言防止queueing与rejecting被同时指定static_assert(!has_policyqueueing, Policy::value || !has_policyrejecting, Policy::value, );从源码结构可以推断rejecting 节点把背压的责任委托给了上游如input_node或buffer_node自己只保留接受或拒绝的二元状态而 queueing 节点则通过内部队列把背压问题吞下去代价是内存占用随队列增长。实战用 rejecting function_node 限制在途大对象数量官方用户指南 use_concurrency_limits.rst 给出了一个非常典型的场景input_node源源不断创建大对象我们希望在流图中同时在途已创建但尚未被删除的大对象数量不超过 3 个。做法是在input_node下游放置一个并发限制为 3 的rejectingfunction_nodegraph g; int src_count 0; int number_of_objects 0; int max_objects 3; // 上游按需生成 big_object最多生成 M 个 input_node big_object * s( g, - big_object* { if ( src_count M ) { big_object* v new big_object(); src_count; return v; } else { fc.stop(); // 告知流图数据源已耗尽停止拉取 return nullptr; } } ); s.activate(); // 构造后需显式激活 input_node 才会开始产出 // 下游并发限制为 3 的 rejecting 节点同时最多执行 3 个 body 实例 function_node big_object *, continue_msg, rejecting f( g, 3, []( big_object *v ) - continue_msg { spin_for(1); // 模拟耗时处理 delete v; // 处理完毕释放大对象 return continue_msg(); } ); make_edge( s, f ); // 建立 s - f 的边 g.wait_for_all(); // 等待整个图执行完毕逐环节拆解运行过程整个图的节拍可以拆成如下循环input_node s被激活后其 body 被调用创建第一个大对象并投递给ff有空闲并发槽位当前执行实例数 3接受消息并启动一个 body 实例槽位占满后f进入拒绝模式当f已有 3 个实例在并发执行时它拒绝s投来的下一条消息input_node收到拒绝信号后暂停调用自己的 body把最后创建的那个大对象暂存在内部等待被重新拉取一旦f的某个实例执行完毕、并发数降到 3 以下它会重新从s拉取新消息s的 body 再次被唤醒。于是整张图中最多同时存在 4 个大对象3 个正在function_node中并发处理1 个缓冲在input_node中等待放行。input_node在这里扮演了上游暂存 背压吸收的角色——这正是它体含flow_control停止/恢复机制的原因。为什么这里必须用 rejecting 而非默认的 queueing如果把策略换成默认的queueingf达到并发上限后仍会接受所有消息并堆进内部队列input_node会继续疯狂创建对象队列无限膨胀同时在途对象不超过 4 个的约束立刻被打破。所以当目标是用并发限制钳制资源占用时rejecting 是唯一正确的选择queueing 更适合消息一个都不能丢、允许排队的场景。相关机制对比limiter_node 与 Token 系统rejectingfunction_node解决的是单个节点的并发问题。若要限制一段子图或整张图的消息总量oneTBB 还提供两种互补方案同属用户指南中的资源限制专题参见 use_limiter_node.rst 与 create_token_based_system.rst。limiter_node带 decrement 端口的闸门limiter_node在指定点设置允许通过的消息总数上限limiter_node( graph g, size_t threshold );它维护一个内部计数放行消息时计数加一当计数达到阈值后开始拒绝上游消息。它额外提供一个decrement端口——受控部分处理完一条消息后把结果投回该端口即可让计数减一、放行下一条limiter_node big_object * l( g, max_objects ); function_node big_object *, continue_msg f( g, unlimited, /* ... */ ); make_edge( l, f ); // 闸门 - 处理节点 make_edge( f, l.decrement ); // 处理完成 - 计数减一 make_edge( s, l ); // 数据源 - 闸门Token 系统用 reserving join_node 做消息-令牌配对更灵活的做法是把可并发数抽象为令牌token预先向buffer_node注入 N 个令牌再用reserving策略的join_node把输入消息与令牌配对——只有当令牌可用时join_node才从input_node拉取新消息消息处理完毕后由function_node把令牌回送到buffer_node循环利用。由于令牌类型完全自定义可以是任意对象甚至大对象本身这套机制甚至可以实现对象池与运行期动态调整并发度。三种方案的选择要点方案限制粒度适用场景rejectingfunction_node单节点并发实例数限制某个计算节点的同时执行数形成背压limiter_node一段图/一条路径上的消息总量保护下游子图需要 decrement 反馈环令牌 join_node全图在途消息总量需要动态增删并发度、复用对象池在 mold 仓库中的位置与验证方式oneTBB 以第三方依赖的形式内嵌于本仓库的 third-party/tbb 目录本主题所在的完整用户指南文档位于 third-party/tbb/doc/main/tbb_userguide与并发限制直接相关的姊妹篇还包括use_input_node.rstinput_node的惰性创建与activate()时机理解背压示例的前提Flow_Graph_Buffering_in_Nodes.rst节点缓冲与消息转发的通用语义use_limiter_node.rst 与 create_token_based_system.rst上文对比的另两种限流方案。本文引用的策略实现源码位于 flow_graph.hfunction_node类模板、_flow_graph_node_impl.hfunction_input_base的队列分配与 _flow_graph_body_impl.h策略标签定义。mold 本体的并行化如 gc-sections.cc、arch-arm32.cc 中的tbb::parallel_for/tbb::parallel_for_each同样复用这份 tbb 依赖可见其作为高性能并行基础库在本项目中的地位。需要提醒的是graph_buffer_policy属于 Flow Graph 接口mold 自身并未在链接器主流程中使用 Flow Graph本文讨论的并发限制能力是 tbb 组件面向所有使用者的通用特性可直接在任意使用 oneTBB 的项目中按上述示例落地。【免费下载链接】moldmold: A Modern Linker 项目地址: https://gitcode.com/GitHub_Trending/mo/mold创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表