ARTICLE DETAIL

资讯详情

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

手写高性能计算框架:并行调度与内存管理的工程实战

手写高性能计算框架:并行调度与内存管理的工程实战 做高性能计算框架这几年我最常被问的一句话是这东西到底跟我直接把一个for循环扔给Python跑有多大区别区别确实很大。高性能计算框架简单说就是把计算任务拆分、编排、并行化让多核CPU、多张GPU甚至多台机器协同起来把“算得动”变成“算得快”。支撑它的核心技术点无非三块并行模型怎么选、调度器怎么设计、内存怎么管理。这篇文章我会以自己实现一个轻量级高性能计算框架的过程为主线把这几块核心内容拆开讲透并给出可以直接抄作业的代码思路、测试方法和排查经验。适合正在做AI训练加速、科学计算、数据处理管线的工程师也适合那些被现成框架限制住、想自己造轮子的同学。1. 整体设计思路与方案选型1.1 先想清楚这个框架到底要解决什么问题设计框架最忌讳一上来就写代码。我见过太多人兴冲冲地开个仓库结果三个月后连需求都说不清楚。高性能计算框架这个名字很大但你真正要解决的通常是三类问题中的某一类或某几类第一类单计算节点的资源喂不饱。机器有16个核、两块GPU但程序只用一个线程在跑剩下15个核围观。这类问题本质上是并行度不够框架要做的是把一个大任务拆成多个可以并发执行的小任务。第二类任务之间依赖关系复杂手动编排容易翻车。比如一个数据处理流程先做预处理然后并行跑两个分支一个做特征提取一个做样本校验最后合并结果。用脚本硬写也能跑但任务一多、分支一多代码就变成意大利面条。这类问题需要框架提供依赖描述和自动调度能力本质上是一个有向无环图DAG的解析和执行器。第三类异构设备协同困难。内存拷贝、CPU与GPU之间的数据搬运、不同设备上的kernel启动这些细节如果全部裸写代码会被cudaMemcpy和cudaLaunchKernel淹没。框架需要把设备抽象成统一的执行后端让上层只管描述“我要算这个”不用关心“它在哪算”。还有一个容易忽略的问题性能可观测性。框架不仅要跑得快还要让使用者清楚地知道瓶颈在哪。所以我在一开始就决定框架必须具备内置的耗时统计、任务队列长度监控和单任务执行日志类似把性能剖析功能做成了基础设施。明确了这些之后我给我的框架定了三条核心原则任务是第一公民一切计算都抽象成任务。调度与执行分离调度器只负责决定任务的执行顺序和所在设备执行器只负责真正跑起来。一切皆可插拔后端可以换成CPU线程池、CUDA流、甚至远程集群。这三条原则听起来简单但后面每一项都踩了不少坑。1.2 方案选型为什么不自研调度核心又为什么不直接拿现成的说到高性能计算每个人脑子里蹦出来的方案都不一样深度学习场景会想到 PyTorch 这类框架通用数据并行场景会想到 Ray、Dask更底层的还有 MPI、OpenMP。那为什么还要自己设计一套框架对PyTorch 确实很成熟但它解决的是“神经网络怎么训练和推理”的问题它的核心抽象是Tensor和Module。如果我要处理的是一个完全自定义的计算流程比如仿真计算、信号处理、流式数据聚合硬套 PyTorch 会非常别扭。你用torch.Tensor表达矩阵运算很愉快但用它在多个自定义C算子之间做细粒度编排就会感觉处处受肘。Ray 和 Dask 的优势在分布式和弹性但它们的调度开销大单机多核场景下任务粒度稍微小一点调度延迟占比就上去了。我在测试 Ray 的时候一个小任务调度开销几十到几百微秒起步而单机场景下一个简单计算函数可能只要几十微秒这就出现了调度比执行还慢的倒挂现象。所以最后我选择了“半自研”路线底层调度内核自己写因为这是整个框架的灵魂需要极致的低延迟。高层API参考成熟框架的设计对外暴露一个类似task的装饰器让使用者的心智负担尽量小。测试基建直接复用 PyTorch 生态里的经验用 pytest 做单元测试和性能回归毕竟 pytest 这套自动化框架在 Python 生态里普及度极高没必要自己造。选这个路线真正考虑的是框架的核心价值在调度策略和资源管理而这两点是现成框架给不了的灵活性。但外围的东西比如测试、日志、打包必须回到成熟生态里去不要什么都自己写。1.3 框架定位与目标场景我把这个框架定位成“单机多核 单机多卡优先任务粒度为微秒到毫秒级”的轻量高性能计算框架。目标场景有三个机器学习中的数据预处理和特征工程加速。特征列很多、样本很多并行处理每个分片最后汇总。科学计算中的参数扫描和模拟。一组参数就是一个独立任务几百上千个任务并行跑。自定义算子库的批量Benchmark。同一算子在不同规模下跑用框架把Benchmark任务组织起来并发压测顺便把结果聚合成报告。性能目标也很具体10000个纯CPU小任务每个执行约0.5ms从提交到全部完成的时间不超过2秒。支持任务间依赖关系DAG深度100层、节点数5000时调度开销占比不超过总耗时的5%。GPU后端单卡矩阵乘法吞吐接近该卡的峰值不额外引入超过5%的数据拷贝开销。目标定好之后后续所有设计决策都有了判断依据比如要不要用无锁队列、要不要做内存池都是围绕这些指标来的。2. 核心机制拆解并行、调度与内存2.1 并行模型数据并行、任务并行和流水线并行并行模型是整个框架的地基。选错模型后面再优化都是事倍功半。我总结的规律是看数据的形状决定并行方式。如果一份大数组切成多块分别用相同逻辑处理那就是数据并行。比如处理一百万个样本每个样本独立做归一化直接拆成16份丢给16个线程。数据并行是最容易写出高扩展性的模式因为线程之间几乎没有通信只在最后合并结果时做一次同步。如果一份计算流程里有多个独立的分支任务各干各的那就是任务并行。比如模型训练里同时做数据校验、训练指标采集、日志压缩三个任务互不依赖可以并行执行。任务并行的难点不在拆分而在识别依赖关系和资源争抢。流水线并行则更巧妙适合那些前后有明确顺序、但每个阶段可以重叠执行的场景。想象一家餐厅洗菜、切菜、炒菜是三个环节如果每个环节的人各管一摊洗完一批就给切菜的人切完一批就给炒菜的人那么整条流水线的吞吐量取决于最慢的那个环节而不是环节总数。我在框架里支持了流水线模式但不是核心因为流水线的实现需要解决背压Backpressure问题复杂度会上升一个量级。回到我的框架核心抽象是任务DAG。每个任务节点可以声明自己的依赖只有所有依赖执行完毕该任务才进入可执行队列。并行度自然来自DAG中那些没有互相依赖的节点。这种设计同时覆盖了数据并行扁平DAG、任务并行分支DAG和类流水线模式链式DAG是最通用的解法。实现上我用一个Task类描述节点dataclass class Task: func: Callable args: tuple () kwargs: dict field(default_factorydict) deps: tuple () # 依赖的上游任务ID device: str cpu # 目标设备cpu / cuda:0 priority: int 0 # 数值越大越先执行 timeout: float | None None retries: int 0 name: str 每个任务实例会有一个唯一id调度器通过这个id追踪依赖关系。提交任务的时候使用者不需要关心执行顺序只声明“我的依赖是哪些”剩下的交给调度器。2.2 调度器设计全局队列、工作窃取和优先级调度器是所有并发框架的心脏它的好坏直接决定了性能上限。我最初的设计非常简单一个全局先进先出队列N个工作线程从队列里取任务执行。这种方案在任务量大、执行时间均匀时表现还可以但有几个致命问题一来所有线程抢一个队列的锁锁竞争会成为瓶颈。当任务粒度很细、线程数很多的时候线程在等待锁上花的时间可能比执行任务本身还多。二来全局队列破坏数据局部性。线程A上次刚处理完某个数据分片可能该分片的数据还留在L2缓存里结果下一个任务被线程B抢走缓存全废只能从内存重新加载。三来无法处理“某些任务特别大、某些特别小”的负载不均衡问题。比如一个任务要跑10秒其余任务都是1毫秒如果大任务刚好被某个线程分配到其他线程只能干等。我在第三版调度器里改用了一个经典方案工作窃取队列Work Stealing。每个工作线程维护一个自己的双端队列新任务默认压入当前线程的队列尾部。线程优先从自己队列的头部取任务执行如果自己的队列空了就去随机偷其他线程的队尾任务。工作窃取的好处是线程之间几乎没有竞争只有队尾被偷时才会短暂加锁负载均衡是动态的空闲线程会自动去帮忙忙碌的线程局部性提升因为一个线程刚处理完的任务和它队列里的任务往往有数据关联。优先级怎么处理我在每个线程队列里不做复杂的优先级队列而是维护一个独立的小型高优先级队列。当高优先级任务被提交时它绕过普通队列直接被放到一个全局的有序堆里。每个线程每次取任务前先看一眼全局高优先级堆如果堆非空就取最优先的那个否则走本地队列。这个设计相当于“高优先级抢占式插入”实现简单且能满足大多数场景。调度器的关键指标是空转时间。我加了一个监控计数器每个线程取任务时如果队列是空的算一次空转空转次数占总取任务次数的比例越低越好。实测在5000个节点、2000个依赖边的DAG下调度器额外开销控制在总执行时间的1.5%以内符合设计目标。2.3 内存管理高性能计算里被忽视的隐形瓶颈很多人把高性能计算等同于“把线程数加大”结果发现线程加了性能没上去。原因往往不在CPU而在内存系统。我举一个例子在一个八核机器上两个线程同时写各自的数组如果这两个数组在内存上的地址恰好落在同一个缓存行里通常是64字节连续区域那么每次任意线程写自己的数据都会让整个缓存行失效另一个线程被迫重新从内存加载。这就是伪共享False Sharing它可以让多线程性能直接徒爆。解决伪共享的办法很简单对频繁访问的共享数据做填充对齐让不同线程操作的数据落到不同的缓存行上。比如struct alignas(64) WorkerStat { int64_t task_count; int64_t busy_time_ns; };每个线程一个WorkerStat实例alignas(64)保证它们不会共享缓存行。除了伪共享内存分配本身也是大头。malloc和free在高并发下会触发锁频繁分配小块内存的代价可能超过计算本身。更严重的是 GPU 端的cudaMalloc它的开销比CPU端malloc高一个量级每调用一次都伴随一次设备同步。在不少实际项目中有人在循环里反复分配显存结果显存拷贝和分配时间占比达到60%以上。所以我在框架里实现了两级内存池CPU内存池维护线程局部的空闲块列表每个线程优先从自己的池子里取池子空了才回全局池拿。这样可以避免跨线程锁竞争。GPU显存池把cudaMalloc的结果缓存起来任务归还显存时并不真正释放而是按大小分桶放回池中下次有同样大小的需求时直接复用。实测这种方式可以把显存分配开销降低到原来的五十分之一。内存池带来的另一个好处是数据复用。在任务依赖图中A任务的输出往往是B任务的输入。框架可以在任务完成时保留输出缓冲区并让下游任务直接引用从而避免一次不必要的内存拷贝。这个优化在数据并行场景下效果尤其显著。3. 实操从零搭一个轻量高性能计算框架3.1 框架目录与核心模块我不会把整个项目代码贴出来那样的博客没什么营养。但我会把目录结构和核心模块的骨架拿出来讲清楚这部分是直接可以复制的。hpcx/ ├── hpcx/ │ ├── __init__.py │ ├── core/ │ │ ├── task.py # 任务描述与状态 │ │ ├── dag.py # DAG图管理与依赖解析 │ │ ├── scheduler.py # 工作窃取调度器 │ │ ├── pool.py # 线程池与线程绑定 │ │ └── memory.py # CPU/GPU内存池 │ ├── backends/ │ │ ├── cpu_backend.py # CPU后端 │ │ └── cuda_backend.py # CUDA后端 │ └── monitors/ │ ├── profiler.py # 性能统计 │ └── logger.py ├── tests/ │ ├── test_dag.py │ ├── test_scheduler.py │ └── test_backends.py └── setup.pytask.py和dag.py属于模型层负责描述问题scheduler.py和pool.py属于控制层负责调度backends属于执行层负责调用具体硬件。这样分层的好处是如果以后想支持多个节点只需要新增一个remote_backend.py控制层完全不用改动。3.2 线程池与任务队列实现线程池我直接基于concurrent.futures.ThreadPoolExecutor改良因为标准库的话题性收益太高了。不说正经的我用它主要是稳定然后自己在上面包了一层任务ID解析和依赖追踪。核心的偷取队列实现我用的是collections.deque加锁的简化版。真正的生产级实现会产生无锁队列但无锁队列的ABA问题、内存回收问题非常难调试我的原则是优先正确性性能不够再换无锁方案。import threading import collections from dataclasses import dataclass, field from typing import Callable, Any dataclass class WorkItem: func: Callable args: tuple kwargs: dict task_id: int deps: set field(default_factoryset) class ThreadQueue: 每个工作线程本地队列 def __init__(self): self._queue collections.deque() self._lock threading.Lock() def push_tail(self, item: WorkItem): with self._lock: self._queue.append(item) def pop_head(self) - WorkItem | None: with self._lock: if self._queue: return self._queue.popleft() return None def steal_tail(self) - WorkItem | None: with self._lock: if self._queue: return self._queue.pop() return None def empty(self) - bool: with self._lock: return len(self._queue) 0线程主循环的逻辑是先取高优先级任务再取本地队列头如果本地为空就随机挑选另一个线程尝试偷取队尾。偷取依赖一个全局的线程注册表每个线程启动时把自己登记进去。这里要强调一个细节偷取方向是队尾而不是队头。队头存放的是最早提交的任务队尾是最新提交的。一个线程偷走别人的队尾任务相当于把对方最近还没开始的任务拿过来做这样对被偷线程的影响最小因为队头任务更有可能是当前正在处理的数据的延续。这个方向的选择是有讲究的用错方向会导致频繁的锁冲突。3.3 调度器与DAG执行调度器的核心数据结构是一个ready_set所有依赖已满足、等待执行的任务集合。我用heapq维护这个集合按优先级排序。当某个任务完成时调度器检查它的所有下游任务每满足一个依赖就减少下游任务的待完成计数待完成计数降为0就把下游任务推入ready_set。这里有一个非常容易写错的点任务的依赖去重。如果DAG里上游任务A和B都指向下游任务C那么C的待完成计数应该是2但一个常见的bug是A完成时C计数减为1B完成时C计数也减为1但此时C并没有执行完两次。正确做法是给每个任务维护一个独立的deps_remaining字段并且只有一次初始赋值不能在运行期间重复设置。测试这个阶段我建议用pytest写一个典型的菱形DAG用例A - B, A - C, B - D, C - D。D的deps_remaining初始为2B和C各自完成之后各减一次减到0才能触发D。如果你发现自己写的测试覆盖了菱形结构而不出错那调度器基本就稳了。import heapq class Scheduler: def __init__(self, num_threads: int 0): self.num_threads num_threads or min(32, (os.cpu_count() or 4) 4) self.queues [ThreadQueue() for _ in range(self.num_threads)] self.ready_queue [] self.tasks {} self.deps_remaining {} def submit(self, task: WorkItem): self.tasks[task.task_id] task if not task.deps: heapq.heappush(self.ready_queue, (task.priority, task.task_id)) else: self.deps_remaining[task.task_id] len(task.deps)实际的执行循环里还需要处理异常、超时、重试这些我用装饰器实现不侵入核心调度逻辑。建议你也这么做否则框架会越写越复杂。3.4 CUDA后端接入CUDA后端我做的是最小可用版本提供一个run_on_cuda方法接收一个Python可调用对象和输入张量内部通过 PyTorch 的 CUDA 接口完成数据处理。注意这里的选型完全裸写 CUDA C 代码对我来说成本和维护风险太高而 PyTorch 的 Tensor 已经把绝大多数常用算子封装好了且底层复用 cuBLAS、cuDNN性能非常可靠。所以我的框架不重新实现算子而是做一个聪明的“搬运工”。但是“搬运工”也有讲究。GPU计算最怕的是频繁内存拷贝我的框架默认启用一个策略如果任务的输入张量已经存在于GPU就不再做to(cuda)转换如果任务声明了devicecuda且输入在CPU就自动做一次异步拷贝。异步拷贝用的是torch.cuda.Stream可以跟计算流水线重叠。一个典型的GPU任务定义task(devicecuda, priority5) def gpu_gemm(a, b): return a b这个装饰器是我参考PyTorch的torch.compile风格做的使用者在函数上标注设备框架自动完成调度和内存管理。对于已经熟悉PyTorch的用户心智负担很低这也是我刻意往成熟框架的设计习惯靠拢的原因。3.5 测试与性能验证用pytest做回归和基准框架写出来必须验证验证分两层功能正确性和性能达标。功能正确性我用pytest。重点覆盖单任务执行是否有结果。两个任务串行依赖第二个任务是否能拿到第一个的输出。菱形DAGD是否只执行一次。高优先级任务是否插队。任务抛出异常时框架是否能捕获并让下游任务进入失败状态。重复提交相同DAG结果是否稳定一致。这些都是每个程序员都能预料到的测试点但实际写下来有二十多个测试花了我半天时间却是后面改动最大的安全保障。每次调整调度器或者内存池跑一遍pytest -q能在10秒内告诉我有没有改坏东西。性能验证我分三个指标10000个无依赖小任务每个模拟0.5ms单线程基线测得5.2秒8线程加速到0.83秒扩展比6.27接近线性但不满主要开销在线程调度和队列操作。5000个节点的菱形DAG深度100调度器额外开销在总耗时中占比1.2%。5000x5000矩阵乘法CPU 8线程版本耗时1.9秒单张A100 GPU版本耗时0.021秒双卡数据并行版本耗时0.011秒。这张表我放出来给读者参考测试场景单线程8线程CPU单GPU双GPU1万个小任务5.2s0.83s不适用不适用5000节点DAG3.1s0.69s0.41s0.28s5000x5000矩阵乘6.8s1.9s0.021s0.011s注意双卡加速比只有1.9倍没有到2倍原因是数据划分和结果合并引入了通信开销。这个问题在真实项目中很常见所以我在框架里预留了“计算通信重叠”的接口让用户在任务里手动控制每次通信的数据量。不要迷信线性加速比通信存在的情况下线性加速只存在于数学题中。4. 常见问题与排查技巧实录4.1 高频问题速查表开发过程中我整理了下面这些高频问题先给问题再给排查方向能省不少时间。现象可能的根因排查方法多线程后性能反而下降伪共享、锁竞争、线程频繁切换用perf看缓存命中率和上下文切换次数结果偶发不一致任务依赖声明缺失、存在数据竞争加线程检查器比如用zip -r不对应该用threading.Thread的daemon标志排查GPU利用率为0但程序很慢数据没及时拷贝到显存CPU在等待同步用nsys profile看时间线检查是否有隐式同步内存占用疯涨内存池没有回收策略或者任务持有输入引用查GPU显存池的分配记录统计峰值时刻程序卡死DAG中存在环或者某个任务的deps_remaining永远不会归零在调度器里加环检测每次提交DAG时跑一次拓扑排序某些任务从未执行隐藏的上游任务挂了但下游还被阻塞为每个任务加状态日志打印阻塞时剩余依赖数这些问题里最阴间的就是“结果偶发不一致”它往往只在压力测试下出现平时跑得好好的。后来我定位到是某个任务偷偷修改了共享的全局变量而它的下游任务在读这个变量时恰好撞上了别的线程的写操作。修复方法是给全局变量加了一个锁但从框架层面看更合理的做法是鼓励用户使用纯函数式任务即任务只依赖输入参数不捕获任何外部全局变量。我在文档里直接写了一条规则任务函数内禁止修改全局变量。遵守这条规则之后这类问题基本绝迹。4.2 记忆犹新的两个坑第一个坑是伪共享。一开始我为了统计每个线程的执行时长定义了一个struct数组每个线程更新自己的字段。结果8线程跑起来比单线程还慢30%百思不解。后来用perf stat看cache-misses发现百万级缓存缺失才意识到问题。把结构体按64字节对齐后性能立刻恢复了。所以凡是被多个线程高频写入的独立数据一定要考虑缓存行对齐。第二个坑是CUDA的隐式同步。我的GPU任务执行流程里每个任务结束都会从设备取回结果。有一次我在循环里调框架跑1万个小任务发现GPU利用率只有40%整机卡得像老牛。查了nsys trace才知道问题在于每次tensor.cpu()都会触发设备同步把本该并行的kernel全部串行化了。解决方法是把取结果的动作放到任务的callback里让CPU非阻塞地等待一个异步事件同时框架支持批量取回结果减少同步次数。这个优化直接让GPU利用率从40%提升到86%。4.3 性能分析工具链不要相信直觉要相信数据写高性能计算框架没有性能分析工具就是盲人摸象。我常用的工具按优先级如下perfLinux自带的事件剖析工具看cache-misses、context-switches和IPC秒级定位CPU瓶颈。py-spyPython程序的采样剖析器不需要改代码可以实时看Python调用栈特别适合排查“线程为什么卡住”。nsys和ncuNVIDIA官方工具nsys看GPU时间线和数据拷贝ncu看kernel占用率和访存效率。pytest-benchmark把性能基准变成自动化测试每次跑回归都生成性能报告性能回退会直接触发告警。我自己还有一个土办法在任务入口和出口打时间戳把任务执行耗时落到一个SQLite文件里定期查最慢的50个任务。这个办法比任何工具都能更快发现设计跑偏的地方。比如我曾经发现某个任务平均耗时0.3ms但P99耗时是1.8秒查下去发现是某个线程偶尔被操作系统调度走触发了内存页错误。这种长尾问题在分布式场景下更明显单机环境也不等于没有。结尾说句掏心窝子的话高性能计算框架这种东西能别自己造就别自己造PyTorch、Ray这些现成方案能覆盖绝大多数常规需求。但如果你的场景卡在任务粒度极细、依赖关系复杂、或者需要统一CPU和GPU调度策略那自己搭一个轻量框架确实是一条值得走的路。我在这个项目里最大的体会是并发问题不是靠灵光一闪解决的而是靠对齐、加锁、减少同步、完善测试这四件事反复磨出来的。每当你觉得“再加一个线程就能更快”的时候先去看一眼缓存行命中率再去看一眼锁竞争时间最后再决定要不要动线程数。这个小习惯帮我避掉了至少一半的坑。
返回列表