CANN Runtime 多 Stream 内存语义同步实战:aclrtValueWait 与 aclrtValueWrite 详解
【免费下载链接】runtime本项目提供CANN运行时组件和维测功能组件。项目地址: https://gitcode.com/cann/runtime
导读
本文以 CANN Runtime 仓库中的官方样例 9_multistream_sync_memory 为主线,系统讲解如何在两个独立 Stream(Stream A 与 Stream B)之间通过 Device 内存上的**值同步(Value Wait/Write)**机制实现跨 Stream 的内存语义同步。读者学完后将掌握aclrtValueWait、aclrtValueWrite两个核心接口的语义、等待模式 flag 的取值、与 Event/Notify 同步机制的差异,以及如何搭建双线程 + 双 Stream 的同步场景并完成编译、运行与结果自校验。
一、场景背景:为什么需要“内存语义同步”
在异构计算中,多个 Stream 上的任务天然是异步并发的。当 Stream A 上的任务需要等待 Stream B 上的任务写入某个数据后才能继续执行时,开发者通常有 Event、Notify 等同步手段可选。但这两类同步机制有一个共同限制:同步双方是流上的任务/主机侧逻辑,算子本身很难直接作为同步参与方。
CANN Runtime 提供了基于通用 Device 内存的内存语义同步机制来补齐这一缺口。正如官方开发指南 03-07 内存语义同步 所述:该机制允许用户基于通用 Device 内存实现同步,并且支持算子作为同步参与方——即算子可以在执行过程中与另一条流进行同步,这是 Event/Notify 机制不具备的能力。
而本文要讲的样例9_multistream_sync_memory,正是用纯 Host 侧 API(aclrtValueWait/aclrtValueWrite)演示了这一机制的最小可运行版本:一条流上的等待任务持续阻塞,直到另一条流上的写任务把指定内存的值写入目标值。
二、样例描述与产品支持
2.1 样例行为
按照 README.md 的描述,本样例会触发两个线程:
- 线程 A(等待方):等待指定内存中的数据满足一定条件后解除阻塞;
- 线程 B(写入方):向指定内存中写入数据。
在线程 B 写入满足条件的数据之前,线程 A 将持续阻塞。这正好模拟了流水线/生产者-消费者模型中“依赖前序数据就绪”的典型场景。
2.2 产品支持情况
| 产品 | 是否支持 |
|---|---|
| Ascend 950PR / Ascend 950DT | √ |
| Atlas A3 训练系列产品 / Atlas A3 推理系列产品 | √ |
| Atlas A2 训练系列产品 / Atlas A2 推理系列产品 | √ |
三、核心原理:基于内存的值等待与值写入
3.1 接口声明与语义
两个核心接口的声明位于头文件 include/external/acl/acl_rt.h:
/** * @brief mem write value * @param [in] devAddr dev addr * @param [in] value write value * @param [in] flag reserved, must be 0 * @param [in] stream asynchronized task stream */ aclError aclrtValueWrite(void* devAddr, uint64_t value, uint32_t flag, aclrtStream stream); /** * @brief mem wait value * @param [in] devAddr dev addr * @param [in] value expect value * @param [in] flag wait mode * @param [in] stream asynchronized task stream */ aclError aclrtValueWait(void* devAddr, uint64_t value, uint32_t flag, aclrtStream stream);两个接口均为异步下发:它们向指定 Stream 中下发一个 wait 任务或 write 任务,任务在 Device 侧按 Stream 顺序执行。
aclrtValueWrite:向devAddr指向的 Device 内存写入value。其flag为保留参数,必须传 0。aclrtValueWait:阻塞等待,直到devAddr指向的 Device 内存中的值满足flag指定的条件。
3.2 等待模式 flag 详解
等待模式由头文件 acl_rt.h 中的宏定义:
| 宏 | 值 | 语义 |
|---|---|---|
ACL_STREAM_WAIT_VALUE_GEQ | 0x00000000U | 等待内存值≥期望值(Greater-or-Equal) |
ACL_STREAM_WAIT_VALUE_EQ | 0x00000001U | 等待内存值==期望值(Equal) |
ACL_STREAM_WAIT_VALUE_AND | 0x00000002U | 等待内存值按位与期望值后结果非 0(按位 AND 命中) |
ACL_STREAM_WAIT_VALUE_NOR | 0x00000003U | 等待内存值按位或期望值后结果为 0(NOR 命中) |
本样例使用的是ACL_STREAM_WAIT_VALUE_EQ,即等待内存值恰好等于目标值(valueCompare = 100)。
3.3 与 Event/Notify 机制的关键差异
从 03-07 内存语义同步 可知,内存语义同步机制有两大特性值得关注:
- 算子可以作为同步参与方:Device 侧算子核函数可以通过读/写同一块 Device 内存参与同步。例如开发指南中展示的 Device 侧示例——
myKernel1向syncMem写 1 并通过dcci指令刷新缓存,myKernel2则用volatile指针配合dcci轮询阻塞直到内存值变为 2。这与 Host 侧样例中的aclrtValueWait/aclrtValueWrite可互相配合使用。 - 同步基于通用 Device 内存:由于没有额外的硬件同步对象,同步用的内存可通过
aclrtMemset/aclrtMemsetAsync初始化和清除,使用上非常灵活。
四、编译与运行
4.1 前置条件
- 已安装 CANN 软件包(默认安装根目录为
/usr/local/Ascend); - 主机环境为 Linux(x86_64 或 aarch64),并具备上述支持列表中的昇腾产品;
- 已下载本仓库源码。
4.2 编译运行步骤
按照 README.md 的步骤:
第 1 步:切换到样例目录
cd ${git_clone_path}/example/1_basic_features/memory/9_multistream_sync_memory其中${git_clone_path}为本仓库克隆到本地的根目录。
第 2 步:设置环境变量
# ${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其中 set_sample_env.sh 是样例公共脚本:它根据当前机器架构(x86_64/aarch64)和已安装的 CANN 软件自动探测并导出SOC_VERSION(昇腾 AI 处理器型号,如 Ascend910_9362、Ascend910B2 等)与ASCENDC_CMAKE_DIR(AscendC 编译器ascendc.cmake所在路径)。
第 3 步:运行样例
bash run.shrun.sh 脚本内部完成“构建 → 安装 → 运行 → 校验”的完整闭环:
set -euo pipefail _ASCEND_CANN_PATH="${ASCEND_HOME_PATH:-}" # ... 检查 ASCEND_HOME_PATH 是否已设置 source "${_ASCEND_CANN_PATH}/bin/setenv.bash" echo "[INFO]: Current compile soc version is ${SOC_VERSION}" rm -rf build && mkdir -p build cmake -B build -DASCEND_CANN_PACKAGE_PATH="${_ASCEND_CANN_PATH}" cmake --build build -j cmake --install build构建使用的 CMakeLists.txt 要点如下:
- 通过
include(${ASCENDC_CMAKE_DIR}/ascendc.cmake)引入 AscendC 构建规则; include_directories(${ASCEND_CANN_PACKAGE_PATH}/include)引入 CANN 头文件(acl/acl.h等);link_directories(${ASCEND_CANN_PACKAGE_PATH}/lib64)定位运行库;target_link_libraries(main PRIVATE ascendcl Threads::Threads):链接 AscendCL 运行时库与系统线程库(样例依赖std::thread)。
五、源码逐段剖析
样例主程序为 main.cpp,下面按执行流程拆解。
5.1 主线程:初始化与资源准备
aclInit(nullptr); int32_t deviceId = 0; aclrtSetDevice(deviceId); uint64_t size = 1 * 1024 * 1024; void* devPtrA; CHECK_ERROR(aclrtMalloc(&devPtrA, size, ACL_MEM_MALLOC_HUGE_FIRST)); INFO_LOG("Allocate memory on the device successfully"); aclrtStream streamA = nullptr; CHECK_ERROR(aclrtCreateStream(&streamA)); aclrtStream streamB = nullptr; CHECK_ERROR(aclrtCreateStream(&streamB));关键点:
aclInit(nullptr)使用默认配置完成 Runtime 初始化,对应aclFinalize()做去初始化;aclrtSetDevice(deviceId)指定 Device 0 作为运算设备;aclrtMalloc申请1 MiBDevice 内存,分配策略为ACL_MEM_MALLOC_HUGE_FIRST(优先申请大页内存,降低 TLB miss);这块内存就是后续用于跨 Stream 同步的“信令内存”;- 创建两条独立 Stream:
streamA(等待方)与streamB(写入方)。
5.2 双线程编排:等待方与写入方
const char* filePath = "file/flag.txt"; uint64_t valueCompare = 100; constexpr uint32_t waitTime = 1000000; std::thread threadA(ThreadWait, streamA, deviceId, devPtrA, valueCompare, filePath); (void)usleep(waitTime); uint64_t valueWrite = 100; std::thread threadB(ThreadWrite, streamB, deviceId, devPtrA, valueWrite, filePath); threadB.join(); threadA.join();- 先启动等待线程
threadA,随后主线程usleep(1000000)(1 秒),确保线程 A 中的 wait 任务已下发到 Stream A 并处于阻塞状态; - 再启动写入线程
threadB,向devPtrA写入目标值 100(与valueCompare一致,配合ACL_STREAM_WAIT_VALUE_EQ即可解除等待); - 两个线程
join后,主线程依次执行aclrtDestroyStreamForce、aclrtFree、aclrtResetDeviceForce、aclFinalize完成资源回收。
5.3 等待线程 ThreadWait
int ThreadWait(aclrtStream stream, int32_t deviceId, void* devPtr, uint64_t valueCompare, const char* filePath) { aclrtSetDevice(deviceId); CHECK_ERROR(aclrtValueWait(devPtr, valueCompare, ACL_STREAM_WAIT_VALUE_EQ, stream)); INFO_LOG("Stream A: wait for data at virtual memory = %p to meet the condition", devPtr); CHECK_ERROR(aclrtSynchronizeStream(stream)); INFO_LOG("Stream A: the data in the specified memory has met the condition, all tasks are complete"); int32_t waitFlag = 0; memory::ReadFileEx(filePath, &waitFlag, sizeof(waitFlag)); INFO_LOG("Flag value read by the waiting thread: %d", waitFlag); return 0; }- 线程在
aclrtSetDevice后立即通过aclrtValueWait(devPtr, 100, ACL_STREAM_WAIT_VALUE_EQ, streamA)向 Stream A 下发等待任务。由于此时devPtrA中的值尚未被写入,该任务在 Device 侧持续阻塞,线程 A 不会返回; aclrtSynchronizeStream阻塞等待 Stream A 上所有任务完成——只有 wait 条件满足、等待任务结束后才会继续;- 解除阻塞后,线程 A 通过
memory::ReadFileEx读取写入线程写下的标志文件file/flag.txt,用于验证等待期间确实发生了阻塞(若未阻塞,读到的是初始值 0 而非 123)。
5.4 写入线程 ThreadWrite
int ThreadWrite(aclrtStream stream, int32_t deviceId, void* devPtr, uint64_t valueWrite, const char* filePath) { int32_t writeFlag = 123; memory::WriteFileEx(filePath, &writeFlag, sizeof(writeFlag)); INFO_LOG("Flag value after the writing thread starts: %d", writeFlag); aclrtSetDevice(deviceId); CHECK_ERROR(aclrtValueWrite(devPtr, valueWrite, 0, stream)); INFO_LOG("Stream B: write data at virtual memory = %p", devPtr); CHECK_ERROR(aclrtSynchronizeStream(stream)); INFO_LOG("Stream B: the data in the specified memory has met the condition, all tasks are complete"); return 0; }- 写入线程先写标志文件(写入整数 123),再下发写任务。这样做的意图是:等待线程解除阻塞后读取该文件,只有读到 123 才能证明它确实是在写入线程开始之后、写任务完成之后才继续执行的;
aclrtValueWrite(devPtr, 100, 0, streamB):flag传 0(符合“保留参数必须为 0”的约束),向devPtrA写入值 100,从而解除 Stream A 上 wait 任务的阻塞;aclrtSynchronizeStream(streamB)确保写任务真正完成。
5.5 标志文件读写工具
样例通过文件系统辅助验证阻塞行为,其实现位于 file_ops.cpp:
WriteFileEx:以O_WRONLY | O_CREAT | O_TRUNC打开文件并写入指定大小数据,写入字节数与请求不符时记录Partial write错误;ReadFileEx:以O_RDONLY打开并读取数据,同样对部分读做错误处理。
这是纯 Host 侧的“实验检测手段”,与 Device 内存同步本身无关,但为验证“等待线程确实被阻塞”提供了可观测证据。
5.6 错误处理宏
样例复用了仓库公共工具 utils.h 中的CHECK_ERROR宏:任何 AscendCL 接口返回非ACL_SUCCESS时,打印错误码并立即返回 -1,保证失败路径可见、可定位。
六、涉及的关键 CANN RUNTIME API
以下是本样例涉及的关键功能点及接口清单(与 README.md 一致,并补充说明):
| 功能类别 | 接口 | 作用 |
|---|---|---|
| 初始化 | aclInit | 初始化配置(此处传nullptr使用默认配置) |
| 初始化 | aclFinalize | 去初始化,释放 Runtime 全局资源 |
| Device 管理 | aclrtSetDevice | 指定用于运算的 Device |
| Device 管理 | aclrtResetDeviceForce | 强制复位当前 Device,回收 Device 上的资源 |
| Stream 管理 | aclrtCreateStream | 创建 Stream |
| Stream 管理 | aclrtSynchronizeStream | 阻塞等待 Stream 上任务全部完成 |
| Stream 管理 | aclrtDestroyStreamForce | 强制销毁 Stream,丢弃所有任务 |
| 内存管理 | aclrtValueWait | 等待指定内存数据满足条件后解除阻塞 |
| 内存管理 | aclrtValueWrite | 向指定内存写入数据 |
| 内存管理 | aclrtMalloc | 申请 Device 内存 |
| 内存管理 | aclrtFree | 释放 Device 内存 |
| 数据传输 | aclrtMemcpy | 通过内存复制实现数据传输(样例用于数据搬运类场景的基础接口) |
说明:
aclrtMemcpy在样例中作为数据传输基础接口被列出;本样例的同步核心是aclrtValueWait/aclrtValueWrite与aclrtSynchronizeStream的组合。
七、运行结果与自动校验
7.1 示例输出
运行成功后的典型输出(来自 README.md):
[INFO] Allocate memory on the device successfully [INFO] Create Stream A successfully [INFO] Create Stream B successfully [INFO] Stream A: wait for data at virtual memory = 0x... to meet the condition [INFO] Start writing to file/flag.txt [INFO] Flag value after the writing thread starts: 123 [INFO] Stream B: write data at virtual memory = 0x... [INFO] Stream B: the data in the specified memory has met the condition, all tasks are complete [INFO] Stream A: the data in the specified memory has met the condition, all tasks are complete [INFO] Flag value read by the waiting thread: 123注意输出顺序的语义:Stream A: wait ...先于Stream B: write ...,而Stream A: the data ... all tasks are complete出现在Stream B完成之后——这正是“等待线程阻塞直到写入方满足条件”的直接证据。
7.2 脚本自动校验逻辑
run.sh 在运行后会对输出做断言:
wait_value=$(awk -F':' '/Flag value read by the waiting thread:/ {gsub(/^ +| +$/, "", $2); print $2; exit}' "${file_path}") write_value_after=$(awk -F':' '/Flag value after the writing thread starts:/ {gsub(/^ +| +$/, "", $2); print $2; exit}' "${file_path}") if [[ -n "${wait_value}" && "${wait_value}" = "${write_value_after}" ]]; then echo "[SUCCESS] Memory semantics synchronization across multiple streams is successful" else echo "[FAILURE] Memory semantics synchronization across multiple streams failed" exit 1 fi校验思路:等待线程读到的标志值(Flag value read by the waiting thread)必须等于写入线程写入的标志值(Flag value after the writing thread starts,即 123)。若相等,说明等待线程确实被阻塞到写入方完成之后才读取文件,跨 Stream 内存语义同步生效;否则脚本以非零退出码报[FAILURE]。通过解析日志文本并比对两个关键数值,使样例具备可自动回归验证的能力。
八、源码与测试层面的佐证
8.1 单元测试对接口行为的约束
仓库单元测试 tests/ut/acl/testcase/acl_runtime_unittest.cpp 覆盖了这两个接口的入参与转调行为:
TEST_F(UTEST_ACL_Runtime, aclrtValueWrite_failed_with_invalid_args) { auto ret = aclrtValueWrite(nullptr, 100, 0, nullptr); EXPECT_EQ(ret, ACL_ERROR_INVALID_PARAM); // devAddr 为空 → 参数错误 // ... mock rtsValueWrite 返回 ACL_ERROR_RT_PARAM_INVALID 时透传 } TEST_F(UTEST_ACL_Runtime, aclrtValueWrite_success) { auto devAddr = reinterpret_cast<void*>(0x1000U); const auto ret = aclrtValueWrite(devAddr, 100, 0, nullptr); EXPECT_EQ(ret, ACL_SUCCESS); } TEST_F(UTEST_ACL_Runtime, aclrtValueWait_failed_with_invalid_args) { auto ret = aclrtValueWait(nullptr, 100, ACL_STREAM_WAIT_VALUE_GEQ, nullptr); EXPECT_EQ(ret, ACL_ERROR_INVALID_PARAM); // ... } TEST_F(UTEST_ACL_Runtime, aclrtValueWait_success) { auto devAddr = reinterpret_cast<void*>(0x1000U); const auto ret = aclrtValueWait(devAddr, 100, ACL_STREAM_WAIT_VALUE_GEQ, nullptr); EXPECT_EQ(ret, ACL_SUCCESS); }从测试可见:当devAddr为空时两个接口均返回ACL_ERROR_INVALID_PARAM;底层实现经由rtsValueWrite/rtsValueWait转调(mock 层在 tests/depends/acl_stub.h 与 tests/depends/runtime/src/runtime_stub.cpp 中声明),接口成功路径返回ACL_SUCCESS。这印证了aclrtValueWait/aclrtValueWrite是 AscendCL 对外标准 API,且入参校验严格。
8.2 与算子级内存语义同步的衔接
若要在真实算子场景中复刻本样例的同步逻辑,可参考 docs/zh/dev_guide/03-07_memory_semantic_synchronization.md 中的 Device 侧写法:
- 写入侧算子对同步内存执行
*flag = 1;后调用dcci(flag, 0, 2)刷新 cache,保证写入对另一条流可见; - 等待侧算子使用
volatile指针轮询while (*flag != 2) { dcci(flag, 0, 2); },配合dcci保证读取的是最新值。
Host 侧的aclrtValueWait/aclrtValueWrite与 Device 侧的该写法共享同一套“基于通用 Device 内存 + 缓存一致性维护”的同步语义,二者可以混合编排在一条同步流水线中。
九、注意事项与实战建议
- 等待值必须与写入值、等待模式匹配:样例中
valueCompare = 100、valueWrite = 100,模式为ACL_STREAM_WAIT_VALUE_EQ。若改用GEQ,写入大于等于目标值的任意值即可解除等待;若使用AND/NOR,则按位运算命中即解除。 aclrtValueWrite的flag必须传 0:这是头文件注释中明确规定的保留参数约束(见 acl_rt.h),传非 0 值属于非法用法。- 同步内存建议显式初始化/清除:由于同步基于通用 Device 内存,推荐在开始同步前用
aclrtMemset/aclrtMemsetAsync将内存初始化为已知值,避免读到残留数据产生误判。 - 区分同步粒度的选择:若只是“流内任务完成”级别的串行化,优先考虑 Event;若需要“某一具体数值就绪”这一数据级条件,且希望算子也能参与同步,则内存语义同步(
aclrtValueWait/aclrtValueWrite)更合适。 - 异常路径处理:样例通过
CHECK_ERROR宏对每个接口做失败即退出的处理;生产代码中建议对aclrtValueWait设置合理的超时/错误码处理,防止死等。
十、小结
9_multistream_sync_memory样例以最小的双线程、双 Stream 结构完整演示了 CANN Runtime 内存语义同步机制的核心:aclrtValueWait在等待流上下发阻塞任务,aclrtValueWrite在写入流上下发写值任务,二者通过一块共享的 Device 内存完成数据级同步,且可通过日志数值比对实现自动化验证。结合 03-07 内存语义同步 的算子级示例,开发者可以进一步将这一机制推广到多 Stream 流水线、生产者-消费者算子编排等真实场景中,作为 Event/Notify 之外又一重要的同步原语使用。
【免费下载链接】runtime本项目提供CANN运行时组件和维测功能组件。项目地址: https://gitcode.com/cann/runtime
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考