pyasc 框架自动插入流水同步:基于 TPipe/TQue 的 Add 算子实现详解
【免费下载链接】pyasc本项目为Python用户提供算子编程接口,支持在昇腾AI处理器上加速计算,接口与Ascend C一一对应并遵守Python原生语法。项目地址: https://gitcode.com/cann/pyasc
本指南以 CANN pyasc 开源仓库中的examples/02_add_framework样例为核心,系统讲解如何利用 Ascend C 框架(TPipe/TQue)自动完成流水同步,实现两个向量 x、y 的逐元素加法 z = x + y。通过本指南,读者将掌握 pyasc 框架编程模式下的 copy_in / compute / copy_out 三段式算子结构、TQue 队列 enque/deque 的自动同步原理,以及与手动set_flag/wait_flag同步方式的差异。
样例概述
examples/02_add_framework样例演示了通过 Ascend C 框架自动插入流水同步的 Add 算子:两个向量 x 和 y 的逐元素加法z = x + y。其核心特点是流水同步完全由框架自动完成:
- 将数据搬运和计算拆分为copy_in、compute、copy_out三个子函数;
- 通过 TQue 的
enque/deque操作自动保证流水同步; - 开发者无需手动编写任何同步指令,非常适合学习 Ascend C 框架编程模式。
作为对照,仓库中 examples/01_add 样例实现了手动同步版本的 Add 算子,通过显式调用set_flag/wait_flag指令控制流水同步。建议读者对比学习两个样例,可以直观感受框架自动同步带来的编程体验差异。
运行环境要求
| 类别 | 要求 |
|---|---|
| AI 处理器 | Ascend 910B / 910C |
| CANN 版本 | 社区版 8.5.0.alpha001 及以上 |
注意:
- 样例支持NPU 上板运行(需要 NPU 硬件)和仿真器模式(不需要 NPU 硬件)两种运行方式。仿真器模式运行方式,请参考 运行环境变量配置 完成配置。
- PyTorch 和 torch_npu 的安装,请参考 样例运行验证。
样例规格
| 参数名称 | 输入/输出 | Shape | 数据类型 | 格式 |
|---|---|---|---|---|
| x | 输入 | [8, 2048] | float32 | ND |
| y | 输入 | [8, 2048] | float32 | ND |
| z | 输出 | [8, 2048] | float32 | ND |
整体流程与关键步骤
整体流程
样例的数据流如下:
Global Memory (x_gm, y_gm) │ copy_in: data_copy → TQue.enque ▼ TQue (in_queue_x, in_queue_y) │ compute: TQue.deque → add → TQue.enque ▼ TQue (out_queue_z) │ copy_out: TQue.deque → data_copy ▼ Global Memory (z_gm)数据从 Global Memory 出发,经 copy_in 搬运到 TQue 管理的 Local Memory 队列,compute 阶段从队列取出数据完成加法后推入输出队列,最后由 copy_out 搬运回 Global Memory。
关键步骤
- copy_in—— 使用
asc.data_copy将输入从 Global Memory 搬运到 Local Memory,然后通过enque将 Tensor 推入队列。计算过程中的 Local Memory 通过TQue.alloc_tensor接口获取。 - compute—— 从队列中
deque取出 Tensor,调用asc.add执行逐元素加法,结果通过enque推入输出队列,最后free_tensor释放输入 Tensor。 - copy_out—— 从输出队列
deque取出结果 Tensor,通过asc.data_copy搬运回 Global Memory,最后free_tensor释放。
在此过程中,Ascend C 框架会自动插入对应的同步事件,无需调用set_flag/wait_flag设置同步。
核心接口
| 接口 | 用途 |
|---|---|
asc.TPipe | 统一管理 Device 端内存和同步事件资源,一个 Kernel 函数必须且只能初始化一个 TPipe 对象 |
asc.TQue | 管理流水任务之间的队列通信和同步,支持 alloc_tensor / enque / deque / free_tensor 操作 |
asc.data_copy | 数据搬运(Global Memory↔Local Memory),支持多种搬运场景 |
asc.add | 按元素求和 |
asc.get_block_idx | 获取当前核的索引,用于多核切分 |
源码实现剖析
算子内核主体
add_framework.py 中,Kernel 主体通过@asc.jit装饰器编译,结构如下:
BUFFER_NUM = 2 # BUFFER_NUM should be 1 or 2 USE_CORE_NUM = 8 TILE_NUM = 8 @asc.jit def vadd_kernel(x: asc.GlobalAddress, y: asc.GlobalAddress, z: asc.GlobalAddress, block_length: int, tile_length: asc.ConstExpr[int]): offset = asc.get_block_idx() * block_length x_gm = asc.GlobalTensor() y_gm = asc.GlobalTensor() z_gm = asc.GlobalTensor() x_gm.set_global_buffer(x + offset) y_gm.set_global_buffer(y + offset) z_gm.set_global_buffer(z + offset) pipe = asc.TPipe() in_queue_x = asc.TQue(asc.TPosition.VECIN, BUFFER_NUM) in_queue_y = asc.TQue(asc.TPosition.VECIN, BUFFER_NUM) out_queue_z = asc.TQue(asc.TPosition.VECOUT, BUFFER_NUM) pipe.init_buffer(in_queue_x, BUFFER_NUM, tile_length * x.dtype.sizeof()) pipe.init_buffer(in_queue_y, BUFFER_NUM, tile_length * y.dtype.sizeof()) pipe.init_buffer(out_queue_z, BUFFER_NUM, tile_length * z.dtype.sizeof()) for i in range(TILE_NUM * BUFFER_NUM): copy_in(i, x_gm, y_gm, in_queue_x, in_queue_y, tile_length) compute(z_gm, in_queue_x, in_queue_y, out_queue_z, tile_length) copy_out(i, z_gm, out_queue_z, tile_length)三个子函数
copy_in—— 分配输入队列 Tensor,从 Global Memory 搬运数据并入队:
@asc.jit def copy_in(i: int, x_gm: asc.GlobalAddress, y_gm: asc.GlobalAddress, in_queue_x: asc.TQue, in_queue_y: asc.TQue, tile_length: asc.ConstExpr[int]): x_local = in_queue_x.alloc_tensor(x_gm.dtype) y_local = in_queue_y.alloc_tensor(y_gm.dtype) asc.data_copy(x_local, x_gm[i * tile_length:], tile_length) asc.data_copy(y_local, y_gm[i * tile_length:], tile_length) in_queue_x.enque(x_local) in_queue_y.enque(y_local)compute—— 从输入队列取数、计算并入输出队列,最后释放输入 Tensor:
@asc.jit def compute(z_gm: asc.GlobalTensor, in_queue_x: asc.TQue, in_queue_y: asc.TQue, out_queue_z: asc.TQue, tile_length: asc.ConstExpr[int]): # "z_gm" is passed here to obtain dtype x_local = in_queue_x.deque(z_gm.dtype) y_local = in_queue_y.deque(z_gm.dtype) z_local = out_queue_z.alloc_tensor(z_gm.dtype) asc.add(z_local, x_local, y_local, tile_length) out_queue_z.enque(z_local) in_queue_x.free_tensor(x_local) in_queue_y.free_tensor(y_local)copy_out—— 从输出队列取出结果并搬回 Global Memory:
@asc.jit def copy_out(i: int, z_gm: asc.GlobalTensor, out_queue_z: asc.TQue, tile_length: asc.ConstExpr[int]): z_local = out_queue_z.deque(z_gm.dtype) asc.data_copy(z_gm[i * tile_length:], z_local, tile_length) out_queue_z.free_tensor(z_local)启动与验证
def vadd_launch(x: torch.Tensor, y: torch.Tensor) -> torch.Tensor: z = torch.zeros_like(x) total_length = z.numel() block_length = (total_length + USE_CORE_NUM - 1) // USE_CORE_NUM tile_length = block_length // TILE_NUM // BUFFER_NUM vadd_kernelUSE_CORE_NUM, rt.current_stream() return z启动时通过vadd_kernel[USE_CORE_NUM, rt.current_stream()]指定核数(8 核)与当前流,最终在vadd_custom中使用torch.allclose(z, x + y)完成功能正确性验证。
分块、多核与流水线逻辑
多核切分
- 使用
USE_CORE_NUM = 8个核并行计算。 - 总数据
total_length按核数等分为block_length = (total_length + USE_CORE_NUM - 1) // USE_CORE_NUM。这里采用向上取整的写法,相比 01_add 样例 中total_length // USE_CORE_NUM的整除写法,能更好地处理总数据量不能被核数整除的场景,避免末尾数据遗漏。 - 每个核通过
asc.get_block_idx() * block_length计算自己在 Global Memory 中的偏移量。
分块计算
- 每个核内部将数据进一步切分为
TILE_NUM = 8个 tile。 - 采用双缓冲机制(
BUFFER_NUM = 2),tile_length = block_length // TILE_NUM // BUFFER_NUM。
流水线同步
本样例使用Ascend C 框架自动同步方式。TPipe 通过init_buffer接口为 TQue/TBuf 分配内存,在enque/deque操作过程中自动插入对应的同步事件,从而在双缓冲的配合下实现"搬运与计算在不同 buffer 间流水叠加"。
框架自动同步的底层原理
TPipe / TQue 的前端定义
从 tpipe.py 源码可以看到框架侧的完整接口定义:
- TPipe用于统一管理 Device 端内存等资源,一个 Kernel 函数必须且只能初始化一个 TPipe 对象。其主要功能包括:
- 内存资源管理:通过 TPipe 的
init_buffer接口,可以为 TQue 和 TBuf 分配内存,分别用于队列的内存初始化和临时变量内存的初始化。 - 同步事件管理:通过 TPipe 的
alloc_event_id、release_event_id等接口,可以申请和释放事件 ID,用于同步控制。
- 内存资源管理:通过 TPipe 的
- TQueBind绑定源逻辑位置和目的逻辑位置,根据源位置和目的位置来确定内存分配的位置、插入对应的同步事件,帮助开发者解决内存分配和管理、同步等问题。TQue 是 TQueBind 的简化模式,通常情况下开发者使用 TQue 进行编程。
- TQue继承自 TQueBind,流水任务之间通过队列完成通信和同步。构造时指定逻辑位置(如
TPosition.VECIN/TPosition.VECOUT)和队列深度,例如样例中的asc.TQue(asc.TPosition.VECIN, BUFFER_NUM)。 - TBuf用于管理临时变量的存储空间,存储位置通过模板参数设置为不同的 TPosition 逻辑位置,同样通过 TPipe 的
init_buffer接口初始化。
编译器自动插入同步事件
在@asc.jit编译过程中,pyasc 前端将 Python 代码转换为 ASC-IR(MLIR 方言),随后由编译器 Pass 完成同步事件的自动插入。核心实现在 InsertQueSync.cpp:
enqueueTensors:遍历所有带目的张量(OpWithDst)的操作,若目的张量来自TQueBindAllocTensorOp或TQueBindDequeTensorOp对应的队列,则在该操作之后自动插入TQueBindEnqueTensorOp(即enque);否则插入PipeBarrierOp作为流水屏障。dequeueTensors:基于支配关系(DominanceInfo)分析enque后首次使用该张量的位置,在对应位置自动插入TQueBindDequeTensorOp(即deque),并将后续使用替换为deque返回的 Tensor。canonicalizeBarriers:在函数末尾统一插入PIPE_ALL全流水屏障,并通过 canonicalization 模式化简冗余屏障。syncGetValueOp/syncSetValueOp:对张量的get_value/set_value操作自动包上V_S/S_V事件的SetFlagOp/WaitFlagOp。
也就是说,开发者在 Python 层只需编写alloc_tensor→data_copy→enque→deque→add→enque→deque→data_copy→free_tensor的数据流逻辑,最终的硬件同步指令(set_flag/wait_flag)由编译器在 IR 层面自动生成,这正是"框架自动插入流水同步"的实现本质。
手动同步 vs 框架自动同步
| 对比维度 | 01_add(手动同步) | 02_add_framework(框架自动同步) |
|---|---|---|
| 同步方式 | 显式set_flag/wait_flag | 编译器自动插入同步事件 |
| 同步位置 | 每轮迭代显式插入 MTE2_V、V_MTE3、MTE3_MTE2 三对事件 | enque/deque过程中自动生成 |
| Local Memory 管理 | 手工指定LocalTensor逻辑位置与长度 | 通过TQue.alloc_tensor分配 |
| 编程难度 | 需理解流水线硬件事件模型 | 只需关注数据流,适合框架编程模式学习 |
两个样例中 Kernel 的循环结构高度一致,均采用TILE_NUM * BUFFER_NUM次迭代、双缓冲流水叠加,区别仅在于同步指令由谁书写,对比阅读 01_add/add.py 与 02_add_framework/add_framework.py 可快速理解两种模式的差异。
编译执行
环境配置请参考 quick_start.md(环境准备)、运行环境变量配置(仿真器模式需配置LD_LIBRARY_PATH与LD_PRELOAD=libruntime_camodel.so)以及 样例运行验证(PyTorch/torch_npu 安装)。完成环境配置后,执行如下命令可进行功能验证:
cd pyasc/examples/02_add_framework python3 add_framework.py -r [RUN_MODE] -v [SOC_VERSION]其中脚本参数说明如下:
- RUN_MODE:编译执行方式,可选择 NPU 仿真、NPU 上板,对应参数分别为
Model/NPU。 - SOC_VERSION:昇腾 AI 处理器型号。如果无法确定具体的 SOC_VERSION,则在安装昇腾 AI 处理器的服务器执行
npu-smi info命令进行查询,在查询到的 "Name" 前增加Ascend信息,例如 "Name" 对应取值为xxxyy,实际配置的 SOC_VERSION 值为Ascendxxxyy。
示例如下,Ascend910B1请替换为实际的 AI 处理器型号:
# 仿真器模式 python3 add_framework.py -r Model -v Ascend910B1 # NPU 上板模式 python3 add_framework.py -r NPU -v Ascend910B1脚本内部通过 asc.runtime.config 完成后端(Backend)与平台(Platform)的校验与设置:-r参数仅接受Model/NPU,-v参数需匹配config.Platform枚举值,非法输入会抛出明确的 ValueError。
执行成功后输出:
[INFO] start process sample add_framework. [INFO] Sample add_framework run success.小结
本样例是学习 pyasc 框架编程模式的入门示例:通过 TPipe/TQue 三段式结构(copy_in → compute → copy_out),开发者以纯数据流视角编写算子,流水同步由框架在编译期自动完成。若要进一步理解同步事件的生成细节,可阅读 InsertQueSync.cpp 及对应的 IR 测试用例 insert-que-sync.mlir、erase-sync.mlir;仓库 examples 目录下还提供了 matmul、gelu、rmsnorm 等更复杂的算子样例,可作为进阶学习路径。
【免费下载链接】pyasc本项目为Python用户提供算子编程接口,支持在昇腾AI处理器上加速计算,接口与Ascend C一一对应并遵守Python原生语法。项目地址: https://gitcode.com/cann/pyasc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考