
Mojo 项目 AsyncRT 并发原语解析AsyncValue 的类型擦除、状态机与延续传递式异步编程【免费下载链接】mojoThe Modular Platform (includes MAX Mojo)项目地址: https://gitcode.com/GitHub_Trending/mo/mojo本篇技术指南聚焦 Modular 平台Mojo/MAX底层异步运行时 AsyncRT 的核心原语M::AsyncRT::AsyncValue及其配套智能指针AsyncValueRefT/AnyAsyncValueRef完整讲解其与std::future的本质差异、类型擦除与类型注册机制、四种状态含就绪状态机与 Indirect AsyncValue 间接解析模型并结合仓库源码给出可直接复用的andThenSync链式编程范式与规避 C 求值顺序陷阱的实战写法。概述AsyncValue 在 AsyncRT 中的定位AsyncRT 是 Modular 平台中面向领域无关的并行 CPU 计算的底层库目标是为高性能运行时以及 MLIR、LLD 等并行编译器基础设施提供支撑见 AsyncRT/docs/README.md。M::AsyncRT::CPUDevice是其中管理系统资源的低层并发库它采用库式设计、可与应用内其他线程模型协作、并将线程池等关键策略抽象出来见 AsyncRT/docs/AsyncRTRuntime.md。而AsyncValue正是构建在CPUDevice之上的异步值承载单元是 AsyncRT 中最常用的数据流原语。说明AsyncValue 的完整 API 声明位于 AsyncRT/include/AsyncRT/Runtime/AsyncValue.h其核心实现含 waiter 列表、状态迁移、引用计数销毁路径位于 AsyncRT/lib/Runtime/AsyncValue.cpp。本文所有代码与行为描述均以当前仓库源码为准。与std::future的关键差异从概念上说AsyncValue与 C 标准库的std::future类似——它代表一个将来才可用的值。但它与std::future有两点根本性不同不允许调用方阻塞等待。std::future::get()/wait()会阻塞当前线程直到值就绪而AsyncValue不提供任何阻塞接口。相反调用方通过AsyncValue::andThenSync注册一个闭包continuation当值可用时该闭包被调度执行。这是一种典型的**延续传递风格continuation passing style**编程AsyncValue::emplace负责在值就绪时触发所有已注册闭包。内置一等错误处理。AsyncValue除了可以被一个未来值完成还可以被一个错误值完成——该错误值同时携带位置信息EncodedDiagnostic见 AsyncRT/include/AsyncRT/Support/Diagnostic.h。所有使用者都必须正确处理并传播错误而不是假设等到值就绪就等于拿到了合法值。从源码看AsyncValue的注释将其定位为 a lightweight and generic future type that can be fulfilled by an asynchronously provided value or an errorAsyncValue.h任意 C 类型都可作为其负载包括不可拷贝甚至不可移动的类型。生命周期管理堆分配 引用计数AsyncValue实例是堆分配的并采用原子引用计数管理生命周期。因此官方文档强烈建议使用 Support/include/Support/RCRef.h 中的RCRefT或使用AsyncValueRefT见 AsyncRT/include/AsyncRT/Runtime/AsyncValueRef.h。RCRef是 move-only 的智能指针避免隐式拷贝引发原子计数开销但提供显式的.copy()方法AsyncValueRefT则是对AnyAsyncValueRef的模板特化假定目标AsyncValue存放类型TAsyncValueRef.h。引用计数本身是原子std::atomicuint32_t refcountAsyncValue.hdropRef对最后一次引用释放做了优化当释放计数恰好等于当前引用数时仅需 acquire 屏障即可安全销毁AsyncValue.h。类型擦除AnyAsyncValueRef与AsyncValueRefTAsyncValue最终会解析为某个 C 类型的值但这是动态的且发生在构造之后。因此AsyncValue本身是**类型擦除type-erased**的用户可以在不知道最终类型的情况下用AsyncValue::andThenSync()注册闭包只有在访问内部数据时才需要类型信息例如AsyncValue::getT()或AsyncValue::emplaceT()。仓库为此提供两个智能指针层级类型适用场景AnyAsyncValueRef持有无类型的AsyncValue智能指针面向类型擦除场景是 AsyncRT 中操作AsyncValue的主要 APIAnyAsyncValueRef.hAsyncValueRefT明确知道元素类型T时使用get()/emplace()无需再传模板参数AsyncValueRefT可隐式转换为AnyAsyncValueRef。官方建议知道强类型就用强类型而动态类型泛型代码type-generic code有时做不到这一点。另外AsyncValueRefDerived到AsyncValueRefBase的隐式转换也是允许的AsyncValueRef.h。类型注册registerTypeT()与isTypeT()AsyncValue可以承载任意 C 类型包括 move-only 甚至不可移动类型但所有类型在使用前都必须通过AsyncValue::registerTypeT()注册。注册逻辑带来两个好处允许对 payload 和数据进行紧凑存储dense storage支持有限的类型反射例如-isTypeT()谓词。在实现层面每个AsyncValue携带一个 16 位的TypeID typeIDAsyncValue.hisTypeT()即比较typeID TypeID::getT()AsyncValue.h。对于IndirectAsyncValuetypeID在未解析前为空解析后才被动态设置AsyncValue.h。从 AsyncValue 取回 CPUDevicegetRuntime()与CompactCPUDevicePtrAsyncRT 运行时被设计为同一进程内可同时存在多个实例因此某些操作例如分配一个新的AsyncValue需要一个M::AsyncRT::CPUDevice在手边。这很别扭——因为CPUDevice就像MLIRContext一样几乎是一种全局状态到处传递非常麻烦。好在每个AsyncValue实例始终知道自己来自哪个CPUDevice。可以通过asyncVal-getRuntime()方法获取一个CompactCPUDevicePtr见 AsyncRT/include/AsyncRT/Runtime/CompactCPUDevicePtr.h该类型可以与CPUDevice互换使用提供operator-与operator*。CompactCPUDevicePtr的设计颇具匠心它是一个压缩到 8 位一个字节的CPUDevice*指针借助全局单例CPUDeviceTable维护 CPUDevice 索引到指针的映射索引255表示无 CPUDevice实现CompactCPUDevicePtr.h。正是由于这种紧凑编码每个AsyncValue都能以极小的内存开销携带指向分配它的 Runtime 的回指指针并可通过该 Runtime 的分配器归还内存。结论只要手里有一个AsyncValue就等价于拥有了需要的CPUDevice。测试代码中大量利用了这一点例如AsyncValueTest在构造、emplace、addTask 中直接以*cpuDevice传入AsyncRT/unittests/AsyncValueTest.cpp。用andThenSync串联异步工作构建一系列异步计算时最常见的需求是当某个值可用后入队后续工作。AsyncValue通过andThenSync让这变得非常简单void printWhenReady(AnyAsyncValueRef input) { input-andThenSync([]() { // This prints whenever input becomes ready. printf(input is ready!); }); }行为规则如果andThenSync执行时AsyncValue已经就绪lambda 会被立即执行否则 lambda 被入队等到值可用emplace触发时再执行。从实现看andThenSync是模板方法andThenIsAsync的IsAsyncfalse特化AsyncValue.h。andThen的同步路径会先检查状态若已就绪则通过runWaiterNow在调用方栈上直接执行 waiterAsyncValue.h否则走andThenOutOfLine入队。与之相对andThenAsyncIsAsynctrue总是将 waiter 作为一个独立任务通过WorkQueue调度runWaiterLater→workQueue-addLocalTask见 AsyncValue.h若触发线程是未处于 await 循环中的外部foreign线程同步 waiter 也会被当作异步任务执行。捕获状态延续闭包的妙处与陷阱这种模式最棒的地方在于lambda 的捕获列表可以携带任意状态并且该捕获在 lambda 执行期间一直存活——这意味着你捕获的任何其他RCRef都会在此期间保持引用计数不会提前析构/// When the specified int32_t becomes available, add it to the refcounted /// table. void addToTableWhenReady(AsyncValueRefint32_t input, RCRefTableOfValues tablePtr) { // Watch out for order of evaluation, std::move will corrupt our input // argument. AsyncValue *inputPtr input.getPointer(); inputPtr-andThenSync([input std::move(input), tablePtr std::move(tablePtr)]() { tablePtr-addValue(input.get()); }); }这里存在一个 C 特有的求值顺序order of evaluation脚枪不能直接写input-andThenSync([input std::move(input), ...]因为编译器可能先求值std::move(input)、再加载基表达式中的input从而破坏参数。这种andThenSync风格的另一个缺点它捕获的是被等待值的指针这会增大 lambda 的体积提高其脱离内联表示out-of-line representation的概率。为同时解决这两点可以使用值可用时以引用传入的andThenSync重载/// When the specified int32_t becomes available, add it to the refcounted /// table. void addToTableWhenReady(AsyncValueRefint32_t input, RCRefTableOfValues tablePtr) { // Note that use of input. vs input-: input.andThenSync([tablePtr std::move(tablePtr)] (const AsyncValueRefint32_t input) { tablePtr-addValue(input.get()); }); }要点这里不再std::move掉input参数而是在 lambda 内引入它的影子副本从而缩小捕获列表、消除求值顺序脚枪参数类型可以是const AsyncValueRefint32_t 或const AnyAsyncValueRef 取决于手头是AsyncValueRef还是无类型RCRef由于参数以 const 引用传入若想延长其生命周期需要显式.copy()copy()会原子地增加底层AsyncValue的引用计数。源码中这一重载被命名为ConsumingWaiterllvm::unique_functionvoid(AsyncValueRefT ref)它会消费调用处的引用、将其捕获进闭包并在触发时把引用移交move给 waiterAsyncValueRef.h、AnyAsyncValueRef.h。andThenSync同步触发与andThenAsync异步任务触发都提供普通 waiter 与 ConsumingWaiter 两种形式。单元测试SyncConsuming/AsyncConsumingAsyncValueTest.cpp验证了消费语义waiter 触发时引用已唯一isUnique()证明 producer 的引用已在emplace时被消费不会出现多余的悬挂引用。AsyncValue 的四种状态AsyncValue可能处于四种状态其中后两种value available 与 error被统称为就绪ready状态——它们表示 future 已被解析为值或错误。所有 waiter 在转入就绪状态时被通知且一旦进入就绪状态便不可再转出。状态说明进入方式Unconstructed未构造负载尚未构造所有andThenSync请求排队等待转入就绪状态AsyncValue::allocateT或AsyncValueRefT::allocate静态方法Unconstructed (with inline waiter)带内联 waiter 的未构造与普通未构造状态行为一致但第一个 waiter 被保存在 payload 字段中AsyncValue的普通使用者无需关心此状态内部由运行时自动进入Value Available值可用持有已构造完成的 C 值所有andThenSyncwaiter 已被通知AsyncValue::createReadyT/AsyncValueRefT::createReady直接创建更常见的是先创建未构造状态再用emplace(...)转入Error错误表示产生值的计算出错同时携带诊断含位置信息createError直接创建更典型的是发现未构造值有问题后用setToError转入状态机的底层实现虽然文档层面抽象为四种状态但源码内部的状态枚举更加精细AsyncValue.hState仅占用3 个 bitkUnconstructed0初始状态可迁移到内联 waiter 初始化、kAvailable、kErrorkUnconstructedInitializingInlineWaiter1第一个 waiter 初始化时的瞬态状态感知代码应自旋等待其迁移到4kUnconstructed{1,2,3,4}ValidOOLWaiterSlots2-5用于跟踪第一个出线out-of-lineWaiterListNode中已占用的槽位数1~4 个有效 waiter编码在状态整数中使原子分配槽位可用 compare/exchange 完成kAvailable6与kError7终态不可再迁移。关键实现机制是waiter 列表与状态被压缩进同一个原子字using WaitersAndState llvm::PointerIntPairWaiterListNode *, 3, State, ...并用compare_exchange_weak/exchange原子更新AsyncValue.h。同时维持不变式状态就绪时 waiter 列表必须为 null。ConcreteAsyncValueT的存储布局同样经过精心设计AsyncValue基类固定为16 字节kAsyncValueSize保证派生类中负载的16 字节对齐负载区域是一个 unionEncodedDiagnostic diagnostic错误态/Waiter waiter未构造态的内联 waiter避免为第一个 waiter 额外堆分配/T payload可用态AsyncValue.h构造函数通过static_assert(offsetof(ConcreteAsyncValueT, payload) AsyncValue::kAsyncValueSize, ...)硬性保证负载紧随基类、无填充AsyncValue.h。这也印证了文档所说ConcreteAsyncValueT子类把元数据与负载数据连续存放从而减少分配次数、改善缓存效率而IndirectAsyncValue增加一层间接允许负载类型稍后解析AsyncValue.h。AsyncValue::emplaceAndDecRef的完整流程AsyncValue.h为取出内联 waiter → 在 payload 位置原地构造T→ 通过notifyReadyAndDecRef迁移状态并通知 waiter → 断言旧状态必然处于带内联 waiter 的未构造态保证与andThen并发时的互锁。注意emplace系列的AndDecRef语义在触发任何既有 waiter 之前会移除一个引用计数因此AsyncValue只剩余一个引用也是合法的——waiter 甚至可能在AsyncValue被删除之后才被触发AsyncValue.h。Indirect Async Values先建容器后定类型除了上述四种核心状态你还会遇到一种情况在知道AsyncValue将包含什么 C 类型之前就需要先创建它。此时可以用AsyncValue::createIndirect创建一个特殊的 indirect AsyncValue随后用resolveIndirect方法解析它。顾名思义这增加了一层间接允许先创建一个AsyncValue之后再用一个具体类型的AsyncValue去兑现它。例如某个类型泛型代码需要根据输入类型决定输出类型// This works with both integer and string values forming xx or concat(x,x) // depending on what the argument resolves to. AnyAsyncValueRef genericAsyncDouble(AnyAsyncValueRef input) { // Must create this value before knowing what type input is. AnyAsyncValueRef result AsyncValue::createIndirect(input-getRuntime()); input.andThenSync(result result.copy() { AnyAsyncValueRef newVal; if (input.isTypeint32_t()) newVal AsyncValue::createReadyint32_t(input.getint32_t()*2); else { assert(input.isTypestd::string() unexpected type); const std::string str input.getstd::string(); newVal AsyncValue::createReadystd::string(strstr); } result-resolveIndirect(std::move(newVal)); }); return result; }实现层面IndirectAsyncValue的 union 中存放的是RCRefAsyncValue value解析后或内联 waiter未解析时其析构函数保证若已就绪则释放指向被解析值的引用若仍处于未构造且无 waiter则可安全销毁AsyncValue.h。AsyncValue还提供emplaceIndirect间接值直接携带具体类型负载AsyncValue.h以及resolveIndirect的消费式重载resolveIndirectAndDecRef。测试IsUniqueAsyncValueTest.cpp覆盖了 indirect 值解析前后的唯一性判定未解析的 indirect 值不允许测试唯一性与完成过程存在竞态已解析的 indirect 值需要自身与目标值引用计数都为 1 才算唯一isUniqueSlow见 AsyncValue.h。实战组合Typed 与 Any 的生产者-消费者模式将上述 API 组合起来即可写出 AsyncRT 风格的标准异步流水线。仓库测试 AsyncRT/unittests/AsyncValueTest.cpp 给出了可直接借鉴的范式AsyncValueRefint typedProducer(CPUDevice cpuDevice) { auto result AsyncValueRefint::allocate(cpuDevice); addTask(cpuDevice, [result result.copy()]() mutable { std::move(result).emplace(1); }); return result; } int typedConsumer(AsyncValueRefint result) { return *result 1; }AnyAsyncValueRef anyProducer(CPUDevice cpuDevice) { auto result AnyAsyncValueRef::allocateint(cpuDevice); addTask(cpuDevice, [result result.copy()]() mutable { std::move(result).emplaceint(1); }); return result; } int anyConsumer(AnyAsyncValueRef result) { return result.getint() 1; }这两个范式揭示了emplace的重要语义AnyAsyncValueRef.h同步 producer通常会保留一个引用继续使用因此需要先拷贝再 emplace——ref.copy().emplace(...)异步 producer通常把拷贝放进addTask闭包中任务完成时再 emplace——addTask(cpuDevice, [ref ref.copy()]() mutable { ... ref.emplace(...); })emplace在触发 waiter 前消费掉这个引用从而保证基于 move-vs-copy 或 copy-on-write 优化的 waiter 不会看到 producer 遗留的额外引用。AsyncValueTest中还有更多值得阅读的用例StressAndThen用 5 轮 × 500 个值构造多层依赖 DAG 做压力测试故意超订线程以暴露竞态EmplacingFromTask_DeadlockOnFailure/EmplaceOnForeignThread_DeadlockOnFailure验证在任务内部或外部线程 emplace 时 waiter 不会被同步执行否则会死锁TupleAndThenSync/ArrayCopyingSync等验证 AsyncRT/include/AsyncRT/Runtime/Algorithms.h 中andThenSync/andThenAsync的元组与数组变体AddTaskOverflow_DeadlockOnFailure验证任务队列溢出时不会丢失任务。使用建议与注意事项结合文档与源码使用 AsyncValue 时值得牢记以下要点优先强类型知道类型就用AsyncValueRefT其get()/emplace()省去模板参数并带isCompatibleT()断言动态类型泛型代码再退回到AnyAsyncValueRef。用RCRef/AsyncValueRef管理生命周期不要裸持有AsyncValue*需要延长生命周期时显式.copy()不要依赖隐式拷贝。正确选择andThenSync与andThenAsyncandThenSync在值已就绪时于调用方栈同步执行否则由触发线程或 WorkQueue 任务执行andThenAsync总是作为独立任务入队AsyncValue.h。留意求值顺序脚枪不要在被等待对象自身的成员调用中std::move该对象优先使用带const AsyncValueRefT参数ConsumingWaiter的重载。错误必须显式处理isError()/getDiagnosticIfPresent()/takeDiagnostic()AsyncValue.h用于检查与提取错误未就绪时调用getT()会触发断言。同一个 Runtime 内的分配与回收通过AsyncValue自带的getRuntime()CompactCPUDevicePtr访问所属CPUDevice让负载数据与计算共同绑定在同一个执行上下文包括 NUMA 亲缘性与分配器策略见 AsyncRTRuntime.md。AsyncValue 的调试辅助还包括printDebug非线程安全的内部表示打印、getRefCountForDebugging、以及仅在MODULAR_DEBUG构建下启用的存活实例计数getNumAllocatedInstancesAsyncValue.h便于在测试与排障中定位引用泄漏。【免费下载链接】mojoThe Modular Platform (includes MAX Mojo)项目地址: https://gitcode.com/GitHub_Trending/mo/mojo创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考