Effect 事务内存(STM)模块演进:TxDeferred、TxPubSub、TxReentrantLock 与显式Effect.tx事务边界
【免费下载链接】effectBuild production-ready applications in TypeScript项目地址: https://gitcode.com/GitHub_Trending/ef/effect
导读
本文基于 Effect 开源仓库中的 changeset 变更记录(.changeset/pre/add-missing-tx-modules.md),系统梳理软件事务内存(Software Transactional Memory, STM)在 v4 中重构后的完整图景:新增的五个事务模块(TxDeferred、TxPriorityQueue、TxPubSub、TxReentrantLock、TxSubscriptionRef)各自解决什么并发问题、Effect.atomic/Effect.atomicWith为何被移除、显式Effect.tx事务边界如何工作,以及TxSubscriptionRef.changes竞态修复与TxRandom模块移除的动机。读完本文,你将掌握在 Effect 中正确编写"全有或全无"事务代码的完整方法论,并能根据源码理解事务提交、重试与组合的底层机制。
变更总览:一次事务模型的全面重构
该 changeset 标记为"effect": patch,属于一次"增补 + 破坏性重构"的组合变更,核心内容可归纳为四条:
- 新增五个事务(Tx)模块:
TxDeferred、TxPriorityQueue、TxPubSub、TxReentrantLock、TxSubscriptionRef,补齐了 STM 生态中的"一次性完成原语""优先级队列""发布订阅""可重入读写锁""可订阅事务引用"。 - 重构事务模型:移除
Effect.atomic与Effect.atomicWith;所有 Tx 操作统一返回Effect<A, E, Transaction>,要求在边界处显式调用Effect.tx(...)。 - 暴露组合原语:新增
TxPubSub.acquireSubscriber/releaseSubscriber,使订阅的注册与注销可以嵌入更大的事务中;同时修复TxSubscriptionRef.changes的竞态条件,保证当前值一定先于后续更新送达订阅者。 - 清理模块:移除
TxRandom模块。
其中第 2 条是语义上的破坏性变更(breaking change),也是理解整套新 API 的关键钥匙。
事务模型重构:从Effect.atomic到显式Effect.tx
旧模型的局限
在旧模型中,开发者通过Effect.atomic/Effect.atomicWith显式圈定原子区域,事务语义与 Effect 本身耦合较紧。这种设计的痛点在于:事务边界只能由"是否包在 atomic 里"隐式推断,难以将多个事务操作(读、写、阻塞等待、发布)自由组合成一个更大范围的原子工作流,也难以在类型层面表达"这段代码必须在事务中运行"的约束。
新模型:Effect<A, E, Transaction>类型约束
重构后的模型将事务能力建模为 Effect 环境中的一个服务标签Transaction。所有 Tx 操作的返回类型都带有Transaction需求,例如:
TxRef.get/TxRef.set/TxRef.modifyTxQueue.offer/TxQueue.takeTxPubSub.publish/TxPubSub.subscribeTxDeferred.await/TxDeferred.succeed
它们的类型签名形如Effect<A, E, Transaction>,意味着:不使用Effect.tx包裹,这些操作无法在普通(非事务)上下文中运行——类型系统会直接报错,要求开发者显式声明事务边界。
Effect.tx边界的源码实现
Effect.tx定义于 packages/effect/src/Effect.ts 的 Transactions 小节(@category transactions,@since 4.0.0),其核心逻辑为:
export const tx = <A, E, R>(effect: Effect<A, E, R>): Effect<A, E, Exclude<R, Transaction>> => withFiber((fiber) => { let state = Context.getOrUndefined(fiber.context, Transaction) if (state) { // 已处于事务中:组合进当前事务,复用其 journal 与 retry 状态 return effect as Effect<A, E, Exclude<R, Transaction>> } // 仅在最外层边界创建事务状态 state = { journal: new Map(), retry: false } ... })关键语义(与源码注释一致):
- 内层
tx是组合而非嵌套:如果调用时已经处于某个事务中,tx直接复用当前事务的 journal(未提交变更记录)与 retry 标志,而不是新建边界。因此把多个 Tx 操作分别包上Effect.tx,最终会合并成同一个"全有或全无"的事务。 - 最外层
tx负责提交/回滚:只有最外层的tx调用创建事务状态(journal: new Map(), retry: false),执行事务体,然后进入提交或清理流程。 - 乐观并发 + 版本校验:提交前通过
isTransactionConsistent检查 journal 中所有TxRef的版本号是否与事务读取时一致;只要有任何引用被其他事务抢先修改,当前事务立即清空重试(见 Effect.ts 中的isTransactionConsistent、commitTransaction与clearTransaction)。 - 可中断性控制:事务体在
uninterruptibleMask下运行,避免提交过程中被外部中断导致状态不一致;commitTransaction中会递增版本号并通过fiber.currentDispatcher.scheduleTask(pending, 0)唤醒所有等待该引用的 pending 事务。
手动重试原语Effect.txRetry
与Effect.tx配套的是Effect.txRetry(同样位于 Effect.ts 的 transactions 小节)。它是一个Effect<never, never, Transaction>,作用是把当前事务标记为需要重试并中断事务体:
export const txRetry: Effect<never, never, Transaction> = flatMap( Transaction, (state) => { state.retry = true return interrupt } )典型用法是"等待条件成立":在事务体内读取TxRef,若不满足条件则调用Effect.txRetry;事务被挂起,当任何被访问的TxRef因其他事务提交而改变时,该事务会被唤醒并重新执行。awaitPendingTransaction会在所有被访问引用上注册 pending 回调,一旦某个引用被提交,就解除等待并重跑整个事务。
Transaction服务本身
Transaction是定义在 Effect.ts 中的一个Context.Service,承载两个字段:
retry: boolean:记录当前事务是否应当重试;journal: Map<TxRef<any>, { version: number; value: any }>:记录尚未提交的 TxRef 变更及其读取时的版本。
开发者极少直接构造它(示例中仅用于教学演示Effect.provideService),但它解释了整个事务机制的状态来源。
新增事务模块逐个解析
以下五个模块全部标注@since 4.0.0,均以"事务内可安全使用 + 可通过类型系统强制边界"为设计原则。
1. TxDeferred:事务内的一次性完成原语
TxDeferred<A, E>(packages/effect/src/TxDeferred.ts)是一个"写一次"单元格,其完成结果以Result<A, E>形式存储在TxRef<Option<Result<A, E>>>中:
make:创建空 deferred;succeed/fail/done:完成它,只有第一次写入返回true,后续写入是 no-op(返回false),保证"恰好完成一次";await:在事务内等待完成——若尚未完成则调用Effect.txRetry阻塞重试,完成时携带成功值或类型化失败继续;poll:非阻塞地查看当前完成状态,未完成返回Option.none()。
await的实现直接体现了 STM 的阻塞模型:读到Option.none()就Effect.txRetry,靠版本号唤醒机制在别的提交完成后自动重跑。适用于"事务内协调等待某个一次性事件"的场景,例如多个事务等待某个初始化完成。
2. TxPriorityQueue:带排序的优先级事务队列
TxPriorityQueue<A>(packages/effect/src/TxPriorityQueue.ts)的状态存放在TxRef<Chunk<A>>中,元素按构造时传入的Order<A>升序排列:
empty(ord):按指定排序创建空队列;offer:通过二分插入(insertSorted)保持有序;take:取出最小元素(按排序定义),peek只观察不删除;- 空队列上的
peek/take会事务性地等待(自动重试),因此可以与其他事务读写组合成一个原子工作流。
适合多 fiber 通过共享队列协调、且出队/入队需要与其它事务状态变更保持原子性的场景。
3. TxPubSub:事务化发布订阅
TxPubSub<A>(packages/effect/src/TxPubSub.ts)是事务语义下的发布/订阅中枢。每个订阅者持有一个TxQueue,发布时把消息投递给"发布时刻已注册"的所有订阅者队列。
- 四种容量策略(
bounded/dropping/sliding/unbounded),与TxQueue一一对应:bounded:订阅者队列满时,发布者所在事务会重试直到腾出空间(背压);dropping:队列满时丢弃新消息,publish返回false;sliding:队列满时挤掉最旧消息;unbounded:容量为Number.POSITIVE_INFINITY,总是接受。
publish/publishAll:向当前所有订阅者广播,返回是否全部投递成功;hub 关闭后返回false。subscribe:基于Effect.acquireRelease的作用域订阅,返回TxQueue,作用域关闭时自动注销并关闭队列。shutdown/awaitShutdown/isShutdown/size/isEmpty/isFull:生命周期与容量查询,全部是事务操作(awaitShutdown在未关闭时靠Effect.txRetry挂起等待)。
import { Effect, TxPubSub, TxQueue } from "effect" const program = Effect.gen(function*() { const hub = yield* TxPubSub.unbounded<string>() return yield* Effect.scoped( Effect.gen(function*() { const sub = yield* TxPubSub.subscribe(hub) yield* TxPubSub.publish(hub, "hello") return yield* TxQueue.take(sub) }) ) }) await Effect.runPromise(program) // => "hello"4. TxReentrantLock:可重入的事务读写锁
TxReentrantLock(packages/effect/src/TxReentrantLock.ts)允许"多个读者并发持有,或单个写者独占"。其内部状态LockState由两部分组成:
readers: HashMap<number, number>:按 fiberId 记录的读锁持有数;writer: Option<[fiberId: number, count: number]>:当前写者及其重入计数。
可重入语义:锁归属按 fiber 记录,已持有锁的 fiber 可以再次获取(读锁、写锁均可嵌套),并需要对应次数的释放。事务性等待:无法立即获取锁时,获取操作会事务性重试,直到锁可用。模块同时提供手动(acquireReadLock/releaseReadLock等)、作用域(withReadLock/withWriteLock)两种风格,适合在事务中保护"读多写少"的共享资源访问。
5. TxSubscriptionRef:可订阅的事务引用(含竞态修复)
TxSubscriptionRef<A>(packages/effect/src/TxSubscriptionRef.ts)把TxRef<A>(当前值)与事务化 pub/sub 通道组合:订阅者先收到当前值,再收到后续每次已提交更新的新值。
- 读写族:
make、get、set、update、modify、getAndSet、getAndUpdate、updateAndGet——所有变更都在提交时原子地发布给订阅者; - 订阅族:
changes(返回作用域内的TxQueue)与changesStream(返回Stream)。
竞态修复的源码证据:changeset 明确指出"修复TxSubscriptionRef.changes竞态条件,确保当前值先送达"。看 TxSubscriptionRef.ts 中changes的实现:
export const changes = <A>( self: TxSubscriptionRef<A> ): Effect.Effect<TxQueue.TxQueue<A>, never, Scope.Scope> => Effect.acquireRelease( Effect.tx( Effect.gen(function*() { const sub = yield* TxPubSub.acquireSubscriber(self.pubsub) const current = yield* TxRef.get(self.ref) yield* TxQueue.offer(sub, current) return sub }) ), (queue) => Effect.tx(TxPubSub.releaseSubscriber(self.pubsub, queue)) )关键点在于:订阅者注册与"读取当前值 + 投递当前值"发生在同一个事务中。这正是 changeset 中"ExposeTxPubSub.acquireSubscriber/releaseSubscriber"的目的——如果不把它们暴露出来,changes只能先注册订阅再在事务外读当前值,中间就可能漏掉或错序一个已提交的更新;而现在"注册 → 读当前值 → 投递当前值"原子完成,后续所有更新都从该事务的提交点之后开始计数,从根本上消除了"订阅者错过更新"或"收到顺序错乱"的竞态。changesStream则基于changes构建,把每次TxQueue.take包上Effect.tx后重复执行形成Stream。
边界组合实战:用acquireSubscriber组合自定义事务
changeset 将TxPubSub.acquireSubscriber/releaseSubscriber从subscribe的内部实现中独立暴露,使其成为一等 API(@since 4.0.0):
acquireSubscriber(self):创建订阅者队列并注册到 hub,返回Effect<TxQueue<A>, never, Transaction>;releaseSubscriber(self, queue):从 hub 移除该队列并将其关闭,返回Effect<void, never, Transaction>;subscribe本身正是Effect.acquireRelease(Effect.tx(acquireSubscriber(...)), (q) => Effect.tx(releaseSubscriber(...)))的组合。
因此开发者可以在单个事务中完成"注册订阅 + 读取当前值 + 执行其它事务读写 + 注销订阅",例如为TxSubscriptionRef.changes这类复合原语自定义行为。注意releaseSubscriber的 Gotcha:传入的队列必须是本 hub 通过acquireSubscriber获取的,因为释放时会对该队列执行TxQueue.shutdown。
TxRandom的移除与其它配套调整
changeset 同时移除了TxRandom模块。结合变更方向可以推断:随机数生成本质上是"无状态副作用",与"可重试的乐观事务"语义存在根本冲突——事务重试会导致随机序列被重复消费,破坏一致性,因此从 STM 生态中移除。这一移除与Effect.atomic/atomicWith的移除同属"收窄事务边界、强化事务原语纯度"的整体设计取向(属于基于源码结构与变更说明的推断,仓库中已不存在 packages/effect/src/TxRandom.ts 文件)。
此外,本次新增的五个模块均已在 packages/effect/src/index.ts 中导出,使用时直接从"effect"包导入即可:
import { Effect, TxDeferred, TxPriorityQueue, TxPubSub, TxReentrantLock, TxSubscriptionRef } from "effect"迁移指南:从旧模型升级到显式事务边界
对正在使用旧事务 API 的开发者,迁移路径非常直接:
- 删除
Effect.atomic/Effect.atomicWith:不再存在这两个入口(源码中已无定义,所有 Tx 操作的类型都带Transaction环境需求)。 - 把事务体包进
Effect.tx:在事务的最外层边界调用Effect.tx(effect),所有内部的 Tx 操作自然组合进同一事务;嵌套的Effect.tx是组合而非新建边界。 - 利用类型系统排查遗漏:任何 Tx 操作若不处于
Effect.tx内,会产生Transaction环境未满足的类型错误,编译器会精确指出缺失的边界。 - 改用新模块:需要一次性完成原语用
TxDeferred,需要可订阅事务状态用TxSubscriptionRef,需要优先级调度用TxPriorityQueue,需要广播与背压用TxPubSub,需要读写互斥用TxReentrantLock。 - 注意行为细节:
publish的返回布尔值在不同策略下含义不同(bounded 背压重试、dropping 返回false、sliding 挤旧值);releaseSubscriber会关闭队列;TxDeferred只能成功完成一次。
总结
这次 changeset 标志着 Effect v4 的 STM 设计走向成熟:五个新模块补齐了事务生态的能力拼图,Effect.tx显式边界配合Effect<A, E, Transaction>类型约束,把"必须在事务中运行"从约定提升为编译期保证;TxPubSub.acquireSubscriber/releaseSubscriber的暴露使订阅注册、当前值投递与注销可以原子组合;TxSubscriptionRef.changes的竞态修复则展示了"把注册与初值投递放进同一事务"这一关键模式。对于编写高并发、强一致 TypeScript 应用的开发者,这套 API 提供了从"锁 + 手动协调"到"乐观事务 + 自动重试"的范式升级。
【免费下载链接】effectBuild production-ready applications in TypeScript项目地址: https://gitcode.com/GitHub_Trending/ef/effect
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考