ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

oneTBB aggregator 类详解:不建模 Mutex Concept 的互斥批量执行聚合器

oneTBB aggregator 类详解:不建模 Mutex Concept 的互斥批量执行聚合器 并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载本篇文章以 oneTBB 参考文档中aggregator类Basic Interface 与 Expert Interface为核心系统讲解这个“提供互斥执行但不建模 Mutex Concept”的特殊同步原语包括基础接口execute的用法、专家接口aggregator_ext与操作节点的处理协议并深入到仓库源码 include/oneapi/tbb/detail/_aggregator.h 剖析其基于“待处理邮箱 单处理器线程”的实现原理。读者读完可以掌握何时用 aggregator 替代 mutex、如何用两种接口保护非并发容器、以及如何在 oneTBB 容器与流图内部看到它的真实应用。一、aggregator 是什么与 mutex 的定位差异oneTBB 官方参考文档 aggregator_cls.rst 对aggregator的定位非常明确Class for mutual exclusion that does not model the Mutex Concept.即它提供互斥执行的能力但并不建模 Mutex Concept。这意味着它没有lock()/unlock()、没有 RAII 风格的scoped_lock也不满足标准互斥量的语义要求关于标准互斥量要求可参见 mutex_cls.rst。它的接口完全不同操作函数体或 lambda 表达式通过execute方法传给aggregator由 aggregator 负责互斥地执行传给同一个 aggregator 对象的操作会被互斥执行mutually exclusiveexecute方法在传入的函数执行完成之后才返回。换句话说传统 mutex 是“你锁住临界区然后自己执行代码”aggregator 是“你把代码交给我我保证它们在同一个时刻只有一个在执行”。这种“任务投递”式的设计天然适合将多个线程对共享数据结构如非并发 STL 容器的操作串行化并且为后续的批量聚合处理把多个操作合并到一次处理中打下了基础。二、基础接口Basic Interface2.1 语法与头文件class aggregator;使用前需要开启预览宏并包含头文件#define TBB_PREVIEW_AGGREGATOR 1 #include oneapi/tbb/aggregator.h关于头文件的现状说明以上头文件路径是参考文档中声明的公开接口。从当前仓库快照看include/oneapi/tbb/下尚未发布独立的aggregator.h聚合机制的核心实现位于 include/oneapi/tbb/detail/_aggregator.h并被concurrent_priority_queue.h、concurrent_lru_cache.h以及 flow_graph 内部模块直接使用详见第六节。因此在实际使用中若你的 oneTBB 版本尚未导出公开头文件可以通过包含这些内部头文件所在的公开容器头如oneapi/tbb/concurrent_priority_queue.h来间接获得该机制或者等待后续版本公开TBB_PREVIEW_AGGREGATOR对应的正式头文件。2.2 类成员namespace oneapi { namespace tbb { class aggregator { public: aggregator(); templatetypename Body void execute(const Body b); }; } // namespace tbb } // namespace oneapi成员说明成员说明aggregator()构造一个aggregator对象。templatetypename Body void execute(const Body b)将b提交给 aggregator使其以互斥方式执行b执行完成后返回。Body可以是函数对象或 lambda只要其可被调用即可调用形式如b()。2.3 示例用 aggregator 保护非并发 priority_queue以下示例源自 basic_interface.rst使用aggregator对非并发的std::priority_queue进行安全操作typedef priority_queuevalue_type, vectorvalue_type, compare_type pq_t; pq_t my_pq; aggregator my_aggregator; value_type elem 42; // push elem onto the priority queue my_aggregator.execute( [my_pq, elem](){ my_pq.push(elem); } ); // pop an elem from the priority queue bool result false; my_aggregator.execute( [my_pq, elem, result](){ if (!my_pq.empty()) { result true; elem my_pq.top(); my_pq.pop(); } } );要点所有对my_pq的修改都被包进 lambda交给同一个my_aggregator互斥执行因此无需为my_pq单独加锁第二个 lambda 通过按引用捕获的result与elem把“是否弹出了元素”和“弹出的值”传回调用线程execute是同步的返回时lambda 已执行完毕捕获的变量对调用线程可见。三、专家接口Expert Interface对于需要更精细控制的场景oneTBB 提供aggregator_ext专家接口见 expert_interface.rst。它的设计目标是把控制权交给用户不再传函数体而是传“操作数据”并由用户自定义的 handler 决定如何处理这些数据。3.1 语法与头文件templatetypename handler_type class aggregator_ext;头文件与基础接口相同#define TBB_PREVIEW_AGGREGATOR 1 #include oneapi/tbb/aggregator.h3.2 类成员namespace oneapi { namespace tbb { class aggregator_operation { public: enum aggregator_operation_status {agg_waiting0,agg_finished}; aggregator_operation(); void start(); void finish(); aggregator_operation* next(); void set_next(aggregator_operation* n); }; templatetypename handler_type class aggregator_ext { public: aggregator_ext(const handler_type h); void process(aggregator_operation *op); }; } // namespace tbb } // namespace oneapi成员说明成员说明aggregator_ext(const handler_type h)构造一个使用 handlerh处理操作的aggregator_ext对象。void process(aggregator_operation* op)将op描述的操作数据提交给aggregator_ext互斥执行op被处理完成后返回。aggregator_operation::aggregator_operation()构造一个基础aggregator_operation对象。void aggregator_operation::start()准备该操作对象被处理处理前必须调用。void aggregator_operation::finish()准备该操作对象被释放回其发起线程处理完成后调用。aggregator_operation* aggregator_operation::next()返回跟在this之后的下一个aggregator_operation。void aggregator_operation::set_next(aggregator_operation* n)将n设为this之后的下一个操作对象。3.3 示例用 aggregator_ext 保护 priority_queue以下示例同样操作非并发的std::priority_queue但改用专家接口typedef priority_queuevalue_type, vectorvalue_type, compare_type pq_t; pq_t my_pq; value_type elem 42; // The operation data, derived from aggregator_node class op_data : public aggregator_node { public: value_type* elem; bool success, is_push; op_data(value_type* e, bool pushfalse) : elem(e), success(false), is_push(push) {} }; // A handler to pass in the aggregator_ext template class my_handler_t { pq_t *pq; public: my_handler_t() {} my_handler_t(pq_t *pq_) : pq(pq_) {} void operator()(aggregator_node* op_list) { op_data* tmp; while (op_list) { tmp (op_data*)op_list; op_list op_list-next(); tmp-start(); if (tmp-is_push) pq-push(*(tmp-elem)); else { if (!pq-empty()) { tmp-success true; *(tmp-elem) pq-top(); pq-pop(); } } tmp-finish(); } } }; // create the aggregator_ext and initialize with handler instance aggregator_extmy_handler_t my_aggregator(my_handler_t(my_pq)); // push elem onto the priority queue op_data my_push_op(elem, true); my_aggregator.process(my_push_op); // pop an elem from the priority queue bool result; op_data my_pop_op(elem); my_aggregator.process(my_pop_op); result my_pop_op.success;注示例中操作节点基类写作aggregator_node与成员列表中的aggregator_operation名称不一致属于文档历史命名二者指同一概念——用户自定义操作类型需从该基类派生在当前仓库实现中对应的内部基类是 detail/_aggregator.h 中的aggregated_operationDerived。3.4 handler 必须遵守的协议专家接口的关键在于 handler 算法必须符合文档约定的协议否则行为未定义handler 收到的是一个aggregator_node链表它必须处理链表中的全部节点节点处理顺序由用户决定但在 handler 返回之前所有节点都必须被处理完链表的所有操作遍历、拼接都应通过next()与set_next()完成不要直接操作内部指针处理单个节点时先调用该节点的start()再执行与该节点关联的操作完成后调用finish()finish()会把节点释放回它的发起线程从而使该线程对process()的调用返回。上面的my_handler_t::operator()正是这一协议的最简实现遍历链表、逐个start()→ 处理 →finish()。四、源码级实现剖析detail/_aggregator.h聚合机制的核心实现在 include/oneapi/tbb/detail/_aggregator.h位于tbb::detail::d1命名空间其设计可以概括为“待处理邮箱mailbox 单处理器线程”。4.1 操作节点基类 aggregated_operationtemplate typename Derived class aggregated_operation { public: // Zero value means wait status, all other values are user specified values std::atomicuintptr_t status; std::atomicDerived* next; aggregated_operation() : status{}, next(nullptr) {} };两个要点status为 0 表示“等待处理”非 0 表示“用户自定义状态”next是原子指针用于在链表/邮箱中串联节点。这正好对应文档专家接口中start/finish与next/set_next的底层支撑。4.2 aggregator_generic邮箱与处理器竞争template typename OperationType class aggregator_generic { public: aggregator_generic() : pending_operations(nullptr), handler_busy(false) {} template typename HandlerType void execute( OperationType* op, HandlerType handle_operations, bool long_life_time true ); private: std::atomicOperationType* pending_operations; // 待处理操作链表邮箱 std::atomicuintptr_t handler_busy; // 是否有线程正在处理 };execute的逻辑对应 源码第 64-98 行先读op-status必须在把操作插入邮箱之前读取因为操作被执行后可能立即失效通过 CAScompare_exchange_strong把op插入pending_operations链表如果插入后链表原先为空自己是第一个则该线程成为“处理器”调用start_handle_operations处理整条链表如果不是第一个则说明已有处理器在跑调用线程只需自旋等待op-status变为非 0处理完成。start_handle_operations对应 源码第 102-133 行保证同一时刻只有一个处理器自旋等待handler_busy 0然后置 1拿到处理权用pending_operations.exchange(nullptr)一次性取走整条待处理链表调用用户 handlerhandle_operations(op_list)处理完释放handler_busy 0。这种“取走整条链表再批量处理”的模式正是 aggregator 相比普通互斥锁的优势所在多个线程提交的操作可以合并成一次批量处理减少调度开销。4.3 long_life_time 语义execute还有一个关键参数long_life_timelong_life_time true操作对象在执行之后仍可被访问默认值因此可以安全等待其完成long_life_time false操作对象可能在执行过程中被销毁执行期间对它的一切访问包括检查 status、等待完成都是未定义行为。文档基础接口的execute对用户隐藏了这一层总是同步等待而内部容器在使用时可以利用“短生命周期”优化。4.4 辅助设施aggregating_functor源码第 156-169 行 提供的aggregating_functorAggregatingClass, OperationList是一个适配器它保存对象指针调用my_object-handle_operations(op_list)。这让容器类可以把“处理操作链表”的成员函数直接包成一个 handler与aggregatorHandlerType, OperationType第 141-152 行配合使用——后者内部持有一个 handler 实例并提供initialize_handler与execute(OperationType* op)。五、仓库中的真实使用场景aggregator 机制并非孤立存在它在 oneTBB 内部多个组件中被实际采用5.1 concurrent_priority_queueinclude/oneapi/tbb/concurrent_priority_queue.h 在内部大量使用聚合机制第 21 行#include detail/_aggregator.h针对 push/pop 等不同操作分别构造 functor 并调用my_aggregator.initialize_handler(functor{this})第 57-124 行通过my_aggregator.execute(op_data)提交操作第 175、183、201 行第 395 行声明容器内嵌的聚合器类型using aggregator_type aggregatorfunctor, cpq_operation;可见并发优先队列正是把每个“操作”包装为继承自aggregated_operation的节点再交给内部的aggregator串行执行从而在无锁队列结构之上保证操作的整体一致性。5.2 concurrent_lru_cache 与 flow_graphinclude/oneapi/tbb/concurrent_lru_cache.h 第 27 行同样包含detail/_aggregator.h用聚合机制串行化对 LRU 结构的管理操作流图内部也复用了该机制例如 include/oneapi/tbb/detail/_flow_graph_indexer_impl.h 第 153 行声明d1::aggregatorhandler_type, indexer_node_base_operation my_aggregator;用于把并发到达的 indexer 输入操作串行化处理。这些内部用法印证了文档的定位aggregator 是一个经过实战检验的“操作聚合串行化”原语而不只是一个 API 草案。六、使用前提与注意事项预览特性文档中的公开接口由TBB_PREVIEW_AGGREGATOR宏控制属于预览功能API 未来可能调整相关预览宏机制见 preview_features.rst。头文件现状当前仓库快照未发布oneapi/tbb/aggregator.h公开头文件聚合机制实现位于 include/oneapi/tbb/detail/_aggregator.h。内部模块通过包含容器头如oneapi/tbb/concurrent_priority_queue.h间接使用。语义边界aggregator 不建模 Mutex Concept不提供lock/unlock与 RAII 作用域锁它只保证“提交给同一 aggregator 的操作互斥执行”。处理协议专家接口handler 必须处理完链表全部节点、按start → 操作 → finish的顺序处理每个节点、用next/set_next操作链表违反协议属于未定义行为。生命周期内部execute的long_life_time参数决定了操作对象执行后是否仍可访问文档级接口总是同步等待但理解该参数有助于解释内部容器的行为。适用场景当临界区操作可以封装为“数据 处理函数”的形态、且希望把多次互斥执行合并为批量处理时典型如并发容器内部操作aggregator 比逐次加锁更契合若只是简单临界区保护标准互斥量见 mutex_cls.rst仍是更直观的选择。七、总结aggregator是 oneTBB 中一个“另类”的同步原语它不建模 Mutex Concept却通过execute基础接口与process 自定义 handler专家接口提供了互斥且可批量聚合的执行模型。其底层实现detail/_aggregator.h以原子邮箱 处理器自旋竞争的方式保证同一时刻只有一个线程在处理提交的操作链表并在 oneTBB 的concurrent_priority_queue、concurrent_lru_cache和 flow_graph 内部得到实际应用。理解它等于掌握了一种“把临界区操作变成可聚合任务”的设计思路在实现高性能并发数据结构时非常有用。赞分享并发编程高性能计算【免费下载链接】oneTBBoneAPI Threading Building Blocks (oneTBB)项目地址https://gitcode.com/gh_mirrors/on/oneTBB点击查看免费下载相关推荐oneTBB aggregator 类全解析互斥执行的新范式与操作聚合实战oneTBB aggregator 类全解析互斥执行的新范式与操作聚合实战 导读 aggregator 是 oneTBBoneAPI Threading B开发工具构建工具系统编程oneTBB aggregator 基础接口详解基于操作聚合的互斥执行机制oneTBB aggregator 基础接口详解基于操作聚合的互斥执行机制 本文以 oneTBBoneAPI Threading Building Bloc开发工具构建工具系统编程oneTBB parallel_pipeline 之 filter 类模板强类型过滤器的构造、组合与执行模式详解oneTBB parallel_pipeline 之 filter 类模板强类型过滤器的构造、组合与执行模式详解 filter 是 oneAPI Thread并发编程高性能计算上一篇使用 ANTLR4 Dart 运行时从安装、代码生成到解析调试的完整指南下一篇使用 Python imaplib 与 BeautifulSoup 将 IMAP 邮件批量导出为 CSVStore_emails_in_csv 实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进