ARTICLE DETAIL

资讯详情

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

CANN Runtime 多 Device 跨进程 IPC Event 同步实战:生产者-消费者模型从句柄导出到分布式协调

CANN Runtime 多 Device 跨进程 IPC Event 同步实战:生产者-消费者模型从句柄导出到分布式协调 CANN Runtime 多 Device 跨进程 IPC Event 同步实战生产者-消费者模型从句柄导出到分布式协调【免费下载链接】runtime本项目提供CANN运行时组件和维测功能组件。项目地址: https://gitcode.com/cann/runtime导读本文基于 CANN runtime 开源仓库中的1_ipcevent_multi_device样例README_en.md深入讲解如何利用IPC Event在多个 Device 上的多个进程之间实现跨进程任务同步。核心场景是一个生产者进程运行在 Device 0创建并导出 IPC Event 句柄多个消费者进程分别运行在 Device 1、Device 2……导入句柄、在各自的 Stream 中等待事件从而构成一对多的分布式任务协调模型。读完本文你将掌握aclrtCreateEventExWithFlag、aclrtIpcGetEventHandle、aclrtIpcOpenEventHandle、aclrtStreamWaitEvent等关键接口的完整调用链以及如何组织多进程 IPC 同步的工程代码与运行验证。场景与架构一个生产者同步多个消费者在分布式训练、多卡推理等场景中经常需要一个主控进程发布信号多个工作进程分卡执行的协调方式。IPC Event 允许 Event 对象跨进程共享生产者创建一个带ACL_EVENT_IPC标志的 Event将其句柄写入文件每个消费者进程独立启动、读取句柄文件在自己所在的 Device 上打开同一个 IPC Event并在自己的 Stream 中等待该事件。本样例的进程拓扑如下生产者进程proc_a运行在 Device 0创建 IPC Event → 记录事件 → 导出句柄到文件 → 等待所有消费者进程完成。消费者进程proc_b运行在 Device 1、Device 2 等其他 Device每个消费者读取事件句柄 → 打开事件 → 在 Stream 中等待事件 → 模拟自身工作 → 记录事件 → 创建完成标志文件通知生产者。通过这种方式一个生产者可以同时同步多个运行在不同 Device 上的消费者实现分布式任务协调。与单进程版样例 0_ipcevent两个进程、一个消费者相比本样例的关键差异在于消费者数量可配置run.sh中CONSUMER_NUM生产者需要等待全部消费者完成生产者通过轮询每个消费者创建的consumer_i.done完成标志文件来判断整体进度运行前增加 Device 数量预检precheck可执行程序设备不足时输出[SKIP]并以退出码 0 结束便于自动化框架识别为环境不满足而非用例失败。产品支持情况产品是否支持Ascend 950PR/Ascend 950DT支持Atlas A3 训练系列产品/Atlas A3 推理系列产品支持Atlas A2 训练系列产品/Atlas A2 推理系列产品支持上述支持情况与仓库中 Event 管理 API 文档 docs/zh/api_ref/07_event_management.md 记载的产品支持范围一致。需要注意的是IPC 能力的前提是硬件与驱动支持 Event 跨进程共享实际部署时应以当前环境的产品型号为准。构建与运行环境准备样例依赖 CANN 安装环境和示例公共运行脚本。环境安装详情以及通用运行步骤见 example 目录下的 README_en.md。运行前需完成安装 CANN 工具包默认安装在/usr/local/Ascend并设置环境变量ASCEND_INSTALL_PATH确保主机上至少有 2 个可用 Device。运行步骤# ${install_root} 替换为 CANN 安装根目录默认安装在 /usr/local/Ascend source ${install_root}/cann/set_env.sh # 自动识别 SOC_VERSION 和 ASCENDC_CMAKE_DIR source ${git_clone_path}/example/set_sample_env.sh # 编译运行 bash run.sh其中${git_clone_path}为仓库克隆路径。set_sample_env.shexample/set_sample_env.sh负责自动探测 SOC 版本与 AscendC 工具链目录保证 CMake 构建时能正确设置SOC_VERSION与ASCENDC_CMAKE_DIR。消费者数量与 Device 预检机制本样例至少需要 2 个可用 Device。run.sh中CONSUMER_NUM默认配置为 1表示生产者使用 Device 0消费者使用 Device 1如需启动更多消费者请根据实际可用 Device 数量调整CONSUMER_NUM例如有 3 个设备 0、1、2 时可设为 2。run.sh启动进程前会调用precheck程序检查可用 Device 数量详见下文Device 预检实现若设备数不足将输出[SKIP]提示并以退出码 0 结束便于自动化框架识别为环境不满足而非用例失败。run.sh的核心执行流程run.shCONSUMER_NUM1 REQUIRED_DEVICE_COUNT$((CONSUMER_NUM 1)) ./out/precheck ${REQUIRED_DEVICE_COUNT} # precheck 返回 2 表示设备不足直接退出 0其他非 0 为失败 # 后台启动生产者 ./out/proc_a $CONSUMER_NUM 21 | tee producer.log # 逐个后台启动消费者第 i 个消费者运行在 Device i for ((i1; i$CONSUMER_NUM; i)); do ./out/proc_b $i $i 21 | tee consumer_${i}.log done # 等待所有进程结束并检查生产者日志是否出现 finished successfully if grep -q finished successfully producer.log; then echo [SUCCESS] IPC event multi-device synchronization works correctly. fiCMake 构建配置CMakeLists.txt 将三个源码文件编译为三个独立可执行程序并链接 CANN 的 ACL 库add_executable(proc_a proc_a.cpp) add_executable(proc_b proc_b.cpp) add_executable(precheck precheck.cpp) target_link_libraries(proc_a PRIVATE ascendcl) target_link_libraries(proc_b PRIVATE ascendcl) target_link_libraries(precheck PRIVATE ascendcl)头文件路径来自${ASCEND_CANN_PACKAGE_PATH}/include以及样例公共头文件目录ascendcl库则通过${ASCEND_CANN_PACKAGE_PATH}/lib64链接。样例中的错误处理统一使用 example/utils.h 提供的CHECK_ERROR宏——任何接口调用返回非ACL_SUCCESS都会打印具体错误码并立即退出方便快速定位失败点。生产者进程源码剖析proc_a.cpp生产者 proc_a.cpp 的命令行用法为proc_a consumer_num其中consumer_num必须位于 1MAX_CONSUMER_NUM(8) 之间。核心流程如下1. 初始化与创建 StreamCHECK_ERROR(aclInit(nullptr)); int32_t deviceId 0; CHECK_ERROR(aclrtSetDevice(deviceId)); aclrtStream stream nullptr; CHECK_ERROR(aclrtCreateStream(stream));aclInit(nullptr)使用默认配置完成 ACL 初始化aclrtSetDevice(0)将当前进程绑定到 Device 0随后创建一条 Stream用于承载后续 Event 的记录任务。2. 创建 IPC Event 并记录aclrtEvent ipcEvent nullptr; CHECK_ERROR(aclrtCreateEventExWithFlag(ipcEvent, ACL_EVENT_IPC)); CHECK_ERROR(aclrtRecordEvent(ipcEvent, stream)); CHECK_ERROR(aclrtSynchronizeEvent(ipcEvent));这是与普通 Event 最关键的区别创建时通过aclrtCreateEventExWithFlag显式携带ACL_EVENT_IPC标志宏定义为0x00000040U见 include/external/acl/acl_rt.h该 Event 才具备跨进程共享能力。随后aclrtRecordEvent在 Stream 中记录事件aclrtSynchronizeEvent阻塞当前线程直到事件捕获的所有任务执行完成。3. 查询事件状态aclrtEventRecordedStatus status; CHECK_ERROR(aclrtQueryEventStatus(ipcEvent, status)); if (status ! ACL_EVENT_RECORDED_STATUS_COMPLETE) { ... }aclrtQueryEventStatus查询 Event 的录制状态枚举定义见 include/external/acl/acl_rt.htypedef enum aclrtEventRecordedStatus { ACL_EVENT_RECORDED_STATUS_NOT_READY 0, // 未完成 ACL_EVENT_RECORDED_STATUS_COMPLETE 1, // 已完成 } aclrtEventRecordedStatus;样例输出中的event status 1 (1completed)即对应ACL_EVENT_RECORDED_STATUS_COMPLETE。4. 导出 IPC 句柄到文件aclrtIpcEventHandle handle; CHECK_ERROR(aclrtIpcGetEventHandle(ipcEvent, handle)); FILE* fp fopen(./event_handle.bin, wb); fwrite(handle, sizeof(handle), 1, fp); fclose(fp);aclrtIpcGetEventHandle将本进程中的 Event 设置为 IPC Event 并返回其句柄。句柄aclrtIpcEventHandle是固定长度的二进制结构内部为char reserved[ACL_IPC_EVENT_HANDLE_SIZE]见 include/external/acl/acl_rt.h因此样例采用二进制方式整体写入文件./event_handle.bin供消费者进程读取。文件作为句柄在单机多进程间的传输通道是本样例以及 0_ipcevent采用的核心工程手法。5. 等待所有消费者完成生产者通过轮询完成标志文件判断消费者进度while (completed expectedCount elapsed timeoutSec) { completed 0; for (int i 1; i expectedCount; i) { snprintf(flagFile, sizeof(flagFile), ./consumer_%d.done, i); if (access(flagFile, F_OK) 0) completed; } if (completed expectedCount) { sleep(1); elapsed; } }每 1 秒检查一次./consumer_i.done文件是否存在总超时时间为 30 秒超时仍未收齐则报错返回避免进程无限挂死。6. 清理资源与临时文件CHECK_ERROR(aclrtDestroyEvent(ipcEvent)); aclrtDestroyStream(stream); aclrtResetDevice(deviceId); aclFinalize(); remove(./event_handle.bin); // 删除句柄文件 for (int i 1; i consumerNum; i) { remove(./consumer_%d.done); // 删除完成标志文件 }生产者最后会删除所有运行期生成的临时文件保证多次运行时不会残留上一次的句柄或标志文件run.sh在构建阶段也会先清理*.bin、*.done、*.log。消费者进程源码剖析proc_b.cpp消费者 proc_b.cpp 的命令行用法为proc_b device_id consumer_id每个消费者绑定一个不同的 Device。核心流程如下1. 等待句柄文件就绪并读取CHECK_ERROR(aclInit(nullptr)); CHECK_ERROR(aclrtSetDevice(deviceId)); aclrtStream stream nullptr; CHECK_ERROR(aclrtCreateStream(stream)); int timeout 30; while (access(./event_handle.bin, F_OK) ! 0 timeout-- 0) sleep(1);由于生产者和消费者进程是同时由run.sh后台启动的消费者必须先轮询等待生产者写出句柄文件同样 30 秒超时这是进程间通过文件握手的关键一步。2. 打开 IPC EventaclrtIpcEventHandle handle; FILE* fp fopen(./event_handle.bin, rb); fread(handle, sizeof(handle), 1, fp); fclose(fp); aclrtEvent ipcEvent nullptr; CHECK_ERROR(aclrtIpcOpenEventHandle(handle, ipcEvent));aclrtIpcOpenEventHandle在本进程中根据句柄信息返回可用的 Event 指针此后消费者持有的ipcEvent与生产者创建的是同一个底层事件对象。3. 在 Stream 中等待事件并同步CHECK_ERROR(aclrtStreamWaitEvent(stream, ipcEvent)); CHECK_ERROR(aclrtSynchronizeStream(stream));aclrtStreamWaitEvent阻塞指定 Stream 的运行直到生产者记录的那个 Event 完成——这是等待生产者信号的核心同步点aclrtSynchronizeStream再阻塞等待 Stream 上的所有任务执行完。根据接口文档docs/zh/api_ref/07_event_management.mdaclrtStreamWaitEvent支持多个 Stream 等待同一个 Event这正是一对多同步的底层支撑。4. 模拟工作并回写信号usleep(1000000); // 模拟消费者自身工作 1 秒 CHECK_ERROR(aclrtRecordEvent(ipcEvent, stream)); // 记录事件通知其他等待者 CHECK_ERROR(aclrtSynchronizeEvent(ipcEvent));消费者完成自身工作后再次记录事件并同步本样例中已无其他等待者主要是演示 Event 可被反复记录的特性。随后再次用aclrtQueryEventStatus校验事件状态为完成并通过创建./consumer_id.done标志文件通知生产者snprintf(flagFile, sizeof(flagFile), ./consumer_%d.done, consumerId); fp fopen(flagFile, w); fprintf(fp, done); fclose(fp);5. 清理资源CHECK_ERROR(aclrtDestroyEvent(ipcEvent)); aclrtDestroyStream(stream); aclrtResetDevice(deviceId); aclFinalize();与生产者对称地释放 Event、Stream、Device 资源并执行aclFinalize去初始化。Device 预检实现precheck.cppprecheck.cpp 是一个独立的小程序用法为precheck required_device_countCHECK_ERROR(aclInit(nullptr)); uint32_t deviceCount 0; aclError ret aclrtGetDeviceCount(deviceCount); aclFinalize(); if (deviceCount requiredDeviceCount) { INFO_LOG([SKIP] Need at least %d devices, but only %u device(s) are available., ...); return 2; // 环境不满足跳过 } return 0;它通过aclrtGetDeviceCount查询当前可用的 Device 数量与run.sh计算的REQUIRED_DEVICE_COUNT CONSUMER_NUM 1生产者占 1 个 每个消费者各占 1 个比较。返回码约定为0 表示满足、2 表示设备不足run.sh将其转换为退出码 0 实现跳过而非失败、其他非 0 表示预检失败。这种预检先行的设计让样例在设备数量不足的 CI 机器上不会误报失败。CANN Runtime 关键 API 一览本样例涉及的 CANN Runtime 接口按功能划分如下接口声明见 include/external/acl/acl_rt.h详细约束见 docs/zh/api_ref/07_event_management.md功能接口本样例中的作用初始化aclInit/aclFinalize进程启动时初始化、退出前去初始化Device 管理aclrtSetDevice/aclrtResetDevice指定运算 Device退出时复位 Device、回收资源Stream 管理aclrtCreateStream/aclrtSynchronizeStream/aclrtDestroyStream创建承载任务的 Stream阻塞等待任务完成销毁 StreamEvent 管理aclrtCreateEventExWithFlag创建带ACL_EVENT_IPC标志的 IPC EventaclrtRecordEvent在 Stream 中记录 Event异步aclrtSynchronizeEvent阻塞线程直到 Event 捕获的任务完成aclrtQueryEventStatus查询 Event 录制状态0未完成1已完成aclrtIpcGetEventHandle将 Event 设为 IPC Event 并导出句柄aclrtIpcOpenEventHandle在本进程打开 IPC EventaclrtStreamWaitEvent阻塞 Stream 直到指定 Event 完成aclrtDestroyEvent销毁 Event关于 ACL_EVENT_IPC 标志的使用约束根据 docs/zh/api_ref/07_event_management.md 的说明使用ACL_EVENT_IPC创建 Event 时需注意ACL_EVENT_IPC标志不支持与其他 flag 进行位或操作使用该标志创建的 Event 不支持在aclrtResetEvent、aclrtQueryEvent已弃用、aclrtQueryEventWaitStatus、aclrtEventElapsedTime、aclrtEventGetTimestamp、aclrtGetEventId以及模型捕获等场景中使用否则返回报错aclrtIpcGetEventHandle支持同一个 Device 内的多个进程以及跨 Device 的多个进程——这正是本样例能够在多个 Device 上做一对多同步的理论依据。此外本样例属于单机多进程 IPC 场景各进程退出时统一使用aclrtResetDevice释放当前进程的 Device 资源而不是使用aclrtResetDeviceForce强制复位——后者会强制复位并影响同一主机上其他进程的 Device 状态在多进程共存场景下应避免使用。运行结果解读单消费者CONSUMER_NUM1的典型输出如下[INFO] Starting producer on device 0, waiting for 1 consumer(s)... [INFO] Starting consumer 1 on device 1 [INFO] Consumer 1: IPC event opened on device 1 [INFO] Consumer 1: event received, stream synchronized [INFO] Consumer 1: event status 1 (1completed) [INFO] Consumer 1 (device 1) finished successfully. [INFO] Producer: IPC event created on device 0 [INFO] Producer: event status 1 (1completed) [INFO] All 1 consumers have finished. [INFO] Producer: finished successfully. [SUCCESS] IPC event multi-device synchronization works correctly.输出解读与验证点时序上消费者先打印IPC event opened on device 1与event received, stream synchronized表明它成功导入句柄、打开 IPC Event并在自己的 Stream 中等到了生产者记录的事件状态校验生产者与消费者都通过aclrtQueryEventStatus确认事件最终处于1completedACL_EVENT_RECORDED_STATUS_COMPLETE状态握手成功生产者打印All 1 consumers have finished说明轮询到了消费者创建的consumer_1.done标志文件收尾判定run.sh通过grep -q finished successfully producer.log判定整体成功输出[SUCCESS]。若设备不足则会先输出[SKIP]并以退出码 0 结束。设计要点与工程启示句柄文件是单机多进程共享的实用通道IPC Event 句柄是固定长度的二进制结构通过event_handle.bin文件在进程间传递配合消费者轮询等待文件出现的机制规避了进程启动时序问题生产者通过consumer_i.done标志文件聚合多个消费者的完成状态。一对多同步天然成立aclrtStreamWaitEvent支持多个 Stream分布在多个进程、多个 Device等待同一个 IPC Event这是本样例从单消费者扩展到多消费者无需改动底层同步逻辑的根本原因。超时与状态机设计句柄等待、消费者等待、生产者等待均设置了 30 秒超时避免进程因对端异常而无限阻塞每一步关键接口调用后都通过aclrtQueryEventStatus校验状态将异步语义显式化。资源生命周期管理生产者和消费者在退出路径上严格对称——先销毁 Event、再销毁 Stream、复位 Device、最后aclFinalize并清理临时文件保证进程可重复运行在多进程 IPC 场景中坚持使用aclrtResetDevice而非aclrtResetDeviceForce避免影响同机其他进程。延伸阅读单进程版 IPC Event 样例0_ipceventEvent 管理接口全量参考docs/zh/api_ref/07_event_management.md异步任务与 Event 使用指南docs/zh/dev_guide/03-05_event_management.md示例公共工具头文件日志与错误检查宏example/utils.h【免费下载链接】runtime本项目提供CANN运行时组件和维测功能组件。项目地址: https://gitcode.com/cann/runtime创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表