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

跟踪 Fiber

使用 FiberSet 和 FiberMap 跟踪成组的 Fiber。

Effect 提供了两个用于跟踪成组 Fiber 的结构化并发原语:FiberSet 用于无键的集合,FiberMap 用于以任意值为键的集合。当 Fiber 完成时,它们都会自动将其移除;当拥有它们的 Scope 关闭时,它们会中断所有剩余的 Fiber。

使用 FiberSet 跟踪 Fiber

FiberSet<A, E> 会收集 Fiber,以便将它们作为一个整体来观察、join 或中断。你可以用 FiberSet.run 向其中添加 Fiber(fork 一个 effect 并跟踪产生的 Fiber),也可以用 FiberSet.add(跟踪一个你已经 fork 出来的 Fiber);并可以通过 FiberSet.size 或直接遍历它来查看集合。

示例(监控 Fiber 数量)

在这个示例中,我们在计算一个斐波那契数时,周期性地监控应用中正在运行的 Fiber 数量。程序为每个递归步骤 fork 两个子 Fiber,并把它们都加入一个共享的 FiberSet;同时,一个独立的监控 Fiber 会按计划记录该集合的大小,直到计算结束。

import { Effect, Fiber, FiberSet, Schedule } from "effect"

// Main program that monitors fibers while calculating a Fibonacci number
const program = Effect.gen(function* () {
  // Create a FiberSet to track child fibers
  const set = yield* FiberSet.make<number>()

  // Start a Fibonacci calculation, forking every recursive step into the set
  const fibFiber = yield* Effect.forkChild(fib(10, set))

  // Start monitoring the fibers, logging the FiberSet's size every 20ms
  const monitorFiber = yield* Effect.forkChild(
    monitorFibers(set).pipe(Effect.repeat(Schedule.spaced("20 millis"))),
  )

  // Wait for the Fibonacci calculation to finish, then stop the monitor
  const result = yield* Fiber.join(fibFiber)
  yield* Fiber.interrupt(monitorFiber)

  console.log(`fibonacci result: ${result}`)
  // The final result is deterministic even though the intermediate
  // "number of fibers" logs above are racy and vary between runs
  result // => 55
}).pipe(Effect.scoped)

// Function to monitor and log the number of active fibers
const monitorFibers = (set: FiberSet.FiberSet<number>) =>
  Effect.gen(function* () {
    const count = yield* FiberSet.size(set) // Get the current number of tracked fibers
    console.log(`number of fibers: ${count}`)
  })

// Recursive Fibonacci calculation, adding a fiber to the set for each recursive step
const fib = (
  n: number,
  set: FiberSet.FiberSet<number>,
): Effect.Effect<number> =>
  Effect.gen(function* () {
    if (n <= 1) {
      return n
    }
    yield* Effect.sleep("30 millis") // Simulate work by delaying

    // Fork two fibers for the recursive Fibonacci calls, tracked by the FiberSet
    const fiber1 = yield* FiberSet.run(set, fib(n - 2, set))
    const fiber2 = yield* FiberSet.run(set, fib(n - 1, set))

    // Join the fibers to retrieve their results
    const v1 = yield* Fiber.join(fiber1)
    const v2 = yield* Fiber.join(fiber2)

    return v1 + v2 // Combine the results
  })

await Effect.runPromise(program)
/*
Example Output:
number of fibers: 0
number of fibers: 0
number of fibers: 2
number of fibers: 6
number of fibers: 6
number of fibers: 14
number of fibers: 30
number of fibers: 30
number of fibers: 55
number of fibers: 62
number of fibers: 62
number of fibers: 35
number of fibers: 8
number of fibers: 8
*/
Reference documentation

FiberSet 操作的完整列表请参见 FiberSet 模块参考,其中包括 FiberSet.join (只要有被跟踪的 Fiber 失败,就让父 Fiber 失败)和 FiberSet.awaitEmpty (等待直到所有被跟踪的 Fiber 都已完成)。

使用 FiberMap 跟踪 Fiber

FiberMap<K, A, E> 的行为与 FiberSet 类似,但每个 Fiber 都跟踪在一个键之下。当你之后需要查找、替换或中断某个特定的 Fiber 时,这很有用,例如为每个已连接的客户端分配一个 Fiber,并以客户端 ID 为键。为某个已被占用的键设置新的 Fiber 时,会先中断上一个 Fiber。

关于完整的 FiberMap API,请参见 FiberMap 模块参考。