Semaphore
学习在 Effect 中使用 semaphore,精确控制并发、管理资源访问,并高效协调异步任务。
semaphore 是一种同步机制,用于管理对共享资源的访问。在 Effect 中,semaphore 可以帮助控制资源访问,或在异步、并发操作中协调任务。
semaphore 就像一种通用化的互斥锁(mutex),它允许一定数量的 permit(许可)被并发地获取和释放。permit 就像票据,让任务或 Fiber 以受控的方式访问共享资源。当没有可用 permit 时,试图获取 permit 的任务会一直等待,直到有 permit 被释放。
创建 Semaphore
Effect.makeSemaphore 函数会用指定数量的 permit 初始化一个 semaphore。
每个 permit 允许一个任务并发地访问资源或执行操作,而多个 permit 则可以实现可配置的并发级别。
示例(创建一个带 3 个 permit 的 Semaphore)
import { Effect } from "effect"
// Create a semaphore with 3 permits
const mutex = Effect.makeSemaphore(3)
withPermits
withPermits 方法允许你指定运行某个 effect 所需的 permit 数量。一旦所需的 permit 可用,它就会运行该 effect,并在任务完成时自动释放这些 permit。
示例(用一个 permit 的 Semaphore 强制任务串行执行)
在这个示例中,三个任务被并发启动,但它们会串行执行,因为只有一个 permit 的 semaphore 一次只允许一个任务继续执行。
import { Effect } from "effect"
const task = Effect.gen(function* () {
yield* Effect.log("start")
yield* Effect.sleep("2 seconds")
yield* Effect.log("end")
})
const program = Effect.gen(function* () {
const mutex = yield* Effect.makeSemaphore(1)
// Wrap the task to require one permit, forcing sequential execution
const semTask = mutex.withPermits(1)(task).pipe(Effect.withLogSpan("elapsed"))
// Run 3 tasks concurrently, but they execute sequentially
// due to the one-permit semaphore
yield* Effect.all([semTask, semTask, semTask], {
concurrency: "unbounded",
})
})
Effect.runFork(program)
/*
Output:
timestamp=... level=INFO fiber=#1 message=start elapsed=3ms
timestamp=... level=INFO fiber=#1 message=end elapsed=2010ms
timestamp=... level=INFO fiber=#2 message=start elapsed=2012ms
timestamp=... level=INFO fiber=#2 message=end elapsed=4017ms
timestamp=... level=INFO fiber=#3 message=start elapsed=4018ms
timestamp=... level=INFO fiber=#3 message=end elapsed=6026ms
*/
示例(使用多个 permit 控制并发任务的执行)
在这个示例中,我们创建一个带五个 permit 的 semaphore,并使用 withPermits(n) 为每个任务分配不同数量的 permit:
import { Effect } from "effect"
const program = Effect.gen(function* () {
const mutex = yield* Effect.makeSemaphore(5)
const tasks = [1, 2, 3, 4, 5].map((n) =>
mutex
.withPermits(n)(Effect.delay(Effect.log(`process: ${n}`), "2 seconds"))
.pipe(Effect.withLogSpan("elapsed")),
)
yield* Effect.all(tasks, { concurrency: "unbounded" })
})
Effect.runFork(program)
/*
Output:
timestamp=... level=INFO fiber=#1 message="process: 1" elapsed=2011ms
timestamp=... level=INFO fiber=#2 message="process: 2" elapsed=2017ms
timestamp=... level=INFO fiber=#3 message="process: 3" elapsed=4020ms
timestamp=... level=INFO fiber=#4 message="process: 4" elapsed=6025ms
timestamp=... level=INFO fiber=#5 message="process: 5" elapsed=8034ms
*/
withPermits 方法保证每个任务结束后 permit 都会被释放,
即使任务失败或被中断。