Mojo 项目 AsyncRT 并发内核解析:WorkQueue 线程池的设计、任务路由与设备亲和调度
【免费下载链接】mojoThe Modular Platform (includes MAX & Mojo)项目地址: https://gitcode.com/GitHub_Trending/mo/mojo
M::AsyncRT::WorkQueue是 Modular 平台(MAX & Mojo)运行时 AsyncRT 中管理 CPU 并行度的核心抽象,它以一个可替换的线程池接口统一了任务提交、等待与线程捐献机制。本文以 AsyncRT/docs/WorkQueue.md 为骨架,结合仓库内接口声明、实现源码与单元测试,完整讲解 WorkQueue 的创建方式、线程模型、四层任务队列、CPU 亲和性分配、设备任务固定(Device Task Pinning)以及非阻塞设计原则,帮助你理解 Mojo 运行时如何在多核与多 GPU 环境下高效调度工作项,并掌握通过环境变量、modular.cfg与 CLI 参数进行调优的实战方法。
一、WorkQueue 概览:AsyncRT 的 CPU 并行度抽象
M::AsyncRT::WorkQueue是一个用于并发执行工作项(work item)的抽象接口,是 AsyncRT 管理 CPU 并行度的核心抽象。它以线程池的形式将任务分发到可用的 CPU 核上,同时刻意保持接口极简:客户端只通过addTask()(提交任务)、addLocalTask()(提交本地任务)和await()(等待若干值就绪)与之交互。
从源码结构看,这一设计意图体现在 AsyncRT/include/AsyncRT/Runtime/WorkQueue.h 的类声明中:WorkQueue是一个纯虚基类,通过addTask、addLocalTask、await、shutdown等虚函数定义契约,而具体实现(单线程队列、线程池队列、NUMA 分区委托队列)均通过工厂函数创建,客户端不直接构造。
该接口与上层M::AsyncRT::CPUDevice(“god object”,统一组织线程池、内存分配器等运行时资源)解耦,策略可插拔:CPUDeviceOptions::WorkQueueType支持kSingleThread与kThreadPool两种队列类型,具体可见 AsyncRT/include/AsyncRT/Runtime/CPUDevice.h。工作项本身由WorkItem结构承载,内部持有llvm::unique_function<void()>任务函数,并在启用 Tracy 时携带唯一任务 ID 用于性能分析(见 WorkQueue.h)。
创建工作队列(工厂函数)
WorkQueue 不能直接构造,必须通过工厂函数创建:
createSingleThreadWorkQueue(cpuDevicePtr):创建一个只使用调用(捐献)线程的 WorkQueue,零同步开销,适用于单线程平台。对应实现 SingleThreadWorkQueue.cpp 中SingleThreadWorkQueue类——它不派生任何额外线程,但接口本身线程安全,addTask与await可从任意线程调用。createThreadPoolWorkQueue(cpuDevicePtr, numThreads, maxThreads, mainWillDonate, withAffinity, threadBusyWaitTime, poolName):创建多线程 WorkQueue,参数含义如下:
| 参数 | 说明 |
|---|---|
numThreads | 工作线程数量;为 0 时根据系统配置自动决定(见下文“Worker 分配与 CPU 亲和性”) |
maxThreads | 自动探测时的numThreads上限;为 0 时忽略,源码中超过kMaxWorkers(1024)也会被钳制 |
mainWillDonate | 为 true 时,创建线程将在await()期间参与处理工作项(当前默认值) |
withAffinity | 为 true 时,工作线程固定(pin)到特定 CPU 核 |
threadBusyWaitTime | 空闲时自旋(spin)后再休眠的时长 |
poolName | 线程名前缀(在调试器/性能分析器中可见) |
除了文档中的这两个工厂函数,仓库还提供了两个高级变体:WorkQueue.h 中的createPartitionedThreadPoolWorkQueue(将线程池限制在单个 NUMA 节点的 CPU 核上,必须由DelegateThreadPoolWorkQueue托管)与createDelegateThreadPoolWorkQueue(将一组分区队列委托包装为单一统一队列)。
在mainWillDonate语义上,工厂函数有细致考量:若为 false,则创建numThreads个 worker,适合多个请求线程共享同一队列的多线程服务器;若为 true(默认),则只创建numThreads - 1个 worker,假定创建线程最终会调用await并“捐献”自己参与处理,适合 REPL 或执行工具这类由单个主线程驱动的系统。
二、线程模型:Worker、Main 与 Foreign 线程
WorkQueue 区分三类线程(术语注释同样出现在 ThreadPoolWorkQueue.cpp 实现文件头部):
- Worker 线程(Worker threads):由 WorkQueue 创建、运行专用工作处理循环的线程,每个 worker 拥有唯一的
workerID(0 到 N-1)。Worker 启动时会依次设置 TLS 中的localWorkerID、注册当前 CPUDevice、为线程命名(poolName + workerID,通过llvm::set_thread_name)并按需设置 CPU 亲和性(ThreadPoolWorkQueue.cpp)。 - Main 线程:当
mainWillDonate为 true 时,创建 WorkQueue 的线程被指定为 “main” 线程(workerID 0),在await()期间参与工作处理,并且必须是调用shutdown()的线程。实现中 main 线程的WorkQueueThread不创建底层线程,而是在runItemsOnOwningThread中通过runWithThreadAffinity临时设置亲和性后处理工作。 - Foreign 线程:其他任何与 WorkQueue 交互的线程。可以调用
addTask()和await(),但不会捐献自身去处理工作项。在mainWillDonate为 false 时,foreign 线程也可以调用shutdown()。
线程身份的判断通过 TLS 与线程 ID 比对完成:getOwningWorkQueueThread()先读 TLS 中的localWorkerID,再比对worker->threadID != llvm::get_threadid(),不一致即视为 foreign 线程(ThreadPoolWorkQueue.cpp)。
三、Worker 分配与 CPU 亲和性
默认线程数的确定规则
当numThreads为 0 时,createThreadPoolWorkQueue内部调用getThreadAffinityCpuIds()(实现于 AsyncRT/lib/Support/ThreadAffinity.cpp),按以下优先级确定线程数:
- P-core/E-core 不均衡的系统:使用性能核(performance cores)数量。代码注释提到,若物理核与性能核数量不一致,出于对混合架构(如 x86 大小核)调度问题的规避,会优先锁定 P 核;在 macOS(
__APPLE__)上则关闭亲和性、把线程数设为性能核数,交由操作系统调度。 - 开启亲和性(withAffinity):使用物理核数(不含超线程)。
- 未开启亲和性:使用逻辑核数(含超线程)。
CPU 选择与 cgroup 限制
- CPU 选择:亲和性开启时,
CPUSystemInfo::getPreferredCpuIDs()(声明于 Support/include/Support/Threading/HWInfo.h)决定使用哪些 CPU,其启发式优先级为:优先不同虚拟核 → 优先不同物理核 → 优先同一 socket 内的物理核,天然倾向于“性能核优先于能效核、物理核优先于超线程、尽量落在同一 NUMA 节点内”。 - cgroup 限制:容器环境下,线程数会被自动钳制为
max(1, millicores / 1000),即 CPU 限额(千分之一核为单位)对应的核数,保证不会超出容器配额。 maxThreads上限:自动探测的线程数最终还会被maxThreads封顶(maxThreads为 0 或超过 1024 时按 1024 处理)。
亲和性设置
每个 worker 线程启动时调用AsyncRT::setThreadAffinity(cpuID)(对应M::setThreadAffinity),将线程绑定到指定 CPU 核,并在支持的平台上设置内存策略以优先在该 CPU 所在 NUMA 节点的内存上分配(ThreadAffinity.cpp)。kNoAffinity(~0)表示不设置亲和性。
需要说明的是,亲和性默认是关闭的:CPUDeviceOptions::withAffinity默认读取环境变量MODULAR_ENABLE_AFFINITY,仅当其为真值时才开启(见 CPUDevice.h),原因如注释所述——多进程场景下强制亲和性会带来性能问题。
四、任务队列层级:本地 → 亲和 → 全局 → 溢出
WorkQueue 采用多级任务队列来平衡执行效率与工作分发(对应实现见 ThreadPoolWorkQueue.cpp 的WorkQueueThread结构):
- 本地任务列表(
localTaskList):每 worker 一个、无同步的列表,仅由属主线程通过addLocalTask()写入。优先级最高,适合短小的延续任务(如 AsyncValue 的 waiter)——这类任务如果走线程切换,开销会远超执行本身。实现中runItemsImpl每次循环先尽量清空localTaskList(doWork<IsWaiter=true>)。 - 亲和任务列表(
affinityTaskList):每 worker 一个无锁环形缓冲区(LockFreeRingBuffer<WorkItem>,容量为每线程 1024 槽位)。当addTask()传入非负taskId时使用(典型来源是 Mojo 的async_parallelize),任务由taskId指定的 worker 处理,形成缓存友好的执行模式。若环形缓冲区已满,则溢出到受互斥锁保护的localSpillQueue,并在后续从localTaskList恢复执行以维持亲和性。 - 全局任务列表(
taskList):所有 worker 共享的无锁 MPMC 队列(MoodyCamel::ConcurrentQueue),用于addTask()且taskId = kDefaultTaskId(-1)的任务,任何 worker 都可以出队处理。kDefaultTaskId常量定义在 WorkQueue.h,注释明确:所有非async_parallelize来源的任务默认走全局队列。 - 溢出任务列表(
overflowTaskList):互斥锁保护的回退队列,仅当全局队列已满时使用。worker 在即将休眠前检查它。addTask的实现注释详细列出了“队列已满”时的四种可选方案(当前线程内联执行有栈溢出风险且违反“绝不立即执行”契约、推入本地列表破坏负载均衡、无锁队列动态扩容难以保持 push/pop 独立性),最终选择了“推到溢出列表”这一经典方案,将互斥开销仅付给不常见路径(ThreadPoolWorkQueue.cpp)。
工作项按本地 → 亲和 → 全局 → 溢出的优先级处理。runItemsImpl主循环每次迭代依次尝试:本地列表 → 亲和队列 → 全局队列;自旋阶段与预休眠阶段还会再次检查亲和与全局队列,最后才进入休眠(ThreadPoolWorkQueue.cpp)。
单元测试 AsyncRT/unittests/WorkQueueTest.cpp 验证了这套路由语义:
TaskIdRouting测试(第 59 行起):4 个 worker(0-3)下使用 taskId 1、2、3 提交任务,断言同一 taskId 的所有任务都运行在同一线程上,且三个 taskId 落在三个不同线程上——即亲和路由生效且 worker 0 被保守规避;NegativeTaskId测试(第 132 行起):以 -5 提交任务仍能正确执行,验证负数 taskId 走全局队列;TaskIdWithMainWillDonate测试(第 152 行起):在mainWillDonate模式下任务依然全部完成,无需主线程额外 await。
五、所有权与生命周期:创建 → 使用 → shutdown → 销毁
WorkQueue 通常由M::AsyncRT::CPUDevice实例拥有,CPUDevice 根据CPUDeviceOptions创建并管理它。生命周期四阶段:
- 创建:通过
createThreadPoolWorkQueue()或createSingleThreadWorkQueue()。线程池队列的构造函数完成时,所有 worker 线程已启动并进入runItems循环。 - 使用:客户端通过
addTask()/addLocalTask()添加工作,通过await()等待结果。 - shutdown:销毁前必须调用
shutdown()。ThreadPoolWorkQueue::shutdown()(ThreadPoolWorkQueue.cpp)依次完成:- 主线程(
mainWillDonate时)捐献自身帮助排空剩余工作项; - 设置
doneFlag通知 worker 退出; - 对每个 worker 的信号量执行
post(),唤醒所有休眠线程; - 清零
suspendedThreads位向量,防止 in-flight 的andThenSync唤醒正在被 join 的线程; - join 所有 worker 线程。
- 主线程(
- 销毁:
shutdown()返回后即可销毁;析构函数断言全局队列已空(assert(!taskList.try_dequeue(workItem))),并清理线程局部 CPUDevice 指针。
注意shutdown()的调用方约束:mainWillDonate模式下必须由 main 线程调用(实现中有assert),否则必须由 foreign 线程调用;await()返回也不意味着所有依赖资源都可销毁——只有shutdown()能保证所有在途计算完成(WorkQueue.h 的 CAUTION 注释)。
六、空闲行为:指数退避自旋、休眠与唤醒
当 worker 无任务可处理时:
- 忙等阶段:以指数退避(exponential backoff)自旋
busyWaitTime(默认 1ms),期间持续检查新任务。实现使用BusyWaitSpinWaiter,避免在空转时冲击正在做有用工作的线程的内存层级(ThreadPoolWorkQueue.cpp)。 - 溢出检查:休眠前,把
localSpillQueue/overflowTaskList中的任务泵入主队列(无公平性保证,但休眠在即,值得付出互斥锁代价)。 - 休眠:在共享位向量
suspendedThreads中标记自身挂起,然后阻塞在自己的信号量sema上。 - 唤醒:
addTask()发现挂起 worker 时,post 对应信号量。
这里有一个精巧的竞态处理:markSuspended(worker 侧)与takeSuspended(调度侧)之间存在两种交错顺序,实现通过“标记挂起后、休眠前再尝试一次出队”的方式兜底,避免“任务已入队但所有 worker 已休眠”的丢失唤醒问题(ThreadPoolWorkQueue.cpp)。
超过 64 个 worker 的组播方案:挂起位向量是 64 位无符号整数,本限制 64 线程;现代服务器 NUMA 节点常超过 64 核,实现采用“位 = worker 组”的组播(multicast)方案——每个位代表2^multicastFactor个 worker。代价是唤醒时可能向组内所有 worker post 信号量(可能产生虚假唤醒),但换来的是在不知道确切挂起线程时也能保证唤醒正确性;代码注释坦诚表示“模型执行期间不应频繁休眠/唤醒,因此这个代价可以接受”(ThreadPoolWorkQueue.cpp)。
非阻塞await的语义
await()不阻塞:它把客户端线程“捐献”给工作队列去运行任务,直到目标值全部就绪。实现(ThreadPoolWorkQueue.cpp)对每个值注册andThenSync回调递减计数,计数归零时 post 当前 worker 的信号量;worker/main 线程走“边运行边等待”的runItemsOnOwningThread路径,foreign 线程则阻塞在自身独立信号量上(这正是“每线程独有信号量”设计的原因——若用共享信号量或单独的 await 信号量,就无法精准唤醒正在等待特定值的线程)。await可被递归调用,即任务内部可再次await,但文档建议优先使用 AsyncValue 进行同步。
七、关键设计原则:非阻塞、绝不立即执行、线程捐献
- 非阻塞假设:工作项不应阻塞(详见姊妹文档 AsyncRT/docs/WorkQueueNonblocking.md)。阻塞会隐式“抽走”线程池中的线程,导致机器过载或欠载。AsyncRT 的策略不是做自适应线程池(文档分析了 GCD 式自适应池的五个问题:线程数远超核数、资源耗尽边角案例、复杂度失控、遗留代码不协作、缺乏迁移激励),而是:让实现假设工作项不阻塞以保持简单高效;将可能阻塞的操作放到外部运行时,完成后再把完成回调作为工作项提交回 WorkQueue;由上层 CPUDevice 在粒度上平衡 WorkQueue 与外部运行时;以“原生路径更高效”形成迁移激励;同时不做强制检查(
printf、std::mutex理论可阻塞,但微小临界区足够安全)。该文档还指出,AsyncRT 目前缺失一个可移植的异步 I/O 子系统(Linux AIO 不佳、Windows 有边角案例、新内核 io_uring 合适、嵌入式甚至不需要),未来会作为可选组件构建。 - 绝不立即执行:
addTask()永远不在调用线程内联执行任务,任务总是被推迟。这防止栈溢出并保证行为可预测。addLocalTask()同样承诺“绝不立即运行”,只是“尽量在当前线程稍后执行”(WorkQueue.h)。 - 线程捐献:worker/main 线程调用
await()时捐献自身去处理工作项,既避免死锁又最大化 CPU 利用率。shouldRunInlineForTask(taskId)则是“同步内核 + 设备亲和”路径上避免无谓入队的轻量检查——它通过 TLS 中当前线程的workerID与目标taskId比对,命中则直接内联执行(ThreadPoolWorkQueue.cpp)。
八、设备任务固定(Device Task Pinning):将 GPU 任务绑定到同一 NUMA 节点的 CPU
为什么需要固定设备任务
当执行与 GPU 或其他加速器交互的内核时,把任务运行在与设备同 NUMA 节点的固定 CPU 线程上是有利的:
- NUMA 局部性:GPU 挂接在特定 PCIe 总线上,属于特定 NUMA 节点;设备任务运行在同节点 CPU 上可最小化内存访问延迟并最大化 PCIe 带宽。
- 一致的线程亲和性:GPU 驱动上下文(CUDA/HIP context)常有线程局部状态;设备操作始终由同一线程执行可避免驱动内的上下文切换开销。
- 可预测的调度:把设备任务固定到特定 worker,可避免 GPU 绑定工作与 CPU 绑定工作争抢同一线程。
taskId 的确定:三级优先级
当mgp_generic_execute执行引用加速器DeviceContext的内核时,运行时确定一个taskId将任务路由到特定 worker 线程,选择顺序如下:
1. 显式配置(最高优先级)
通过runtime.device_task_cpu_ids配置选项显式指定每个设备对应的 CPU 核:
环境变量:
export MODULAR_RUNTIME_DEVICE_TASK_CPU_IDS="0,32,1,33,2,34,3,35"modular.cfg:
[runtime] device_task_cpu_ids = 0,32,1,33,2,34,3,35列表按设备 ID 索引:以上配置即 设备 0 → CPU 0、设备 1 → CPU 32、设备 2 → CPU 1、设备 3 → CPU 33……(即 worker 的cpuID为对应值)。适用于自动 NUMA 检测失效或需要细粒度控制映射的场景。
2. 自动 NUMA 拓扑检测
若未提供显式配置,且工作队列的 worker 数不少于物理核数,运行时尝试为每个 GPU 推断最优 CPU:
- 查询
NUMATopology::get()获取系统 NUMA 布局; - 通过
device->getPciBusId()获取每个 GPU 的 PCI 总线地址; - 通过
NUMATopology::getNumaNodeForPciBus()把 PCI 总线映射到 NUMA 节点; - 通过
getCpuIdsForNumaNode()获取该 NUMA 节点的 CPU ID 列表; - 为该设备选择该 NUMA 节点第一个可用 CPU。
该映射在启动时计算一次,并缓存在静态gpuToCpuCoreMapping中。底层 API 的声明可在 Support/include/Support/Threading/HWInfo.h 的NUMATopology结构中看到,它维护cpuIdsPerNumaNode、pciBusToNumaNode等映射表,并采用 C++11 静态初始化实现进程内线程安全的首次查询与缓存。
3. 回退:轮询分配(最低优先级)
若 NUMA 检测失败或线程池受限(如 cgroup 约束),回退到简单轮询:
taskId = 1 + (deviceHint % (numWorkers - 1))该公式刻意跳过 worker 0,以避免mainWillDonate为 true 时的潜在停滞(worker 0 是 main 线程,可能不总在处理工作项)。这一“保守规避 worker 0”的约定同样体现在单元测试 WorkQueueTest.cpp 的注释中(“With conservative worker 0 avoidance: taskId = 1 + (hint % 3)”)。
内联 vs. 队外执行
确定taskId后,运行时决定内联执行内核还是派发到亲和队列:
- 无设备亲和性的同步内核(
taskId == kDefaultTaskId):在当前线程内联运行。 - 有设备亲和性的同步内核:检查
shouldRunInlineForTask(taskId)——若已处于正确的 worker 上则内联运行,否则派发到该 worker 的亲和队列并等待。 - 异步内核:总是派发到亲和队列(无亲和性时走全局队列)执行。
调优:autotune_gpu_numa 工具
WorkQueue 文档推荐的调优路径是utils/benchmarking/tools/autotune_gpu_numa工具,它通过基准测试各 NUMA 节点的 GPU 内核启动延迟,输出推荐的 CPU→GPU 映射配置。该工具依赖 Pythonclick模块,需先安装:
bazelw run //utils/benchmarking/tools/autotune_gpu_numa:autotune_gpu_numa脚本会依据 GPU 内核启动基准结果推断最优的 GPU→CPU NUMA 节点(进而到 CPU 核)映射,并给出如何在modular.cfg或环境变量中使用结果的说明。这对多 GPU 系统尤其有用——复杂的 PCIe 拓扑下,默认 NUMA 检测未必最优。需注意,该工具路径存在于官方文档描述中,但未包含在当前仓库快照内,实际使用请以你的工作区为准。
九、运行时调优入口汇总
围绕 WorkQueue,仓库提供了多层次的配置入口,便于在不同场景下调节线程池行为:
环境变量(实现于 CPUDevice.h):
| 环境变量 | 作用 |
|---|---|
MODULAR_THREAD_BUSY_WAIT_US | 覆盖忙等时长(微秒),默认 200 |
MODULAR_ENABLE_AFFINITY | 开启线程 CPU 亲和性(默认关闭) |
MODULAR_RUNTIME_DEVICE_TASK_CPU_IDS | 显式指定设备任务固定到的 CPU 核列表 |
CLI 选项(声明于 AsyncRT/include/AsyncRT/Runtime/RuntimeCLOptions.h,供使用 AsyncRT 的工具共享):
| 选项 | 说明 |
|---|---|
--workqueue {single-thread,thread-pool} | 选择 WorkQueue 类型 |
--num-threads N | 线程数,0 表示启发式自动选择 |
--max-threads N | 自动配置时对 num-threads 的上限 |
--thread-busy-wait-time-us N | 线程休眠前自旋的微秒数,0 表示永不自旋 |
--cpu-affinity | 开启线程亲和性,优先级高于MODULAR_ENABLE_AFFINITY环境变量 |
--allocator {malloc,tcmalloc,leak-checker,profiler,use-after-free} | 选择分配器 |
--time-profile <base> | 输出性能分析文件(JSON 与 CSV) |
编程接口:CPUDeviceOptions提供链式构造方法withNumThreads()、withMaxThreads()、withMainWillNotDonate()、withCPUAffinity()、withSingleThreaded()等(CPUDevice.h),并支持numaPartitioned选项——当为 true 且队列类型为kThreadPool时,为每个 NUMA 节点创建分区 WorkQueue 并包装在DelegateThreadPoolWorkQueue中。
另外,worker 线程栈大小默认 8 MiB(kDefaultWorkerStackSizeBytes),可通过配置runtime.worker_stack_size_mb调整,这一默认值是为避免 macOS 继承的小栈(约 512 KB)在深层内联编译嵌套中溢出(ThreadPoolWorkQueue.cpp)。
十、总结
M::AsyncRT::WorkQueue以极简接口承载了 AsyncRT 的并发内核:四层任务队列兼顾缓存局部性与负载均衡,指数退避自旋与逐线程信号量平衡了响应延迟与功耗,mainWillDonate线程捐献机制避免了等待死锁,而设备任务固定机制则把 GPU 相关工作引导到同一 NUMA 节点的 CPU 上,从 PCIe 带宽、驱动上下文与调度可预测性三个维度优化加速器场景。配合 AsyncRT/docs/WorkQueueNonblocking.md 阐述的非阻塞设计哲学,以及MODULAR_RUNTIME_DEVICE_TASK_CPU_IDS、MODULAR_ENABLE_AFFINITY、--thread-busy-wait-time-us等调优入口,你可以针对单线程嵌入式、多线程服务器、多 GPU 推理等不同形态的负载,精确控制 Mojo 运行时的 CPU 并行行为。
【免费下载链接】mojoThe Modular Platform (includes MAX & Mojo)项目地址: https://gitcode.com/GitHub_Trending/mo/mojo
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考