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

Fiber

了解 Effect 中的 Fiber——轻量级虚拟线程,带来强大并发、结构化生命周期与高效的资源管理。

Effect 是一个由 Fiber 驱动的高并发框架。Fiber 是轻量级虚拟线程,具备资源安全的取消能力,为 Effect 中的诸多特性提供了支撑。

在本节中,你将学习 Fiber 的基础知识,并熟悉一些利用 Fiber 的强大底层操作符。

什么是虚拟线程?

JavaScript 本质上是单线程的,也就是说它按单一指令序列执行代码。不过,现代 JavaScript 环境使用事件循环来管理异步操作,从而营造出多任务并行的假象。在这种语境下,虚拟线程(也就是 Fiber)是由 Effect 运行时模拟出来的逻辑线程。它们允许并发执行,而无需依赖 JavaScript 原生并不支持的真多线程。

Fiber 如何工作

Effect 中的所有 effect 都由 Fiber 执行。如果你没有自己创建 Fiber,那么它要么是由你正在使用的某个操作创建的(如果该操作是并发的),要么是由 Effect 运行时系统创建的。

每当一个 effect 被运行时,就会创建一个 Fiber。当并发运行多个 effect 时,会为每个并发 effect 创建一个 Fiber。

即使你编写的是没有任何并发操作的“单线程”代码,也总会至少存在一个 Fiber:执行你的 effect 的那个“主” Fiber。

Effect 的 Fiber 具有定义良好的生命周期,该生命周期基于它所执行的那个 effect。

每个 Fiber 的退出方式要么是失败,要么是成功,取决于它所执行的 effect 是失败还是成功。

Effect 的 Fiber 具有唯一的标识、局部状态以及状态(例如 done、running 或 suspended)。

总结如下:

  • Effect 是更高层的概念,用于描述一段带副作用的计算。它是惰性且不可变的,这意味着它表示一段可能产生值、也可能失败的计算,但并不会立即执行。
  • 而 Fiber 表示 Effect 正在运行的执行过程。它可以被中断,也可以被等待以获取其结果。可以把它看作一种控制和交互正在进行的计算的方式。

Fiber 数据类型

Effect 中的 Fiber 数据类型表示对某个 effect 执行的“句柄”。

以下是 Fiber 的一般形式:

        ┌─── Represents the success type
        │        ┌─── Represents the error type
        │        │
        ▼        ▼
Fiber<Success, Error>

这个类型表明一个 Fiber:

  • 成功并返回类型为 Success 的值
  • 失败并带有类型为 Error 的错误

Fiber 没有 Requirements 类型参数,因为它们只执行那些依赖需求已经被提供好的 effect。

Fork Effect

你可以通过 fork 一个 effect 来创建新的 Fiber。这会在一个新的 Fiber 中启动该 effect,而你会收到指向该 Fiber 的引用。

示例(Fork 一个 Fiber)

在这个示例中,斐波那契计算被 fork 到它自己的 Fiber 中,使它能够独立于主 Fiber 运行。fib10Fiber 的引用可以在之后用于 join 或中断该 Fiber。

import { Effect, Fiber } from "effect"

const fib = (n: number): Effect.Effect<number> =>
  n < 2
    ? Effect.succeed(n)
    : Effect.zipWith(fib(n - 1), fib(n - 2), (a, b) => a + b)

//      ┌─── Effect<Fiber<number, never>, never, never>
//      ▼
const fib10Fiber = Effect.forkChild(fib(10))

await Effect.runPromise(fib10Fiber.pipe(Effect.andThen(Fiber.join))) // => 55

Join Fiber

对 Fiber 最常见的操作之一是 join。使用 Fiber.join 函数,你可以等待某个 Fiber 完成并获取它的结果。被 join 的 Fiber 要么成功、要么失败,而 join 返回的 Effect 反映了该 Fiber 的结果。

示例(Join 一个 Fiber)

import { Effect, Fiber } from "effect"

