Mojo 项目 AsyncRT 并发原语解析:AsyncValue 的类型擦除、状态机与延续传递式异步编程
2026/9/11 16:43:06 网站建设 项目流程

Mojo 项目 AsyncRT 并发原语解析:AsyncValue 的类型擦除、状态机与延续传递式异步编程

【免费下载链接】mojoThe Modular Platform (includes MAX & Mojo)项目地址: https://gitcode.com/GitHub_Trending/mo/mojo

本篇技术指南聚焦 Modular 平台(Mojo/MAX)底层异步运行时 AsyncRT 的核心原语M::AsyncRT::AsyncValue及其配套智能指针AsyncValueRef<T>/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有两点根本性不同:

  1. 不允许调用方阻塞等待std::future::get()/wait()会阻塞当前线程直到值就绪,而AsyncValue不提供任何阻塞接口。相反,调用方通过AsyncValue::andThenSync注册一个闭包(continuation),当值可用时该闭包被调度执行。这是一种典型的**延续传递风格(continuation passing style)**编程,AsyncValue::emplace负责在值就绪时触发所有已注册闭包。
  2. 内置一等错误处理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 error"(AsyncValue.h),任意 C++ 类型都可作为其负载,包括不可拷贝甚至不可移动的类型。

生命周期管理:堆分配 + 引用计数

AsyncValue实例是堆分配的,并采用原子引用计数管理生命周期。因此官方文档强烈建议:

  • 使用 Support/include/Support/RCRef.h 中的RCRef<T>
  • 或使用AsyncValueRef<T>(见 AsyncRT/include/AsyncRT/Runtime/AsyncValueRef.h)。

RCRef是 move-only 的智能指针(避免隐式拷贝引发原子计数开销),但提供显式的.copy()方法;AsyncValueRef<T>则是对AnyAsyncValueRef的模板特化,假定目标AsyncValue存放类型T(AsyncValueRef.h)。

引用计数本身是原子std::atomic<uint32_t> refcount(AsyncValue.h),dropRef对"最后一次引用释放"做了优化:当释放计数恰好等于当前引用数时,仅需 acquire 屏障即可安全销毁(AsyncValue.h)。

类型擦除:AnyAsyncValueRefAsyncValueRef<T>

AsyncValue最终会解析为某个 C++ 类型的值,但这是动态的,且发生在构造之后。因此AsyncValue本身是**类型擦除(type-erased)**的:

  • 用户可以在不知道最终类型的情况下,用AsyncValue::andThenSync()注册闭包;
  • 只有在访问内部数据时才需要类型信息,例如AsyncValue::get<T>()AsyncValue::emplace<T>()

仓库为此提供两个智能指针层级:

类型适用场景
AnyAsyncValueRef持有(无类型的)AsyncValue智能指针,面向类型擦除场景,是 AsyncRT 中操作AsyncValue的主要 API(AnyAsyncValueRef.h)
AsyncValueRef<T>明确知道元素类型T时使用,get()/emplace()无需再传模板参数

AsyncValueRef<T>隐式转换AnyAsyncValueRef。官方建议:知道强类型就用强类型,而动态类型泛型代码(type-generic code)有时做不到这一点。另外AsyncValueRef<Derived>AsyncValueRef<Base>的隐式转换也是允许的(AsyncValueRef.h)。

类型注册:registerType<T>()isType<T>()

AsyncValue可以承载任意 C++ 类型(包括 move-only 甚至不可移动类型),但所有类型在使用前都必须通过AsyncValue::registerType<T>()注册。注册逻辑带来两个好处:

  1. 允许对 payload 和数据进行紧凑存储(dense storage);
  2. 支持有限的类型反射,例如->isType<T>()谓词。

在实现层面,每个AsyncValue携带一个 16 位的TypeID typeID(AsyncValue.h),isType<T>()即比较typeID == TypeID::get<T>()(AsyncValue.h)。对于IndirectAsyncValuetypeID在未解析前为空,解析后才被动态设置(AsyncValue.h)。

从 AsyncValue 取回 CPUDevice:getRuntime()CompactCPUDevicePtr

AsyncRT 运行时被设计为同一进程内可同时存在多个实例,因此某些操作(例如分配一个新的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是模板方法andThen<IsAsync>IsAsync=false特化(AsyncValue.h)。andThen的同步路径会先检查状态:若已就绪则通过runWaiterNow在调用方栈上直接执行 waiter(AsyncValue.h),否则走andThenOutOfLine入队。与之相对,andThenAsyncIsAsync=true)总是将 waiter 作为一个独立任务通过WorkQueue调度(runWaiterLaterworkQueue->addLocalTask,见 AsyncValue.h);若触发线程是未处于 await 循环中的"外部(foreign)"线程,同步 waiter 也会被当作异步任务执行。

