PyPTO 自动CV并行流水:用 preload 编译期变换把 cube/vector 串行 kernel 改写成并行流水

发布时间:2026/9/20 13:06:55
PyPTO 自动CV并行流水:用 preload 编译期变换把 cube/vector 串行 kernel 改写成并行流水 PyPTO 自动CV并行流水用 preload 编译期变换把 cube/vector 串行 kernel 改写成并行流水【免费下载链接】pyptoPyPTO发音: pai p-t-oParallel Tensor/Tile Operation编程范式。项目地址: https://gitcode.com/cann/pypto自动CV并行流水是 pypto_pro.language.jit 提供的一项编译期变换针对 CV 融合算子中 cube 与 vector 相互依赖、串行执行时空等对方的问题自动把用户手写的串行流水 kernel 改写成并行流水版本并自动插入全部所需的核间同步。读完本文你将掌握如何通过pl.pipeline.stage划分 stage、用pl.make_tile_group声明跨核共享 Buffer配置fwd_ids/bwd_ids、在主循环中按 cube/vector 交替调用 stage以及如何用PipelineConfig(preload...)开启并逐步调优流水。串行流水与并行流水的执行节奏对比功能说明CV 融合算子中cube 和 vector 的计算相互依赖若按串行流水执行一个核工作时另一个核只能空等。自动CV并行流水的思路是让上游核提前若干次迭代计算下游核则处理上游已经算好的较早迭代的数据两个核错开节奏、同时工作从而把彼此的等待时间掩盖掉。图中sN·iM表示第 N 个 stage 正在处理第 M 次迭代的数据格子宽度代表该 stage 的耗时各 stage 耗时不同图中为示意值。自动CV并行流水是pypto_pro.language.jit的一项编译期变换包含两个能力自动流水排布把用户手写的 CV 融合算子串行流水 kernel 代码自动改写成并行流水版本提升性能自动核间同步用户开发的串行版本不需要手写任何 CV 核间同步指令并行流水版本会插入全部所需的核间同步。从源码看整个变换由 python/pypto_pro/runtime/pipeline 目录下的若干模块协作完成__init__.py对外暴露PipelineConfig、stage与transform_pipeline三个入口见 python/pypto_pro/runtime/pipeline/init.py真正的 AST 改写逻辑在 _transformer.py 的transform_pipeline中。收益上限由较慢的那个核决定需要注意的是流水化的收益上限由较慢的那个核决定图中 vector 侧单次迭代的耗时高于 cube 侧稳态节奏就由 vector 侧决定cube 侧会出现等待间隙。stage 划分越均衡两个核的耗时越接近流水填充得越满收益越高。此外流水的建立和排空各需要若干拍迭代次数越多这部分开销占比越低。这与 PipelineConfig 的实现语义一致preload越大掩盖的搬运/计算延迟越多但流水填充fill和排空drain阶段的代价也越长因此需要按 kernel 逐个调优。仓库 docs/zh/guide/figures/pro 目录下还提供了pro_pipeline_cube_bound.pngcube 侧受限的流水节奏与pro_pipeline_with_gaps.png带等待间隙的流水节奏两张示意图与本节讨论的收益边界问题直接对应可作为深入理解流水稳态节奏的参考。使用方法使用自动CV并行流水需要完成四步编写 stage 函数、声明跨核共享 Buffer、编写主循环、开启流水变换。编写 Stage 函数需要用户将算子划分为若干个计算流程每个计算流程对应一个 stage 函数通过pypto_pro.language.pipeline.stage装饰器进行标识。该装饰器在源码层面是一个透明装饰器——它只是给函数打上pipeline_stage True属性标记供流水框架识别不会改变函数行为见 python/pypto_pro/runtime/pipeline/_stage.py 中的stage与is_pipeline_stage。import pypto_pro.language as pl pl.pipeline.stage def stage1(ki, a, b_l1, a_l1_db, left_db, right_db, acc_db, mm1_vec_db): Cubemm1 A_i B cur_a a_l1_db.next() pl.load(cur_a, a, [ki * TILE, 0]) ... # 普通的 Tile/Buffer 操作不用写任何同步声明跨核共享Buffer跨核共享 Buffer 使用make_tile_group接口进行声明tile 数目由用户自主分配通过fwd_ids和bwd_ids参数配置核间正反向同步 id若未配置则不会插入对应的核间同步。其中fwd_ids为正向同步生产者写完通知消费者生产者 stage 之后插入 set、消费者 stage 之前插入 waitbwd_ids为反向同步消费者用完通知生产者可以覆写消费者 stage 之后插入 set、生产者 stage 之前插入 wait。只配置fwd_ids时仅保证消费者读到的是生产者写完的数据不保证生产者下一轮覆写时消费者已经读完。每次迭代实际使用的 id 按ids[迭代序号 % len(ids)]轮转取用。这一点在 make_tile_group 的接口文档中也有明确说明fwd_ids是标记该 group 为生产者→消费者通道的跨核事件 id取值 0~15每个 Tile 一个变换按fwd_ids[i % N]每迭代取用bwd_ids是消费者→生产者方向Buffer 释放的事件 id形状与fwd_ids相同。import pypto_pro.language as pl mm1_vec_db pl.make_tile_group( typepl.TileType(shape[TILE_HALF, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Vec), addrs0x0000, mutex_ids[12, 13], fwd_ids[0, 1], bwd_ids[2, 3], )编写主循环主循环内 stage 函数按照依赖关系顺序进行书写需要保证 stage 之间为 CV 交替。import pypto_pro.language as pl for ki in pl.range(0, N_ITER): with pl.section_cube(): stage1(ki, a, b_l1, a_l1_db, left_db, right_db, acc_db, mm1_vec_db) with pl.section_vector(): stage2(ki, sub_id, mm1_vec_db, relu_vec_db, relu_nz_db, p_mat_db) with pl.section_cube(): stage3(ki, d_l1, p_mat_db, left_db, right_db, acc_db, out_vec_db) with pl.section_vector(): stage4(ki, sub_id, out, out_vec_db)一个 kernel 内可以写多个独立的 for 循环每个循环各自做流水变换互不影响多个循环按源码顺序先后执行preload 可以逐循环配置见下一节。import pypto_pro.language as pl for ki in pl.range(0, N_ITER): # 第一条流水 with pl.section_cube(): stage1(ki, ...) with pl.section_vector(): stage2(ki, ...) pl.system.sync_all(core_typepl.SyncCoreType.MIX) for kj in pl.range(0, M_ITER): # 第二条流水 with pl.section_cube(): stage3(kj, ...) with pl.section_vector(): stage4(kj, ...)开启流水变换在pypto_pro.language.jit接口中通过PipelineConfig参数进行配置参数含义preload表示上游核提前迭代计算的次数类型为 int 或 list[int]。传入单个整数时kernel 内所有流水循环共用该值传入列表时按源码顺序为每个流水循环单独配置。取值为 0 的循环不做流水改写只在原串行循环中插入核间同步。PipelineConfig的源码实现python/pypto_pro/runtime/pipeline/config.py对参数做了严格校验preload默认值为 2接受单个 int 或多个值组成的元组用户写列表也会被归一化为 tuple因为冻结 dataclass 需要可哈希字段每个值必须是 int 且不能是 bool、不能为负数空序列会被拒绝。语义上preload是同一个核上两个相邻 stage 之间的延迟步数对应_compute_delaysctx 环形缓冲区的深度由计算出的延迟决定max_delay 1而不是直接等于 preload。建议流程先配置preload0执行确认串行版本精度正确再逐步调大 preload 开启流水。多个流水循环时可以只把其中一条配成 0如preload[0, 2]单独验证这条流水的串行精度。import pypto_pro.language as pl pl.jit(auto_mutexTrue, pipelinepl.pipeline.PipelineConfig(preload0)) def pipeline_demo_kernel(...): ...import pypto_pro.language as pl pl.jit(auto_mutexTrue, pipelinepl.pipeline.PipelineConfig(preload2)) def pipeline_demo_kernel(...): ...import pypto_pro.language as pl # kernel内有两个流水循环第一条preload1第二条preload2 pl.jit(auto_mutexTrue, pipelinepl.pipeline.PipelineConfig(preload[1, 2])) def pipeline_demo_kernel(...): ...查看生成的代码框架自动生成的并行流水代码保存在编译产物目录下文件名为pipeline_generated.py用户可以在该代码的基础上继续修改调试。使用约束自动CV并行流水对 kernel 写法有明确约束违反时会在编译期被分析器拒绝同一条流水的 stage 调用需要放在同一个 for 循环内。一个 kernel 可以有多个独立的流水循环循环之间需要用户手动插入全核同步两条流水循环不能嵌在同一个外层循环里。流水循环的步长必须为正不支持倒序迭代。preload 取值必须大于等于 0传入列表时列表长度必须与 kernel 内流水循环的个数相同。不支持 stage 嵌套 stage。跨核 Buffer 和核内 Buffer 的make_tile_group声明必须写在 kernel 函数体内不支持在被 stage 调用的普通函数里声明。stage 函数不支持有返回值。stage 调用的普通函数不能返回 tile 或 tile group。每个with pypto_pro.language.section_cube()/pypto_pro.language.section_vector()块里只放单个 stage 调用且 stage 调用需要严格按 cube/vector 交替排列C→V→C→V…不允许连续两个 stage 落在同一个核上。流水循环体内第一个 stage 与最后一个 stage 之间不允许插入其他语句。流水循环含 stage 调用的那个 for 循环体内不允许获取或操作跨核 Buffer以及与跨核 Buffer 地址复用的核内 Buffer请把它们放进 stage 函数体内。在流水循环之外取 Tile 时必须单独写成一条赋值语句slot group.next()不能嵌在更大的表达式里同一个变量只能取一次 Tile需要多个 Tile 请用多个变量。同一个 stage 函数在一条流水循环内只能调用一次不允许重名 stage。允许通过 if 语句判断 stage 执行场景但分支条件必须为编译期常量。从源码看_transformer.py 中的_prune_const_branches会在分析前把编译期常量分支折叠掉使下游分析只看到无条件 stage若运行时条件包裹了 stage 调用则直接报错动态分支 stage 不受支持。fwd_ids/bwd_ids取值范围为 0~15。fwd_ids/bwd_ids只支持两种写法直接写整数列表fwd_ids[0, 1]或写一个在 kernel 外绑定到整数列表的变量名IDS [0, 1]…fwd_idsIDS。元素必须是编译期常量整数不支持切片、拼接、函数调用等表达式形式。fwd_ids/bwd_ids的长度只能等于该 Buffer 的 Tile 数或者等于 1。等于 1 时多个 Tile 共用同一个同步 id交接会被串行化性能下降但结果正确用于同步 id 不够分配的场景。允许跨核 Buffer 之间、跨核 Buffer 与核内 Buffer 之间进行地址复用但最多允许两块 Buffer 复用且复用双方的 Tile 数需要一致。地址复用的 Buffer 之间必须声明相同的mutex_ids。mutex 锁的是地址不同的 id 等于没有互斥。一个跨核 Buffer通过fwd_ids/bwd_ids标记需要恰好被两个 stage 使用且这两个 stage 分别在 cube 和 vector 上构成一对一的生产者/消费者关系。跨核 Buffer 的 Tiles 必须随迭代顺序轮转。stage 函数如果有结构体参数请使用pypto_pro.language.struct()声明不支持pypto_pro.language.struct_array()。跨核 Buffer 仅支持 UB 和 L1 Buffer。以_pl_开头的变量名为框架保留kernel 内不要使用。调用示例以下完整示例来自文档单个 1:2 CV 执行组relu(A B) D两段 matmul 的 CV 融合演示了上述全部要素的组合4 个 stage 依次完成mm1 A_i Bcube→relu并 insert 回 L1vector→mm2 relu(mm1) Dcube→ 写回 GMvector。import os import pypto_pro.language as pl import pytest import torch import torch_npu import pypto ST_DEVICE_ID int(os.environ.get(TILE_FWK_DEVICE_ID, 0)) ST_DEVICE fnpu:{ST_DEVICE_ID} TILE 64 TILE_HALF TILE // 2 # dual mode 沿 M 劈开每个 AIV 处理一半 N_ITER 4 pl.pipeline.stage def stage1(ki, a, b_l1, a_l1_db, left_db, right_db, acc_db, mm1_vec_db): Cubemm1 A_i B cur_a a_l1_db.next() pl.load(cur_a, a, [ki * TILE, 0]) b_slot b_l1.current() left left_db.next() right right_db.next() acc acc_db.next() pl.move(left, cur_a) pl.move(right, b_slot) pl.matmul(acc, left, right) mm1_vec mm1_vec_db.next() # DualModeSplitM[TILE, TILE] 的 acc 沿 M 劈开每个 AIV 拿 [TILE_HALF, TILE] pl.move(mm1_vec, acc, acc_to_vec_modepl.AccToVecMode.DualModeSplitM) pl.pipeline.stage def stage2(ki, sub_id, mm1_vec_db, relu_vec_db, relu_nz_db, p_mat_db): Vector对本 AIV 那半做 relu转 NZ 后按行偏移 insert 回 L1 Buffer mm1_vec mm1_vec_db.next() relu_vec relu_vec_db.next() pl.relu(relu_vec, mm1_vec) relu_nz relu_nz_db.next() pl.move(relu_nz, relu_vec) # ND - NZinsert要求源Tile为NZ p_mat p_mat_db.next() pl.insert(p_mat, relu_nz, [sub_id * TILE_HALF, 0]) pl.pipeline.stage def stage3(ki, d_l1, p_mat_db, left_db, right_db, acc_db, out_vec_db): Cubemm2 relu(mm1) D d_slot d_l1.current() p_mat p_mat_db.next() left left_db.next() right right_db.next() acc acc_db.next() pl.move(left, p_mat) pl.move(right, d_slot) pl.matmul(acc, left, right) out_vec out_vec_db.next() pl.move(out_vec, acc, acc_to_vec_modepl.AccToVecMode.DualModeSplitM) pl.pipeline.stage def stage4(ki, sub_id, out, out_vec_db): Vector每个 AIV 写回本迭代结果的一半 out_vec out_vec_db.next() pl.store(out, out_vec, [ki * TILE sub_id * TILE_HALF, 0]) pl.jit(auto_mutexTrue, pipelinepl.pipeline.PipelineConfig(preload2)) def pipeline_demo_kernel( a: pl.Tensor[[N_ITER * TILE, TILE], pl.DT_FP32], b: pl.Tensor[[TILE, TILE], pl.DT_FP32], d: pl.Tensor[[TILE, TILE], pl.DT_FP32], out: pl.Tensor[[N_ITER * TILE, TILE], pl.DT_FP32], ): # 跨核共享 Buffer声明 fwd_ids/bwd_ids框架据此自动插同步 mm1_vec_db pl.make_tile_group( typepl.TileType(shape[TILE_HALF, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Vec), addrs0x0000, mutex_ids[12, 13], fwd_ids[0, 1], bwd_ids[2, 3], ) p_mat_db pl.make_tile_group( typepl.TileType(shape[TILE, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Mat, layoutpl.NZ), addrs0x8000, mutex_ids[14, 15], fwd_ids[4, 5], bwd_ids[6, 7], ) out_vec_db pl.make_tile_group( typepl.TileType(shape[TILE_HALF, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Vec), addrs0x10000, mutex_ids[16, 17], fwd_ids[8, 9], bwd_ids[10, 11], ) # Cube 侧局部 Buffer with pl.section_cube(): a_l1_db pl.make_tile_group( typepl.TileType(shape[TILE, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Mat, layoutpl.NZ), addrs0x0000, mutex_ids[0, 1], ) b_l1 pl.make_tile_group( typepl.TileType(shape[TILE, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Mat, layoutpl.NZ), addrs0x10000, mutex_ids[2], ) d_l1 pl.make_tile_group( typepl.TileType(shape[TILE, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Mat, layoutpl.NZ), addrs0x14000, mutex_ids[3], ) left_db pl.make_tile_group( typepl.TileType(shape[TILE, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Left, layoutpl.NZ), addrs0x0000, mutex_ids[4, 5], ) right_db pl.make_tile_group( typepl.TileType(shape[TILE, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Right, layoutpl.ZN), addrs0x0000, mutex_ids[6, 7], ) acc_db pl.make_tile_group( typepl.TileType( shape[TILE, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Acc, layoutpl.NZ, fractal1024, ), addrs0x0000, mutex_ids[8, 9, 10, 11], ) # B、D 在循环外一次搬入全程复用 b_slot b_l1.current() d_slot d_l1.current() pl.load(b_slot, b, [0, 0]) pl.load(d_slot, d, [0, 0]) # Vector 侧局部 Buffer with pl.section_vector(): sub_id pl.get_subblock_idx() relu_vec_db pl.make_tile_group( typepl.TileType(shape[TILE_HALF, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Vec), addrs0x8000, mutex_ids[18, 19], ) relu_nz_db pl.make_tile_group( typepl.TileType( shape[TILE_HALF, TILE], dtypepl.DT_FP32, target_memorypl.MemorySpace.Vec, layoutpl.NZ, ), addrs0x18000, mutex_ids[20, 21], ) # 流水循环按 cube/vector 交替调用 4 个 stage for ki in pl.range(0, N_ITER): with pl.section_cube(): stage1(ki, a, b_l1, a_l1_db, left_db, right_db, acc_db, mm1_vec_db) with pl.section_vector(): stage2(ki, sub_id, mm1_vec_db, relu_vec_db, relu_nz_db, p_mat_db) with pl.section_cube(): stage3(ki, d_l1, p_mat_db, left_db, right_db, acc_db, out_vec_db) with pl.section_vector(): stage4(ki, sub_id, out, out_vec_db) pytest.mark.soc(950) pypto.options(pass_options{enable_slice: False}) def test_pipeline_demo_kernel(): device ST_DEVICE torch.npu.set_device(device) torch.manual_seed(0) a torch.randn(N_ITER * TILE, TILE, devicedevice, dtypetorch.float32) b torch.randn(TILE, TILE, devicedevice, dtypetorch.float32) d torch.randn(TILE, TILE, devicedevice, dtypetorch.float32) out torch.zeros(N_ITER * TILE, TILE, devicedevice, dtypetorch.float32) pipeline_demo_kernelNone, 1 torch.npu.synchronize() ref torch.relu(a b) d torch.testing.assert_close(out, ref, rtol1e-2, atol1e-2)[!NOTE]说明本示例按单个 1:2 CV 执行组设计因此block_dim设置为 1。扩展为多个执行组时需要通过pypto_pro.language.get_block_idx()划分各执行组处理的 GM 数据和输出范围。block_dim的完整说明参见 Kernel核函数。深入原理编译期变换与同步依赖图理解了使用方式之后再从源码层面看自动CV并行流水的实现机制可以帮助你更好地理解约束的成因与 preload 的调优方向。变换入口从 AST 到并行流水transform_pipelinepython/pypto_pro/runtime/pipeline/_transformer.py接收 kernel 函数的 AST 与配置工作流程大致为在 AST 深拷贝上工作先折叠编译期常量分支_prune_const_branches保证下游分析只看到无条件的 stage 调用analyze_pipeline分析串行结构按源码顺序返回每个流水循环的PipelineInfo对每个流水循环若 preload 非 0则调用_compute_delays计算各 stage 的延迟beat 偏移同一核上第一个 stage 取上游跨核 stage 的延迟 1同一核上的后续 stage 取前一个同核 stage 的延迟 preload依据延迟计算 ctx 环形缓冲区深度max_delay 1把串行循环替换为带 ctx 槽位轮转、延迟 stage 与合法性守卫is_valid的并行循环循环结束后追加 drain beats让延迟的 stage 追完剩余任务preload0 时走_transform_serial分支保持串行循环结构只自动插入跨核同步。多个流水循环共享同一个函数作用域因此框架引入的所有变量_pl_ctx_arr、_pl_task_id、_pl_ctx_0/_pl_ctx_negN、drain 循环变量等都会按流水序号加后缀并通过validate_names与用户变量做碰撞检查——这也是以_pl_开头的变量名为框架保留约束的由来。同步图模型一次建模覆盖串行与并行所有核间同步的插入都来自 _sync_graph.py 中构建的同步依赖图它用一个模型统一覆盖 preload0、preloadN 与地址复用三种场景nodeop 级访问Access(stage, buffer, role, pipe)region物理内存区域对地址重叠的 Buffer 做并查集合并共址 Buffer 落入同一 regionlane一个 (region, 物理槽位)同一内存上所有竞争访问按时间排序边只会在 lane 内部形成因此地址复用不需要特殊处理edgeRAW读等最近的前驱写与 WAR写等前驱写在该槽位上的所有读者后者天然产生地址复用的反向保护边。每个 stage 恰好推进跨核 Buffer 一次slot task % slot_count距离来自时间线的展开而非槽位算术。规划输出分三部分prefire循环前预释放偏斜边头部缺失的 permit、sites循环内每个 stage 前后的 wait/set、drain循环后吸收尾部多余的 set三者必须精确配对因为硬件要求 set 与 wait 计数严格相等。值得注意的两个实现细节用户声明了fwd_ids却没有bwd_ids等价于声明该 Buffer 槽位足够多不需要覆写保护该 WAR 边会被解析为不发射指令——但边仍然存在于图中环与逆时间分析始终能看到真实依赖地址复用产生的跨 Buffer WAR 依赖是用户从未声明过 id 的框架会从空闲事件 id 池0~15中自动分配一组_pl_overlap_ids_N会被声明在循环之前若池不够会按槽位数量升序把部分边降级为单 id 共用交接串行化但结果正确仍不够则报错。环形缓冲区与任务计数并行版本的核心数据结构是 ctx 环形缓冲区pl.struct_array声明每个 beat 通过_pl_ctx_arr[_pl_task_id % depth]选取槽位快照本 beat 各 stage 的参数stage 参数被推导为 ctx 字段延迟 stage 则读取自己延迟拍之前填充的槽位配合is_valid守卫区分有效任务与排空拍。任务计数_pl_task_id按帧递增drain beats 在循环结束后按depth - 1次继续推进让最后几拍的任务在各自延迟的 stage 上跑完。调优建议与测试验证仓库 python/tests/st/pypto_pro/frontend/auto_pipeline 目录下提供了多个针对自动流水变换的测试用例例如test_auto_pipeline_two_loops.py验证单 kernel 内多条流水循环与test_pipeline_loop_body.py验证循环体语句处理可作为理解变换行为与回归验证的参考。结合前文给出如下调优与排查路径先验证串行正确性PipelineConfig(preload0)只插入核间同步不改写循环是精度验证的基线多流水循环时可逐条置 0如preload[0, 2]定位问题出自哪条流水再逐步增大 preload从 1 开始逐级提升观察pipeline_generated.py中 ctx 深度、延迟与同步站点的变化收益上限受较慢核的稳态节奏约束stage 划分越均衡收益越高遇到同步相关报错时优先核对跨核 Buffer 的fwd_ids/bwd_ids是否按约束声明范围 0~15、长度等于 Tile 数或 1、写法合规、地址复用的 Buffer 是否声明了相同的mutex_ids、跨核 Buffer 是否恰好被一个 cube stage 与一个 vector stage 一对一使用、Tiles 是否随迭代轮转。至此你已经掌握自动CV并行流水的完整用法用pl.pipeline.stage划分阶段、用make_tile_group的fwd_ids/bwd_ids声明跨核通道、在jit中配置PipelineConfig(preload...)开启变换并能在生成的pipeline_generated.py基础上继续调试与优化。【免费下载链接】pyptoPyPTO发音: pai p-t-oParallel Tensor/Tile Operation编程范式。项目地址: https://gitcode.com/cann/pypto创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考