const fib = (n: number): Effect.Effect<number> =>
  n < 2
    ? Effect.succeed(n)
    : Effect.zipWith(fib(n - 1), fib(n - 2), (a, b) => a + b)

//      ┌─── Effect<Fiber<number, never>, never, never>
//      ▼
const fib10Fiber = Effect.forkChild(fib(10))

const program = Effect.gen(function* () {
  // Retrieve the fiber
  const fiber = yield* fib10Fiber
  // Join the fiber and get the result
  const n = yield* Fiber.join(fiber)
  console.log(n)
  n // => 55
})

await Effect.runPromise(program) // => undefined

Await Fiber

在处理 Fiber 时,Fiber.await 函数是一个很有用的工具。它允许你等待某个 Fiber 完成,并获取关于它是如何结束的详细信息。结果被封装在一个 Exit 值中,让你了解该 Fiber 是成功、失败还是被中断。

示例(等待 Fiber 完成)

import { Effect, Fiber, Exit } from "effect"

const fib = (n: number): Effect.Effect<number> =>
  n < 2
    ? Effect.succeed(n)
    : Effect.zipWith(fib(n - 1), fib(n - 2), (a, b) => a + b)

//      ┌─── Effect<Fiber<number, never>, never, never>
//      ▼
const fib10Fiber = Effect.forkChild(fib(10))

const program = Effect.gen(function* () {
  // Retrieve the fiber
  const fiber = yield* fib10Fiber
  // Await its completion and get the Exit result
  const exit = yield* Fiber.await(fiber)
  console.log(exit)
  exit // => Exit.succeed(55)
})

await Effect.runPromise(program) // => undefined

中断模型

在开发并发应用时,有几种情况需要我们中断其他 Fiber 的执行,例如:

  1. 父 Fiber 可能启动了一些子 Fiber 来执行某项任务,之后父 Fiber 可能认定它不再需要其中某些或全部子 Fiber 的结果。

  2. 两个或多个 Fiber 相互竞争。结果最先计算出来的 Fiber 胜出,而其他所有 Fiber 都不再需要,应当被中断。

  3. 在交互式应用中,用户可能希望停止某些已经在运行的任务,例如点击“停止”按钮以阻止继续下载文件。

  4. 运行时间超出预期的计算,应当通过超时操作予以中止。

  5. 当我们的应用根据用户输入执行计算密集型任务时,如果用户更改了输入,我们就应当取消当前任务并执行另一个任务。

轮询 vs. 异步中断

在中断 Fiber 方面,一种朴素的做法是允许一个 Fiber 强制终止另一个 Fiber。然而这种做法并不理想,因为如果目标 Fiber 正在修改共享状态,强制终止就可能使该状态处于不一致、不可靠的状态。因此,它无法保证共享可变状态的内部一致性。

相反,有两种流行且有效的方案可以解决这个问题:

  1. 半异步中断(轮询式中断):命令式语言通常采用轮询作为一种半异步信号机制,例如 Java。在这种模型中,一个 Fiber 向另一个 Fiber 发送中断请求。目标 Fiber 持续轮询中断状态,检查自己是否收到了来自其他 Fiber 的中断请求。如果检测到中断请求,目标 Fiber 会尽快终止自身。

    采用这种方案时,临界区由 Fiber 自身处理。因此,如果某个 Fiber 正处于临界区中并收到中断请求,它会忽略该中断,并把对中断的处理推迟到临界区之后。

    然而这种做法的一个缺点是:如果程序员忘记定期轮询,目标 Fiber 就可能变得无响应,从而导致死锁。此外,轮询一个全局标志与 Effect 所遵循的函数式范式并不契合。

  2. 异步式中断:在异步式中断中,允许一个 Fiber 终止另一个 Fiber。目标 Fiber 并不负责轮询中断状态。取而代之的是,在临界区中,目标 Fiber 会禁用这些区域的可中断性。这是一种纯函数式方案,不需要轮询全局状态。Effect 的中断模型采用了这一方案,它是一种完全异步的信号机制。

    这种机制克服了忘记定期轮询的缺点。它也与函数式范式完全兼容,因为在纯函数式计算中,我们可以在任意时刻中止计算,除非处于那些禁用了中断的临界区。

