Semaphore
学习在 Effect 中使用 semaphore,精确控制并发、管理资源访问,并高效协调异步任务。
semaphore 是一种同步机制,用于管理对共享资源的访问。在 Effect 中,semaphore 可以帮助控制资源访问,或在异步、并发操作中协调任务。
semaphore 就像一种通用化的互斥锁(mutex),它允许一定数量的**许可(permit)**被并发地获取和释放。许可就像票据,让任务或 Fiber 以受控的方式访问共享资源。当没有可用许可时,试图获取许可的任务会一直等待,直到有许可被释放。
创建 Semaphore
Semaphore.make 函数会用指定数量的许可初始化一个 semaphore。
每个许可允许一个任务并发地访问资源或执行操作,而多个许可则可以实现可配置的并发级别。
示例(创建一个带 3 个许可的 Semaphore)
import { Effect, Semaphore } from "effect"
// Create a semaphore with 3 permits
const mutex = Semaphore.make(3)
// The semaphore was created with exactly 3 permits available
const acquired = await Effect.runPromise(
mutex.pipe(Effect.flatMap((sem) => Semaphore.take(sem, 3))),
)
acquired // => 3
withPermits
withPermits 方法允许你指定运行某个 effect 所需的许可数量。一旦所需的许可可用,它就会运行该 effect,并在任务完成时自动释放这些许可。
示例(用一个许可的 Semaphore 强制任务串行执行)
在这个示例中,三个任务被并发启动,但它们会串行执行,因为只有一个许可的 semaphore 一次只允许一个任务继续执行。
import { Effect, Semaphore } 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* Semaphore.make(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",
})
})
await Effect.runPromise(program) // => undefined
/*
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
*/
示例(使用多个许可控制并发任务的执行)
在这个示例中,我们创建一个带五个许可的 semaphore,并使用 withPermits(n) 为每个任务分配不同数量的许可:
import { Effect, Semaphore } from "effect"
const program = Effect.gen(function* () {
const mutex = yield* Semaphore.make(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" })
})
await Effect.runPromise(program) // => undefined
/*
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
*/
Permit Release Guarantee
withPermits 方法保证每个任务结束后许可都会被释放,
即使任务失败或被中断。