pyasc 框架自动插入流水同步:基于 TPipe/TQue 的 Add 算子实现详解
2026/9/19 6:20:43 网站建设 项目流程

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_incomputecopy_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]float32ND
y输入[8, 2048]float32ND
z输出[8, 2048]float32ND

整体流程与关键步骤

整体流程

样例的数据流如下:

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_idrelease_event_id等接口,可以申请和释放事件 ID,用于同步控制。
  • 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)的操作,若目的张量来自TQueBindAllocTensorOpTQueBindDequeTensorOp对应的队列,则在该操作之后自动插入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_tensordata_copyenquedequeaddenquedequedata_copyfree_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_PATHLD_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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询