中断 Fiber

如果 Fiber 的结果不再被需要,就可以中断它。该操作会立即停止该 Fiber,并安全地运行所有终结器以释放资源。

Fiber.await 不同——后者返回一个描述 Fiber 如何完成的 Exit 值——Fiber.interrupt 返回的是 Effect<void>。它仍会等待该 Fiber 完全终止(在此期间运行它的所有终结器)后才恢复,但不会把该 Fiber 的 Exit 交还给你。

示例(中断一个 Fiber)

import { Effect, Fiber } from "effect"

const program = Effect.gen(function* () {
  // Fork a fiber that runs indefinitely, printing "Hi!"
  const fiber = yield* Effect.forkChild(
    Effect.forever(Effect.log("Hi!").pipe(Effect.delay("10 millis"))),
  )
  yield* Effect.sleep("30 millis")
  // Interrupt the fiber and wait for it to fully terminate
  const result = yield* Fiber.interrupt(fiber)
  console.log(result)
})

Effect.runFork(program)
/*
Output:
timestamp=... level=INFO fiber=#1 message=Hi!
timestamp=... level=INFO fiber=#1 message=Hi!
undefined
*/

默认情况下,Fiber.interrupt 返回的 effect 会等待该 Fiber 完全终止后才恢复。这确保了在前一批 Fiber 完成之前不会启动新的 Fiber,这种行为被称为“背压”(back-pressuring)。

如果你不需要这种等待行为,可以把这个中断操作本身 fork 出去,让主程序不必等待该 Fiber 终止就能继续执行:

示例(Fork 一个中断操作)

import { Effect, Fiber } from "effect"

const program = Effect.gen(function* () {
  const fiber = yield* Effect.forkChild(
    Effect.forever(Effect.log("Hi!").pipe(Effect.delay("10 millis"))),
  )
  yield* Effect.sleep("30 millis")
  const _ = yield* Effect.forkChild(Fiber.interrupt(fiber))
  console.log("Do something else...")
})

await Effect.runPromise(program) // => undefined
/*
Output:
timestamp=... level=INFO fiber=#1 message=Hi!
timestamp=... level=INFO fiber=#1 message=Hi!
Do something else...
*/

对于“发射后不管”式的中断,Fiber 还暴露了立即执行的 interruptUnsafe 方法。与 Fiber.interrupt 不同,它不会等待该 Fiber 的终结器执行完成。

import { Effect, Fiber } from "effect"

const program = Effect.gen(function* () {
  const fiber = yield* Effect.forkChild(
    Effect.forever(Effect.log("Hi!").pipe(Effect.delay("10 millis"))),
  )
  yield* Effect.sleep("30 millis")
  // const _ = yield* Effect.forkChild(Fiber.interrupt(fiber))
  fiber.interruptUnsafe()
  console.log("Do something else...")
})

await Effect.runPromise(program) // => undefined
/*
Output:
timestamp=... level=INFO fiber=#1 message=Hi!
timestamp=... level=INFO fiber=#1 message=Hi!
Do something else...
*/
Interrupting via Effect.interrupt

你也可以使用高层 API Effect.interrupt 来中断 Fiber。更多细节请参阅 Effect.interrupt 文档

组合 Fiber 结果

Fiber 句柄不能直接组合。先 join 每个 Fiber 得到对应的 Effect,再用诸如 Effect.zip 之类的操作符组合这些 effect。

示例(组合两个 Fiber 的结果)

在这个示例中,两个 Fiber 并发运行,它们的结果被组合成一个元组。

import { Effect, Fiber } from "effect"

const program = Effect.gen(function* () {
  // Fork two fibers that each produce a string
  const fiber1 = yield* Effect.forkChild(Effect.succeed("Hi!"))
  const fiber2 = yield* Effect.forkChild(Effect.succeed("Bye!"))

  // Join both fibers and zip their results into a tuple
  const tuple = yield* Effect.zip(Fiber.join(fiber1), Fiber.join(fiber2))
  console.log(tuple)
  tuple // => ["Hi!", "Bye!"]
})