捕获状态:延续闭包的妙处与陷阱

这种模式最棒的地方在于:lambda 的捕获列表可以携带任意状态,并且该捕获在 lambda 执行期间一直存活——这意味着你捕获的任何其他RCRef都会在此期间保持引用计数,不会提前析构:

/// When the specified int32_t becomes available, add it to the refcounted /// table. void addToTableWhenReady(AsyncValueRef<int32_t> input, RCRef<TableOfValues> 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(AsyncValueRef<int32_t> input, RCRef<TableOfValues> tablePtr) { // Note that use of `input.` vs `input->`: input.andThenSync([tablePtr = std::move(tablePtr)] (const AsyncValueRef<int32_t> &input) { tablePtr->addValue(input.get()); }); }

要点:

  • 这里不再std::moveinput参数,而是在 lambda 内引入它的影子副本,从而缩小捕获列表、消除求值顺序脚枪;
  • 参数类型可以是const AsyncValueRef<int32_t> &const AnyAsyncValueRef &,取决于手头是AsyncValueRef还是无类型RCRef
  • 由于参数以 const 引用传入,若想延长其生命周期,需要显式.copy()copy()会原子地增加底层AsyncValue的引用计数)。

源码中这一重载被命名为ConsumingWaiterllvm::unique_function<void(AsyncValueRef<T> &&ref)>,它会消费调用处的引用、将其捕获进闭包,并在触发时把引用移交(move)给 waiter(AsyncValueRef.h、AnyAsyncValueRef.h)。andThenSync(同步触发)与andThenAsync(异步任务触发)都提供普通 waiter 与 ConsumingWaiter 两种形式。

单元测试SyncConsuming/AsyncConsuming(AsyncValueTest.cpp)验证了消费语义:waiter 触发时引用已"唯一"(isUnique()),证明 producer 的引用已在emplace时被消费,不会出现多余的悬挂引用。

AsyncValue 的四种状态

AsyncValue可能处于四种状态,其中后两种("value available" 与 "error")被统称为就绪(ready)状态——它们表示 future 已被解析(为值或错误)。所有 waiter 在转入就绪状态时被通知,且一旦进入就绪状态便不可再转出

状态说明进入方式
Unconstructed(未构造)负载尚未构造,所有andThenSync请求排队,等待转入就绪状态AsyncValue::allocate<T>AsyncValueRef<T>::allocate静态方法
Unconstructed (with inline waiter)(带内联 waiter 的未构造)与普通未构造状态行为一致,但第一个 waiter 被保存在 payload 字段中;AsyncValue的普通使用者无需关心此状态内部由运行时自动进入
Value Available(值可用)持有已构造完成的 C++ 值,所有andThenSyncwaiter 已被通知AsyncValue::createReady<T>/AsyncValueRef<T>::createReady直接创建;更常见的是先创建未构造状态,再用emplace(...)转入
Error(错误)表示产生值的计算出错,同时携带诊断(含位置)信息createError直接创建;更典型的是发现未构造值有问题后,用setToError转入

状态机的底层实现

虽然文档层面抽象为四种状态,但源码内部的状态枚举更加精细(AsyncValue.h),State仅占用3 个 bit

  • kUnconstructed(0):初始状态,可迁移到内联 waiter 初始化、kAvailablekError
  • kUnconstructedInitializingInlineWaiter(1):第一个 waiter 初始化时的瞬态,状态感知代码应自旋等待其迁移到+4
  • kUnconstructed{1,2,3,4}ValidOOLWaiterSlots(2-5):用于跟踪第一个出线(out-of-line)WaiterListNode中已占用的槽位数(1~4 个有效 waiter),编码在状态整数中使原子分配槽位可用 compare/exchange 完成;
  • kAvailable(6)与kError(7):终态,不可再迁移。

关键实现机制是waiter 列表与状态被压缩进同一个原子字using WaitersAndState = llvm::PointerIntPair<WaiterListNode *, 3, State, ...>,并用compare_exchange_weak/exchange原子更新(AsyncValue.h)。同时维持不变式:状态就绪时 waiter 列表必须为 null

