已发布 上游基线 bf46254 原文 ↗ 在 GitHub 编辑

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 方法保证每个任务结束后许可都会被释放, 即使任务失败或被中断。