await Effect.runPromise(program) // => undefined

同样的原则也适用于回退行为:join 两个 Fiber,并从第一个被 join 的 effect 的失败中恢复。

示例(提供一个回退 Fiber)

import { Effect, Fiber } from "effect"

const program = Effect.gen(function* () {
  // Fork a fiber that will fail
  const fiber1 = yield* Effect.forkChild(Effect.fail("Uh oh!"))
  // Fork another fiber that will succeed
  const fiber2 = yield* Effect.forkChild(Effect.succeed("Hurray!"))
  // If fiber1 fails, fiber2 will be used as a fallback
  const message = yield* Effect.catchCause(Fiber.join(fiber1), () =>
    Fiber.join(fiber2),
  )
  console.log(message)
  message // => "Hurray!"
})

await Effect.runPromise(program) // => undefined

子 Fiber 的生命周期

当我们 fork Fiber 时,根据 fork 方式的不同,子 Fiber 可以有四种不同的生命周期策略:

  1. 带自动监督的 Fork。如果我们使用普通的 Effect.forkChild 操作,子 Fiber 将由父 Fiber 自动监督。子 Fiber 的生命周期与父 Fiber 的生命周期绑定。这意味着这些 Fiber 要么在自然结束时终止,要么在父 Fiber 被终止时终止。

  2. 在全局作用域中 Fork(Daemon)。有时我们想运行长时间运行的后台 Fiber,它们不依附于父 Fiber,而且我们希望在全局作用域中 fork 它们。任何在全局作用域中 fork 的 Fiber 都会成为 daemon Fiber。这可以通过 Effect.forkDetach 操作实现。由于这些 Fiber 没有父 Fiber,它们不受监督;它们会在自然结束时终止,或在我们的应用被终止时终止。

  3. 在局部作用域中 Fork。有时,我们想运行一个不依附于父 Fiber 的后台 Fiber,但我们希望该 Fiber 生活在局部作用域中。我们可以使用 Effect.forkScoped 在局部作用域中 fork Fiber。这类 Fiber 可以比父 Fiber 存活得更久(因此不受父 Fiber 监督),它们会在自身完成时或局部作用域被关闭时终止。

  4. 在指定作用域中 Fork。这与上一种策略类似,但通过在指定作用域中 fork 子 Fiber,我们可以对子 Fiber 的生命周期进行更细粒度的控制。我们可以使用 Effect.forkIn 操作做到这一点。

带自动监督的 Fork

Effect 遵循结构化并发模型,其中子 Fiber 的生命周期与父 Fiber 绑定。简单来说,一个 Fiber 的寿命取决于其父 Fiber 的寿命。

示例(自动监督的子 Fiber)

在这个场景中,parent Fiber 启动了一个 child Fiber,后者每秒重复打印一条消息。 当 parent Fiber 完成时,child Fiber 将被终止。

import { Effect, Console, Schedule } from "effect"

// Child fiber that logs a message repeatedly every second
const child = Effect.repeat(
  Console.log("child: still running!"),
  Schedule.fixed("1 second"),
)

const parent = Effect.gen(function* () {
  console.log("parent: started!")
  // Child fiber is supervised by the parent
  yield* Effect.forkChild(child)
  yield* Effect.sleep("3 seconds")
  console.log("parent: finished!")
})

await Effect.runPromise(parent) // => undefined
/*
Output:
parent: started!
child: still running!
child: still running!
child: still running!
parent: finished!
*/

这一行为可以扩展到任意层级的嵌套 Fiber,确保 Fiber 的生命周期可预测且受控。

在全局作用域中 Fork(Daemon)

