
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 / 910CCANN 版本社区版 8.5.0.alpha001 及以上注意样例支持NPU 上板运行需要 NPU 硬件和仿真器模式不需要 NPU 硬件两种运行方式。仿真器模式运行方式请参考 运行环境变量配置 完成配置。PyTorch 和 torch_npu 的安装请参考 样例运行验证。样例规格参数名称输入/输出Shape数据类型格式x输入[8, 2048]float32NDy输入[8, 2048]float32NDz输出[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—— 从输入队列取数、计算并入输出队列最后释放输入 Tensorasc.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 Memoryasc.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 2tile_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用于同步控制。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-IRMLIR 方言随后由编译器 Pass 完成同步事件的自动插入。核心实现在 InsertQueSync.cppenqueueTensors遍历所有带目的张量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_PRELOADlibruntime_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),仅供参考