ConcreteAsyncValue<T>的存储布局同样经过精心设计:

  • AsyncValue基类固定为16 字节kAsyncValueSize),保证派生类中负载的16 字节对齐
  • 负载区域是一个 union:EncodedDiagnostic diagnostic(错误态)/Waiter waiter(未构造态的内联 waiter,避免为第一个 waiter 额外堆分配)/T payload(可用态)(AsyncValue.h);
  • 构造函数通过static_assert(offsetof(ConcreteAsyncValue<T>, payload) == AsyncValue::kAsyncValueSize, ...)硬性保证负载紧随基类、无填充(AsyncValue.h)。

这也印证了文档所说:ConcreteAsyncValue<T>子类把元数据与负载数据连续存放,从而减少分配次数、改善缓存效率;而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 "x+x" 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.isType<int32_t>()) newVal = AsyncValue::createReady<int32_t>(input.get<int32_t>()*2); else { assert(input.isType<std::string>() && "unexpected type"); const std::string &str = input.get<std::string>(); newVal = AsyncValue::createReady<std::string>(str+str); } result->resolveIndirect(std::move(newVal)); }); return result; }

实现层面,IndirectAsyncValue的 union 中存放的是RCRef<AsyncValue> value(解析后)或内联 waiter(未解析时),其析构函数保证:若已就绪则释放指向被解析值的引用;若仍处于未构造且无 waiter,则可安全销毁(AsyncValue.h)。

AsyncValue还提供emplaceIndirect(间接值直接携带具体类型负载,AsyncValue.h)以及resolveIndirect的"消费式"重载resolveIndirectAndDecRef。测试IsUnique(AsyncValueTest.cpp)覆盖了 indirect 值解析前后的唯一性判定:未解析的 indirect 值不允许测试唯一性(与完成过程存在竞态),已解析的 indirect 值需要自身与目标值引用计数都为 1 才算唯一(isUniqueSlow,见 AsyncValue.h)。

实战组合:Typed 与 Any 的生产者-消费者模式

将上述 API 组合起来,即可写出 AsyncRT 风格的标准异步流水线。仓库测试 AsyncRT/unittests/AsyncValueTest.cpp 给出了可直接借鉴的范式:

AsyncValueRef<int> typedProducer(CPUDevice &cpuDevice) { auto result = AsyncValueRef<int>::allocate(cpuDevice); addTask(cpuDevice, [result = result.copy()]() mutable { std::move(result).emplace(1); }); return result; } int typedConsumer(AsyncValueRef<int> result) { return *result + 1; }
AnyAsyncValueRef anyProducer(CPUDevice &cpuDevice) { auto result = AnyAsyncValueRef::allocate<int>(cpuDevice); addTask(cpuDevice, [result = result.copy()]() mutable { std::move(result).emplace<int>(1); }); return result; } int anyConsumer(AnyAsyncValueRef result) { return result.get<int>() + 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 时值得牢记以下要点:

  1. 优先强类型:知道类型就用AsyncValueRef<T>,其get()/emplace()省去模板参数并带isCompatible<T>()断言;动态类型泛型代码再退回到AnyAsyncValueRef
  2. RCRef/AsyncValueRef管理生命周期,不要裸持有AsyncValue*;需要延长生命周期时显式.copy(),不要依赖隐式拷贝。
  3. 正确选择andThenSyncandThenAsyncandThenSync在值已就绪时于调用方栈同步执行,否则由触发线程(或 WorkQueue 任务)执行;andThenAsync总是作为独立任务入队(AsyncValue.h)。
  4. 留意求值顺序脚枪:不要在被等待对象自身的成员调用中std::move该对象;优先使用带const AsyncValueRef<T>&参数(ConsumingWaiter)的重载。
  5. 错误必须显式处理isError()/getDiagnosticIfPresent()/takeDiagnostic()(AsyncValue.h)用于检查与提取错误;未就绪时调用get<T>()会触发断言。
  6. 同一个 Runtime 内的分配与回收:通过AsyncValue自带的getRuntime()CompactCPUDevicePtr)访问所属CPUDevice,让负载数据与计算共同绑定在同一个执行上下文(包括 NUMA 亲缘性与分配器策略,见 AsyncRTRuntime.md)。

AsyncValue 的调试辅助还包括:printDebug(非线程安全的内部表示打印)、getRefCountForDebugging、以及仅在MODULAR_DEBUG构建下启用的存活实例计数getNumAllocatedInstances(AsyncValue.h),便于在测试与排障中定位引用泄漏。

【免费下载链接】mojoThe Modular Platform (includes MAX & Mojo)项目地址: https://gitcode.com/GitHub_Trending/mo/mojo

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询