你可以使用 Effect.forkDetach 创建一个长时间运行的后台 Fiber。这类 Fiber 被称为 daemon Fiber,它不依附于父 Fiber 的生命周期,其寿命与全局作用域相关联。即使父 Fiber 被终止,daemon Fiber 仍会继续运行,只有当全局作用域被关闭或该 Fiber 自然完成时才会停止。

示例(创建一个 Daemon Fiber)

这个示例展示了 daemon Fiber 如何在父 Fiber 结束后仍然继续在后台运行。

import { Effect, Console, Schedule } from "effect"

// Daemon fiber that logs a message repeatedly every second
const daemon = Effect.repeat(
  Console.log("daemon: still running!"),
  Schedule.fixed("1 second"),
)

const parent = Effect.gen(function* () {
  console.log("parent: started!")
  // Daemon fiber running independently
  yield* Effect.forkDetach(daemon)
  yield* Effect.sleep("3 seconds")
  console.log("parent: finished!")
})

Effect.runFork(parent)
/*
Output:
parent: started!
daemon: still running!
daemon: still running!
daemon: still running!
parent: finished!
daemon: still running!
daemon: still running!
daemon: still running!
daemon: still running!
daemon: still running!
...etc...
*/

即使父 Fiber 被中断,daemon Fiber 也会继续独立运行。

示例(中断父 Fiber)

在这个示例中,中断父 Fiber 不会影响 daemon Fiber,它会继续在后台运行。

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

// Daemon fiber that logs a message repeatedly every second
const daemon = Effect.repeat(
  Console.log("daemon: still running!"),
  Schedule.fixed("1 second"),
)

const parent = Effect.gen(function* () {
  console.log("parent: started!")
  // Daemon fiber running independently
  yield* Effect.forkDetach(daemon)
  yield* Effect.sleep("3 seconds")
  console.log("parent: finished!")
}).pipe(Effect.onInterrupt(() => Console.log("parent: interrupted!")))

// Program that interrupts the parent fiber after 2 seconds
const program = Effect.gen(function* () {
  const fiber = yield* Effect.forkChild(parent)
  yield* Effect.sleep("2 seconds")
  yield* Fiber.interrupt(fiber) // Interrupt the parent fiber
})

Effect.runFork(program)
/*
Output:
parent: started!
daemon: still running!
daemon: still running!
parent: interrupted!
daemon: still running!
daemon: still running!
daemon: still running!
daemon: still running!
daemon: still running!
...etc...
*/

在局部作用域中 Fork

有时我们想创建一个与局部 scope 绑定的 Fiber,也就是说它的生命周期不依赖于父 Fiber,而是绑定到它被 fork 时所处的局部作用域。这可以使用 Effect.forkScoped 操作来完成。

使用 Effect.forkScoped 创建的 Fiber 可以比其父 Fiber 存活得更久,只有当局部作用域本身被关闭时才会被终止。

示例(在局部作用域中 Fork 一个 Fiber)

在这个示例中,child Fiber 在 parent Fiber 的生命周期结束之后仍继续运行。child Fiber 与局部作用域绑定,只有当作用域结束时才会被终止。

import { Effect, Console, Schedule } from "effect"

// Child fiber that logs a message repeatedly every second
const child = Effect.repeat(
  Console.log("child: still running!"),
  Schedule.fixed("1 second"),
)

//      ┌─── Effect<void, never, Scope>
//      ▼
const parent = Effect.gen(function* () {
  console.log("parent: started!")
  // Child fiber attached to local scope
  yield* Effect.forkScoped(child)
  yield* Effect.sleep("3 seconds")
  console.log("parent: finished!")
})

// Program runs within a local scope
const program = Effect.scoped(
  Effect.gen(function* () {
    console.log("Local scope started!")
    yield* Effect.forkChild(parent)
    // Scope lasts for 5 seconds
    yield* Effect.sleep("5 seconds")
    console.log("Leaving the local scope!")
  }),
)

await Effect.runPromise(program) // => undefined
/*
Output:
Local scope started!
parent: started!
child: still running!
child: still running!
child: still running!
parent: finished!
child: still running!
child: still running!
Leaving the local scope!
*/

