oneTBB Flow Graph 资源限制(Resource Limiting):用 resource_limiter 与 resource_limited_node 安全协调共享外部资源
并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载oneTBBoneAPI Threading Building Blocks的 Flow Graph 提供了flow::resource_limiter与flow::resource_limited_node这一对预览特性组件用于让图节点安全、互斥地协调对共享外部资源如数据库连接、线程不安全库、单一硬件句柄的访问。阅读本文后你将掌握该预览特性的启用方式、两类核心类的构造与约束、基于优先级的防饥饿仲裁原理以及如何在真实代码中把并发图与独占资源访问结合起来。预览特性提示本功能属于 oneTBB 的预览特性Preview Feature默认关闭必须显式开启未来可能被修改或移除生产代码中慎用。官方文档明确给出两种启用方式在使用任何头文件之前定义宏TBB_PREVIEW_FLOW_GRAPH_RESOURCE_LIMITING或TBB_PREVIEW_FLOW_GRAPH_FEATURES为 1。注意宏必须位于#include oneapi/tbb/flow_graph.h之前见 预览特性总览 与 示例源码。特性总览Provider 与 Consumer 的分工资源限制特性由两个组件构成它们分别扮演资源提供方与资源消费方的角色flow::resource_limiterProvider管理一组资源句柄resource handles的类。它持有若干个类型为ResourceHandle的资源并向消费者授予对这些资源的独占访问权。flow::resource_limited_nodeConsumer一个消费者节点。只有当它从每一个关联的resource_limiter中都成功获取到资源访问权之后其用户提供的 body 才会被调用执行。一个典型的应用场景是多个 Flow Graph 节点共享同一个数据库连接。节点自身的并发度可以设置为unlimited无限但通过一个只持有单个句柄的resource_limiter所有节点在同一时刻最多只有一个 body 在访问数据库从而避免并发写坏线程不安全的库或连接池。防饥饿仲裁为什么不是随便发如果一个节点需要同时持有来自多个 limiter 的资源而各个 limiter 总是把资源授予只需要单个资源的节点那么该多资源节点就会长期得不到调度——这就是饥饿。为此resource_limiter采用了一种尽力而为best-effort的优先级仲裁策略在相互竞争的消费者之间更早提出请求的消费者优先而不是以未指定的顺序授予访问权。一个未得到满足的请求会保留其在队列中的位置。因此一个正在等待多个资源的消费者会随着等待时间的推移逐渐比后来的请求更受青睐。需要强调的是这是一种尽力而为策略而非硬性保证请求是按被观察到的先后顺序仲裁的因此 limiter 可能在一个更高优先级的请求尚未可见之前就把资源授予了出去。从源码看这个先到先得的优先级由 request_id 类 实现每个请求同时携带一个单调递增的唯一整数和构造时的steady_clock时间戳排序时先比时间戳、再比唯一整数来打破平局请求在其整个存续期内包括被拒绝后重新请求始终持有同一个 id所以反复失败的请求会逐渐超越所有后来的请求——这正是源码注释中所说的a request that keeps losing eventually outranks every later one。启用预览特性宏定义与 ABI 注意事项根据 preview_features.rst 的说明预览特性默认关闭需要在使用任何 oneTBB 头文件之前定义对应的TBB_PREVIEW_前缀宏。启用资源限制特性有两种等价方式#define TBB_PREVIEW_FLOW_GRAPH_RESOURCE_LIMITING 1 // 或 #define TBB_PREVIEW_FLOW_GRAPH_FEATURES 1选用TBB_PREVIEW_FLOW_GRAPH_FEATURES可以一次性打开 Flow Graph 的一组预览特性而TBB_PREVIEW_FLOW_GRAPH_RESOURCE_LIMITING则只打开本特性。宏定义必须发生在任何 include 之前——因为 oneTBB 可能通过其他头文件被间接包含。此外还要注意预览特性的 ABI 约定只有在所有翻译单元都用同一版本库、并启用同一组预览特性编译时才保证 ABI 兼容。resource_limiter资源提供方详解resource_limiterResourceHandle表示一个管理一个或多个ResourceHandle类型资源句柄的 Provider。它向消费者即resource_limited_node实例提供对所管理资源的独占访问权。关于句柄语义文档特别指出对于某些资源类型ResourceHandle本身就是资源对于另一些类型它可能是访问资源所用的轻量级实体例如指向真实数据库连接的唯一指针。构造方式与示例resource_limiter提供了四种构造方式覆盖了从可复制标量到不可复制对象的完整场景// 方式一initializer_list —— 管理 3 个 int 资源 tbb::flow::resource_limiterint int_limiter{1, 2, 3}; // 方式二piecewise_construct —— 为不可复制的句柄原地构造 using db_resource_handle std::unique_ptrDatabase, CloseDatabase; tbb::flow::resource_limiterdb_resource_handle db_limiter{std::piecewise_construct, std::forward_as_tuple(open_database())};上例中int_limiter管理 3 个int资源db_limiter管理单个Database句柄。由于db_resource_handle不可拷贝句柄通过std::piecewise_construct_t构造函数就地构造。源码中resource_limiter内部用一个std::listResourceHandle存放句柄池m_resource_handles配合spin_mutex保护并发访问。所有由同一个resource_limiter管理的句柄都被视为等价消费者具体拿到哪一个句柄是未指定的。完整构造器签名namespace oneapi { namespace tbb { namespace flow { template typename ResourceHandle class resource_limiter { public: using resource_handle_type ResourceHandle; resource_limiter() delete; template typename InputIterator resource_limiter(InputIterator first, InputIterator last); template typename ContainerBasedSequence resource_limiter(ContainerBasedSequence sequence); resource_limiter(std::initializer_listResourceHandle init); template typename Tuple, typename... Tuples resource_limiter(std::piecewise_construct_t, Tuple tuple, Tuples... tuples); ~resource_limiter(); }; } } }各构造器的要求与语义构造器要求语义与注意事项resource_limiter(InputIterator first, InputIterator last)InputIterator满足 ISO C 标准 [input.iterators]ResourceHandle可由std::iterator_traitsInputIterator::reference构造管理序列[first, last)中的句柄每个句柄由对应元素构造。若first last行为未定义源码中通过__TBB_ASSERT(!m_resource_handles.empty(), ...)断言空句柄池resource_limiter(ContainerBasedSequence sequence)序列类型满足 ContainerBasedSequence 要求等价于resource_limiter(std::begin(sequence), std::end(sequence))resource_limiter(std::initializer_listResourceHandle init)—等价于resource_limiter(init.begin(), init.end())。注意std::initializer_list本身要求ResourceHandle可拷贝构造resource_limiter(std::piecewise_construct_t, Tuple handle_args1, Tuples... handle_argsN)每个T及其对应实参满足ResourceHandle可由std::getN(std::forwardT(handle_args))...构造N ∈ [0, std::tuple_sizestd::decay_tT::value)管理1 sizeof...(Tuples)个句柄每个句柄由对应 tuple 的元素转发构造。用于句柄不可拷贝或必须就地构造的场景piecewise_construct构造器的用法示例——管理两个句柄第一个按Handle(arg1, arg2)构造第二个按Handle(arg3)构造tbb::flow::resource_limiterHandle limiter(std::piecewise_construct, std::forward_as_tuple(arg1, arg2), std::forward_as_tuple(arg3));类型约束与析构成员类型using resource_handle_type ResourceHandle;即资源句柄类型的别名。ResourceHandle要求必须满足 ISO C 标准 [moveconstructible] 与 [moveassignable] 要求可移动构造、可移动赋值。源码中句柄在授予消费者时以std::move方式从池中取出extract_handle释放时再放回池前端全程只做移动不做拷贝。析构函数~resource_limiter()销毁 limiter。若仍有消费者引用该 limiter行为未定义——因此请确保所有使用该 limiter 的节点先被销毁或不再使用它。resource_limited_node资源感知的消费者节点resource_limited_nodeInput, OutputTuple在单个输入端口接收消息并在处理消息之前向一个或多个资源提供方请求资源访问权。处理完成后它可能产生一个或多个输出消息通过N个输出端口广播给后继节点其中N std::tuple_sizeOutputTuple::value。节点语义与消息处理流程节点的并发阈值concurrency threshold可以设为预定义值如unlimited或任意std::size_t值。输入消息到达节点时占用一个并发槽位。当并发阈值允许、且每个关联的resource_limiter都授予了所需资源访问权时节点才会对输入消息执行用户 body。body 执行完毕后所有资源被归还给各自的resource_limiter。若并发阈值被超过或有一个或多个所需资源未被授予输入消息被放入内部缓冲队列待并发允许且所有资源到位后再处理。从源码看这一申请—获取—执行—归还的完整闭环由 resource_limited_body_leaf 实现每个输入消息被封装为一个request_data其notify_counter初始值为资源数 1节点向所有 provider 发出请求后每收到一次notify通知计数减一当计数到达 1 时创建一个try_acquire_resources_and_execute_task图任务去尝试获取全部资源。若某个 provider 拒绝获取例如句柄暂时耗尽已获取的资源会被立即归还并重新请求release_and_rerequest_resources_helper请求 id 保持不变从而保留其在仲裁队列中的优先级。节点类型与属性resource_limited_node同时是一个graph_node和一个receiverInput其输出端口类型为std::tuplesenderOutput...Output为OutputTuple中的元素类型。它具备discarding丢弃与broadcast-push广播推送两个属性含义见 转发与缓冲 文档broadcast-push消息会被推送给尽可能多的、愿意接收它的后继节点。discarding如果没有后继节点接收消息消息将被丢弃对图的后续执行无影响对应表格中try_get()为 no。完整类签名namespace oneapi { namespace tbb { namespace flow { template typename Input, typename OutputTuple class resource_limited_node : public graph_node, public receiverInput { public: using output_ports_type /*unspecified*/; template typename Body, typename ResourceLimiter, typename... ResourceLimiters resource_limited_node(graph g, std::size_t concurrency, std::tupleResourceLimiter, ResourceLimiters... resource_limiters, Body body); resource_limited_node(const resource_limited_node other); ~resource_limited_node(); bool try_put(const Input input); }; } } }构造函数template typename Body, typename ResourceLimiter, typename... ResourceLimiters resource_limited_node(graph g, std::size_t concurrency, std::tupleResourceLimiter, ResourceLimiters... resource_limiters, Body body);要求Body必须满足 ResourceLimitedNodeBody 命名要求。ResourceLimiter及ResourceLimiters中的每个类型都必须是oneapi::tbb::flow::resource_limiter的特化。构造出的节点具备以下性质属于图g并发阈值设为concurrency使用body对象消费resource_limiters中每个元素提供的资源。注意传入构造函数的 body 对象会被拷贝之后对成员变量的更新不会影响用于构造节点的原始对象。若需要从节点外部检查 body 内持有的状态可以使用copy_body函数获取更新后的 body。copy_body是定义在头文件oneapi/tbb/flow_graph.h中的函数模板copy_body_func.rsttemplate typename Body, typename Node Body copy_body( Node n );拷贝构造与析构拷贝构造resource_limited_node(const resource_limited_node other)构造一个与other构造时初始状态相同的节点——属于同一个graph对象、具有相同的并发阈值、使用other初始 body 的一份拷贝、消费同一组资源。other的前驱与后继不会被拷贝新 body 是从other构造时提供的原始 body 的拷贝再拷贝而来other构造后对成员变量的修改不会影响新节点。析构~resource_limited_node()销毁节点对象。try_put 与返回值bool try_put(const Input input);将消息input传入节点。一旦并发阈值允许且所有必需资源均被授予节点即在input上执行用户 body。返回值恒为true——注意消息可能因资源或并发槽位不足而进入内部缓冲排队但这并不构成拒绝接收。类型要求Input必须满足 ISO C 标准的DefaultConstructible与CopyConstructible要求。OutputTuple必须是std::tuple的特化。ResourceLimitedNodeBodybody 的命名要求传入resource_limited_node的Body必须满足 ResourceLimitedNodeBody 命名要求核心签名如下Body::Body(const Body other); // 拷贝 body Body::~Body(); // 销毁 body void Body::operator()(const Input input, OutputPortsType ports, ResourceHandle1 resource_handle1, ..., ResourceHandleN resource_handleN);其中rl_node表示构造时传入该 body 的resource_limited_node实例rl_node_type为其类型。要求Input与rl_node_type的Input模板实参相同OutputPortsType与rl_node_type::output_ports_type成员类型相同ResourceHandle1...ResourceHandleN与构造时传给节点的对应resource_limiter::resource_handle_type相同。语义上body 负责处理输入消息并且可以在任意输出端口上多次调用try_put包括同一端口多次。注意 body 的调用签名的特殊性除了input和ports它还按关联 limiter 的顺序接收资源句柄引用这也是它区别于普通function_nodebody 的关键。完整示例独占数据库连接的读写图下面是在官方文档示例基础上整理出的完整可运行程序源自 examples/resource_limiting.cpp两个节点通过同一个resource_limiter共享一个数据库连接句柄。因为db_limiter只持有一个资源句柄所以db_reader与db_writer的 body永远不会被同时调用——尽管两个节点的并发度都设置为unlimited。#define TBB_PREVIEW_FLOW_GRAPH_RESOURCE_LIMITING 1 #include oneapi/tbb/flow_graph.h struct DB_handle { void read() {} void write() {} }; DB_handle* open_database() { static DB_handle global_db_handle; return global_db_handle; } struct processor_body { int operator()(int id) const { return id; } }; auto input_ids {1, 2, 3}; int main() { using namespace tbb::flow; resource_limiterDB_handle* db_limiter{open_database()}; graph g; using resource_limited_node_type resource_limited_nodeint, std::tupleint; // 并发度无限但 db_limiter 确保对数据库的独占访问 resource_limited_node_type db_reader(g, unlimited, std::tie(db_limiter), [](int id, auto ports, DB_handle* db) { db-read(); // 其他基于读出的数据进行的操作 std::get0(ports).try_put(id); }); function_nodeint, int processor(g, unlimited, processor_body{}); resource_limited_node_type db_writer(g, unlimited, std::tie(db_limiter), [](int id, auto ports, DB_handle* db) { // 其他与数据库相关的操作 db-write(); std::get0(ports).try_put(id); } ); make_edge(output_port0(db_reader), processor); make_edge(processor, db_writer); // 其他图节点与边 for (int id : input_ids) { db_reader.try_put(id); } g.wait_for_all(); }这个例子清晰展示了两个要点限流语义来自资源而非并发度两个节点的并发阈值都是unlimited真正的互斥约束由db_limiter的单个句柄提供。这与limiter_node按消息数限流的定位不同。多 limiter 组合构造器接受std::tupleResourceLimiter, ResourceLimiters...即一个节点可以同时关联多个 limiterbody 将按顺序收到多个资源句柄参数。底层实现原理请求生命周期与仲裁细节结合 实现头文件可以还原资源协调的完整调用链发起请求resource_limited_input::apply_body_impl_bypass在获得并发槽位后调用 body 的operator()resource_limited_body_leaf::operator()为消息生成一个带唯一 id 的request_data随后通过request_resources_helper向所有关联 provider 依次发出requestrequest_resources_helper。通知provider 在锁内评估请求优先级。若通知列表未满或新请求优先级足够高则加入通知列表并回调consumer.notify否则进入m_pending等待。notify计数归零时节点调度一个try_acquire_resources_and_execute_task图任务。获取任务执行acquire_resources_helper逐一向各 provider 调用acquire。若某 provider 返回空资源暂时耗尽已获资源通过release_and_rerequest_resources_helper归还并重新请求请求 id 不变——这正是等待多个资源的消费者逐渐获得优先级的机制来源acquire内部用std::nth_element求出当前服务阈值select_lowest_priority_to_serve只让胜过阈值最低者的请求获得句柄。执行与归还全部资源到手后调用用户 body完成后通过release_resources_helper把句柄依次归还release将句柄emplace_front回池中并调用notify_pending唤醒等待者同时释放并发槽位并取出下一条排队消息。取消路径若图任务被取消cancel_request会通过withdraw从各 provider 的列表中撤回请求避免遗留的请求永远占用通知名额源码注释明确指出这一点。resource_limiter的核心数据结构是m_resource_handles句柄池std::list、m_notified已通知请求与m_pending挂起请求所有操作都受一把spin_mutex保护这也是文档所称尽力而为优先级的具体实现。验证与测试仓库中的一致性证据仓库为这一预览特性提供了专门的测试套件 test/tbb/test_resource_limited_node.cpp其中几个测试直接印证了文档所述语义单资源互斥test_single_resource创建 10 个并发度为unlimited的节点共享一个int*句柄body 内循环 1000 次断言counter 1验证单个资源同一时刻只授予一个消费者test_resource_limited_node.cpp#L49-L91。多资源同时授予test_several_resourcesresource_limiterint{0,1,...,9}提供 10 个句柄10 个节点各获得一个后同时执行借助 resumable tasks 挂起/恢复同步验证多个句柄可被并行授予不同消费者test_resource_limited_node.cpp#L95-L155。类型继承关系test_inheritance静态断言resource_limited_node是graph_node与receiverInput的派生类与文档声明一致test_resource_limited_node.cpp#L40-L47。测试文件顶部同样定义了TBB_PREVIEW_FLOW_GRAPH_RESOURCE_LIMITING 1test_resource_limited_node.cpp#L17再次印证启用宏的必要性。此外test/tbb/test_flow_graph_whitebox.cpp 与 test/conformance/conformance_flowgraph.h 也涉及该特性可作进一步参考。关键注意事项速查关注点说明启用宏在 include 任何 oneTBB 头文件之前定义TBB_PREVIEW_FLOW_GRAPH_RESOURCE_LIMITING或TBB_PREVIEW_FLOW_GRAPH_FEATURES为 1预览特性风险未来可能被移除或显著改变不遵循常规弃用流程生产代码强烈不建议使用ABI仅当所有翻译单元使用同一库版本、同一组预览特性时保证 ABI 兼容空句柄池resource_limiter(first, last)在first last时行为未定义源码断言句柄池非空句柄等价性同一 limiter 内所有句柄等价具体分发哪个句柄未指定句柄类型要求ResourceHandle需满足 MoveConstructible 与 MoveAssignableinitializer_list构造额外要求可拷贝limiter 析构仍有消费者引用 limiter 时析构行为未定义body 拷贝构造时 body 被拷贝外部如需读取最新状态可用copy_bodytry_put返回值恒为true消息可能内部排队而非被拒绝节点属性discarding broadcast-push多输出端口数量由OutputTuple决定输入类型要求Input需 DefaultConstructible CopyConstructibleOutputTuple必须是std::tuple特化综上oneTBB 的资源限制特性通过provider 持有句柄、consumer 按优先级请求的架构把并发度与资源互斥两个正交的约束解耦节点仍可按需设置高并发而共享资源的独占性由resource_limiter统一保障。对于数据库连接、线程不安全的第三方库、硬件设备句柄等共享外部资源场景这是一个值得在符合预览特性约束的前提下评估使用的方案。更多 Flow Graph 细节可参考 Flow Graph 参考文档 与头文件 flow_graph.h。赞分享并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载相关推荐oneTBB Flow Graph 资源限制节点 API 全解resource_limiter 与 resource_limited_node 的设计、演进与源码实现oneTBB Flow Graph 资源限制节点 API 全解resource_limiter 与 resource_limited_node 的设计、演进与并发编程高性能计算oneTBB Flow Graph 资源限制previewresource_limiter 类完全指南oneTBB Flow Graph 资源限制previewresource_limiter 类完全指南 本文是 oneAPI Threading Buil并发编程高性能计算Apache Pulsar Functions 全面指南在 Pulsar 内构建轻量级流处理函数Apache Pulsar Functions 全面指南在 Pulsar 内构建轻量级流处理函数 Pulsar Functions 是 Apache Puls并发编程高性能计算上一篇git-extras 之 git release 命令实战指南一键完成提交、打标签与推送的版本发布流程下一篇Cimoc漫画阅读器深度解析开源Android漫画聚合平台终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考