- 并发编程
- 高性能计算
【免费下载链接】oneTBB
oneAPI Threading Building Blocks (oneTBB)
导读
oneapi::tbb::task_scheduler_observer是 oneAPI Threading Building Blocks(oneTBB)任务调度器对外暴露的一组回调钩子:它允许用户在某个任务调度 arena 中,精确感知"哪些线程开始/停止参与任务处理",从而在线程进入与离开的瞬间注入自定义逻辑。本文以官方参考文档为骨架,结合本仓库头文件、调度器实现源码与测试用例,完整讲解其接口语义、生命周期管理、底层代理(proxy)机制,并给出线程亲和性(pinning)绑定的可运行实战示例。读完本文,你将能够独立编写自己的 observer 子类,并正确把握observe()、on_scheduler_entry、on_scheduler_exit的调用时机与并发约束。
一、task_scheduler_observer 是什么
task_scheduler_observer表示"线程对任务调度服务的兴趣"。它不是一个在后台默默运行的监视器,而是一个由用户派生、注册进指定 arena 的观察者对象:当线程进入或离开该 arena 的任务处理时,调度器会回调其虚方法。
官方参考文档 task_scheduler_observer_cls.rst 给出的类骨架如下:
// Defined in header <oneapi/tbb/task_scheduler_observer.h> namespace oneapi { namespace tbb { class task_scheduler_observer { public: task_scheduler_observer(); explicit task_scheduler_observer( task_arena& a ); virtual ~task_scheduler_observer(); void observe( bool state=true ); bool is_observing() const; virtual void on_scheduler_entry( bool is_worker ) {} virtual void on_scheduler_exit( bool is_worker ) {} }; } // namespace tbb } // namespace oneapi典型用法是:从task_scheduler_observer派生自己的观察者类,重写on_scheduler_entry或on_scheduler_exit,然后调用observe(true)激活观察。观察在对象创建时默认处于禁用状态,必须显式调用observe()才会生效。
需要注意一个重要的安全约定:在重写的方法中抛出且未被捕获的异常属于未定义行为。底层调度器在调用回调时不会拦截异常,相关注释也明确写道"不拦截可能从回调中逃逸的任何异常,让它们要么被 TBB 调度器处理、要么交给调试器"(见 observer_proxy.cpp)。
二、两种观察者:本地观察者与 arena 观察者
根据构造函数的不同,observer 分为两类,它们的回调触发范围不同:
| 构造方式 | 观察范围 | 触发时机 |
|---|---|---|
task_scheduler_observer()(无参) | 本地观察者,观察当前线程所属 arena | 每当一个 worker 线程加入/离开"所有者线程所在 arena"时触发 |
explicit task_scheduler_observer(task_arena& a) | arena 观察者,观察指定 arena | 每当一个线程加入/离开该 arena 时触发 |
两种构造函数都会将对象置于"非活动(观察禁用)"状态。区别在于:无参版本的对象在激活后,仅观察创建它的那个线程所在 arena 的进出事件;带task_arena&参数的版本则与某个具体 arena 绑定,与创建线程无关。
从实现看,这个绑定关系存放在头文件内部的my_task_arena指针中(见 task_scheduler_observer.h),无参构造时该指针为nullptr,表示"跟随当前线程的 arena"。
三、观察的生命周期:observe() 与 is_observing()
3.1 observe(bool state)
observe(true)启用观察,observe(false)停用观察。它有几个从源码中可以确认的行为细节:
- 重复调用同状态是 no-op:头文件注释明确说明 "Repeated calls with the same state are no-ops"(见 task_scheduler_observer.h)。
- 并发调用不安全:头文件警告 "concurrent invocations of this method are not safe"(见 task_scheduler_observer.h),即不要从多个线程同时对同一个 observer 调用
observe()。 - 本地观察者(无参版本)的调用前提:只能在"当前线程已初始化任务调度器或已附着到某个 arena"时使用
observe()(见 task_scheduler_observer.h)。
3.2 is_observing() const
返回true表示观察已启用,false表示未启用。其实现非常轻量:直接检查内部代理指针my_proxy是否为非空(见 task_scheduler_observer.h)。也就是说,观察状态与代理对象(proxy)是否挂载在 arena 的观察者列表上严格对应。
3.3 析构函数与并发安全
析构函数会自动停用观察:若内部my_proxy非空,则先调用observe(false)(见 task_scheduler_observer.h)。参考文档进一步说明:析构会等待所有正在执行的on_scheduler_entry/on_scheduler_exit调用完成之后才销毁实例。这依赖头文件中的my_busy_count原子计数:调度器每回调一次就递增,回调结束递减,observe(false)通过spin_wait_until_eq(tso.my_busy_count, 0)自旋等待归零(见 observer_proxy.cpp)。
因此头文件建议:在派生类析构开始之前就停用观察,否则可能出现"对象已部分析构、回调仍在并发触发"的危险场景(见 task_scheduler_observer.h)。
四、回调语义详解
4.1 on_scheduler_entry(bool is_worker)
调度器在以下时机为每个线程调用一次:
- 线程开始参与 oneTBB 任务处理时;
- 观察启用之后,线程进入 arena 时;
- 对于观察启用时已经在执行任务的线程,会在其启用后执行第一个被窃取(stolen)任务之前调用。
is_worker标志的含义:为true表示该线程由 oneTBB 自己创建(worker 线程);为false表示外部线程(如应用主线程或用户自建线程)。
文档给出了一个重要保证:如果一个线程启用观察后派生任务,那么该任务及其创建的所有任务,都会由"已经调用过on_scheduler_entry"的线程来执行(见 task_scheduler_observer_cls.rst)。这正是很多场景(如设置线程亲和性)能依赖该回调"覆盖所有后续执行者"的理论基础。
默认行为是什么都不做。
4.2 on_scheduler_exit(bool is_worker)
调度器在线程停止参与任务处理或离开 arena时调用。同样携带is_worker标志。默认行为什么都不做。
文档特别给出一个 caution(警告):进程不会等待 worker 线程清理完毕,因此进程可能在on_scheduler_exit被调用之前就终止。这意味着:不能把"必须执行"的收尾逻辑(例如释放独占资源、写日志刷盘)寄托在on_scheduler_exit上——它不保证在进程退出前一定会被调用。
4.3 与 task_arena、this_task_arena 的关系
observer 的回调参数只有is_worker,若需要在回调里知道当前线程在 arena 中的槽位,可以使用oneapi::tbb::this_task_arena::current_thread_index()。这在本仓库源码中就有真实案例:库内部用于 NUMA 绑定的numa_binding_observer在on_scheduler_entry中调用apply_affinity_mask(my_binding_handler, this_task_arena::current_thread_index())(见 arena.cpp)。
相关类与命名空间的完整接口可参阅同目录文档:task_arena_cls.rst 与 this_task_arena_ns.rst。
五、底层实现原理:proxy 与 observer 列表
理解底层机制有助于把握回调的并发语义与性能特征。调度器并不直接把 observer 对象挂在链表上,而是为其创建一个代理对象observer_proxy,再把这些 proxy 组织进 arena 的observer_list(见 observer_proxy.h 的注释:"为维护 observer 的共享列表,调度器先把每个 observer 包装成 proxy,使列表项在用户代码销毁 observer 对象后仍然有效")。
关键数据结构(见 observer_proxy.h):
observer_proxy:持有指向 observer 的指针my_observer、引用计数my_ref_count、前后指针my_next/my_prev;observer_list:用spin_rw_mutex保护的双向链表,my_head/my_tail为链表端点,每个 arena 持有一个自己的my_observers列表。
observe(true)的内部流程(见 observer_proxy.cpp):
- 若
my_proxy为空,new 一个observer_proxy并存入my_proxy; - 确定目标列表:无参构造的本地 observer 取当前线程所属 arena 的
my_observers;带 arena 参数的 observer 取该 arena 的列表(若 arena 尚未初始化则先initialize()); - 把 proxy
insert进列表; - 若当前线程就在该 arena 内,立即补发一次 entry 通知(
notify_entry_observers),保证"激活时已在 arena 中的线程"也能收到回调。
observe(false)的内部流程(见 observer_proxy.cpp):
- 通过
my_proxy.exchange(nullptr)摘掉代理并独占它; - 在写锁下把
proxy->my_observer置空、递减引用计数,必要时从列表移除并销毁 proxy; - 自旋等待
my_busy_count归零,确保其他线程正在执行的回调全部结束。
回调的分发在do_notify_entry_observers/do_notify_exit_observers中完成(见 observer_proxy.cpp)。两个值得注意的实现要点:
- 回调调用时不持有任何列表锁:只在推进链表指针时短暂持锁,真正调用用户代码
tso->on_scheduler_entry(worker)时已释放锁。这避免了用户在回调里做重活时阻塞整个 arena 的观察者列表。 - 引用计数保护:回调执行期间 proxy 和 observer 的存活由
my_ref_count、my_busy_count共同保障,这也是为什么observe(false)/析构可以安全地与进行中的回调并发。
调度器在哪些路径上触发通知?从源码看至少包括:worker 进入 arena 时(arena.cpp)、离开 arena 时(arena.cpp)、外部线程通过task_arena::execute进入/离开时(arena.cpp 与 arena.cpp)、以及线程结束任务处理时(governor.cpp)。
六、实战示例:把 worker 线程固定(pin)到硬件线程
参考文档给出的经典场景,是编写一个 observer,将 oneTBB worker 线程绑定到指定的硬件线程(HW 线程),并在退出时恢复原有亲和性。下面是官方示例的完整版本:
class pinning_observer : public oneapi::tbb::task_scheduler_observer { public: affinity_mask_t m_mask; // HW affinity mask to be used for threads in an arena pinning_observer( oneapi::tbb::task_arena &a, affinity_mask_t mask ) : oneapi::tbb::task_scheduler_observer(a), m_mask(mask) { observe(true); // activate the observer } void on_scheduler_entry( bool worker ) override { set_thread_affinity(oneapi::tbb::this_task_arena::current_thread_index(), m_mask); } void on_scheduler_exit( bool worker ) override { restore_thread_affinity(); } };要点拆解:
- 构造函数里绑定 arena 并立即激活:
task_scheduler_observer(a)将观察范围限定为 arenaa,observe(true)使观察从构造起就生效; - 入口回调设置亲和性:
on_scheduler_entry中依据this_task_arena::current_thread_index()取得当前线程在 arena 中的槽位,结合affinity_mask_t掩码设置线程亲和性; - 出口回调恢复状态:
on_scheduler_exit中调用restore_thread_affinity()还原,避免影响线程后续的其他用途。
这个模式在 oneTBB 库内部有同款实现,可作为对照参考:numa_binding_observer(arena.cpp)在on_scheduler_entry中按this_task_arena::current_thread_index()应用亲和掩码,在on_scheduler_exit中恢复,并在析构时先observe(false)再销毁绑定处理器(arena.cpp)。可见"entry 设置、exit 恢复、析构前停用"是该场景的标准三步曲。
七、测试用例印证:回调语义的验证
仓库测试对 observer 的回调时机与约束做了系统性验证,是理解语义的绝佳材料:
- test_task_arena.cpp 中的
ArenaObserver(第 124 行起)重写on_scheduler_entry,验证:current_thread_index()不超过 arena 的最大并发度;is_worker为 true 时,槽位索引不小于保留槽位数量myNumReservedSlots(即 worker 不会占用被保留的槽位);- 通过线程局部存储记录"上一次进入的 arena id",检测是否发生同一线程重复进入同一 arena 的异常情况。
- 同一文件中的
LocalObserver(第 189 行起)是一个激活顺序回归测试:验证"arena 观察者可以先于本地观察者被激活"的合法场景。 - test_arena_constraints.cpp 则直接使用库内部的 NUMA 绑定 observer 测试在 NUMA 与 core type 约束下亲和性设置与
default_concurrency的正确性。
这些测试共同确认了文档承诺的语义:entry 回调一定发生在该线程执行(被窃取的)任务之前、is_worker标志区分线程来源、以及 observer 与 arena 的绑定关系。
八、使用建议与注意事项汇总
- 创建后必须
observe(true),否则观察不会生效;重复调用同一状态是 no-op。 - 不要在多个线程上并发调用
observe()。 - 回调内抛异常是未定义行为——回调体应保证不抛异常,或自行捕获。
- 不要在
on_scheduler_exit里寄托必须完成的收尾:进程退出前不保证 worker 线程会走完 exit 回调。 - 优先在派生类析构前调用
observe(false),避免对象半析构时回调仍在并发执行。 - 回调执行期间不持有调度器内部锁,但回调本身的耗时仍会串行化在同一 arena 的观察者列表推进上,回调体应保持轻量。
- 需要线程槽位信息时,在回调内配合
this_task_arena::current_thread_index()使用;需要调整全局并发度等调度参数时,可参考 global_control_cls.rst。
九、更多资料
- 官方类参考:task_scheduler_observer_cls.rst
- 头文件实现:include/oneapi/tbb/task_scheduler_observer.h
- 底层实现:observer_proxy.cpp、observer_proxy.h
- 库内真实用例(NUMA 绑定):arena.cpp
- 测试用例:test_task_arena.cpp、test_arena_constraints.cpp
- 相关主题:task_arena_cls.rst、this_task_arena_ns.rst
- 并发编程
- 高性能计算
【免费下载链接】oneTBB
oneAPI Threading Building Blocks (oneTBB)
相关推荐
mold 仓库内嵌 oneTBB 的 task_scheduler_observer 深度解析:线程进入与离开调度器 Arena 的观察机制
mold 仓库内嵌 oneTBB 的 task_scheduler_observer 深度解析:线程进入与离开调度器 Arena 的观察机制 导读 本文以 mo
开发工具构建工具系统编程Java多线程进阶:自定义线程池与任务调度器的实现
Java多线程进阶:自定义线程池与任务调度器的实现 Java多线程编程是现代软件开发的核心技能之一。在并发编程中, 线程池 和 任务调度器 是两个至关重要的概念
文档教程知识库oneTBB static_partitioner 详解:均匀静态划分、确定性线程亲和与源码级实现
oneTBB static_partitioner 详解:均匀静态划分、确定性线程亲和与源码级实现 static_partitioner 是 oneAPI Th
并发编程高性能计算
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考