在指定作用域中 Fork

有些情况下我们需要更细粒度的控制,因此我们想在一个指定作用域中 fork 一个 Fiber。 我们可以使用 Effect.forkIn 操作,它接收目标作用域作为参数。

示例(在指定作用域中 Fork 一个 Fiber)

在这个示例中,child Fiber 被 fork 到 outerScope 中,这使它能够比内部作用域存活得更久,但在 outerScope 被关闭时仍会被终止。

import { Console, Effect, Schedule } from "effect"

// Child fiber that logs a message repeatedly every second
const child = Effect.repeat(
  Console.log("child: still running!"),
  Schedule.fixed("1 second"),
)

const program = Effect.scoped(
  Effect.gen(function* () {
    yield* Effect.addFinalizer(() =>
      Console.log("The outer scope is about to be closed!"),
    )

    // Capture the outer scope
    const outerScope = yield* Effect.scope

    // Create an inner scope
    yield* Effect.scoped(
      Effect.gen(function* () {
        yield* Effect.addFinalizer(() =>
          Console.log("The inner scope is about to be closed!"),
        )
        // Fork the child fiber in the outer scope
        yield* Effect.forkIn(child, outerScope)
        yield* Effect.sleep("3 seconds")
      }),
    )

    yield* Effect.sleep("5 seconds")
  }),
)

await Effect.runPromise(program) // => undefined
/*
Output:
child: still running!
child: still running!
child: still running!
The inner scope is about to be closed!
child: still running!
child: still running!
child: still running!
child: still running!
child: still running!
child: still running!
The outer scope is about to be closed!
*/

Fiber 何时运行?

被 fork 的 Fiber 会在当前 Fiber 完成或让出之后开始执行。

示例(Fiber 启动过晚,只捕获到一个值)

在下面的示例中,changes Stream 只捕获到一个值 2。 这是因为由 Effect.forkChild 创建的 Fiber 在该值被更新之后才启动。

import { Effect, SubscriptionRef, Stream, Console } from "effect"

const program = Effect.gen(function* () {
  const ref = yield* SubscriptionRef.make(0)
  yield* SubscriptionRef.changes(ref).pipe(
    // Log each change in SubscriptionRef
    Stream.tap((n) => Console.log(`SubscriptionRef changed to ${n}`)),
    Stream.runDrain,
    // Fork a fiber to run the stream
    Effect.forkChild,
  )
  yield* SubscriptionRef.set(ref, 1)
  yield* SubscriptionRef.set(ref, 2)
})

await Effect.runPromise(program) // => undefined
/*
Output:
SubscriptionRef changed to 2
*/

如果你使用 Effect.sleep() 添加一个短暂延迟,或者调用 Effect.yieldNow(),就能让当前 Fiber 让出执行权。这样,被 fork 的 Fiber 就有足够的时间在值被更新之前启动并收集到所有值。

Fiber Execution is Non-Deterministic

请记住,Fiber 执行的时机是不确定的,许多因素都会影响 Fiber 何时启动。不要 认为单次让出总能确保你的 Fiber 在某个特定时刻开始执行。

示例(延迟让 Fiber 捕获所有值)

import { Effect, SubscriptionRef, Stream, Console } from "effect"

const program = Effect.gen(function* () {
  const ref = yield* SubscriptionRef.make(0)
  yield* SubscriptionRef.changes(ref).pipe(
    // Log each change in SubscriptionRef
    Stream.tap((n) => Console.log(`SubscriptionRef changed to ${n}`)),
    Stream.runDrain,
    // Fork a fiber to run the stream
    Effect.forkChild,
  )

  // Allow the fiber a chance to start
  yield* Effect.sleep("100 millis")

  yield* SubscriptionRef.set(ref, 1)
  yield* SubscriptionRef.set(ref, 2)
})

await Effect.runPromise(program) // => undefined
/*
Output:
SubscriptionRef changed to 0
SubscriptionRef changed to 1
SubscriptionRef changed to 2
*/