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

Stream 中的错误处理

学习如何处理 Stream 中的错误,实现稳健的恢复、重试与优雅的错误管理,从而保证可靠的流式处理。

从失败中恢复

当处理可能遇到错误的 Stream 时,知道如何优雅地处理这些错误至关重要。Stream.catch 函数是一个强大的工具,它能在失败时恢复,并切换到另一个 Stream。

示例

import { Stream, Effect } from "effect"

const s1 = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.fail("Oh! Error!")),
  Stream.concat(Stream.make(4, 5)),
)

const s2 = Stream.make("a", "b", "c")

const stream = Stream.catch(s1, () => s2)

await Effect.runPromise(Stream.runCollect(stream)) // => [1, 2, 3, "a", "b", "c"]

在这个例子中,s1 遇到了错误,但我们没有终止这个 Stream,而是用 Stream.catch 优雅地切换到 s2。这样即使其中一个 Stream 失败,我们也能继续处理数据。

你也可以在合并两个 Stream 之前,用 Result 数据类型为每一侧打上标记,把 s1 的元素映射为 Result.succeed,把 s2 的元素映射为 Result.fail,从而根据成功或失败来区分两个 Stream 中的元素:

import { Stream, Effect, Result } from "effect"

const s1 = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.fail("Oh! Error!")),
  Stream.concat(Stream.make(4, 5)),
)

const s2 = Stream.make("a", "b", "c")

const stream = Stream.map(s1, Result.succeed).pipe(
  Stream.catch(() => Stream.map(s2, Result.fail)),
)

await Effect.runPromise(Stream.runCollect(stream)) // => [Result.succeed(1), Result.succeed(2), Result.succeed(3), Result.fail("a"), Result.fail("b"), Result.fail("c")]

Stream.catch 相比,Stream.catchFilter 提供了更高级的错误处理能力。借助 Stream.catchFilter,你可以根据所遇到失败的类型和取值来做出决策。

import { Stream, Effect } from "effect"

const s1 = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.fail("Uh Oh!" as const)),
  Stream.concat(Stream.make(4, 5)),
  Stream.concat(Stream.fail("Ouch" as const)),
)

const s2 = Stream.make("a", "b", "c")

const s3 = Stream.make(true, false, false)

const stream = Stream.catch(s1, (error): Stream.Stream<string | boolean> => {
  switch (error) {
    case "Uh Oh!":
      return s2
    case "Ouch":
      return s3
  }
})

await Effect.runPromise(Stream.runCollect(stream)) // => [1, 2, 3, "a", "b", "c"]

在这个例子中,我们有一个 Stream s1,它可能遇到两种不同类型的错误。我们没有像 Stream.catch 那样直接切换到另一个 Stream,而是用 Stream.catch 来精确决定如何处理每一种错误。这种对错误恢复的控制能力,让你可以根据具体的错误情况选择不同的 Stream 或操作。

从 Defect 中恢复

处理 Stream 时,必须为各种失败场景做好准备,包括在 Stream 处理过程中可能出现的 defect。为此,Stream.catchCause 函数提供了一个稳健的解决方案。它让你能够优雅地处理并恢复任何类型的失败。

示例

import { Stream, Effect } from "effect"

const s1 = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.die(new Error("Boom!"))),
  Stream.concat(Stream.make(4, 5)),
)

const s2 = Stream.make("a", "b", "c")

const stream = Stream.catchCause(s1, () => s2)

await Effect.runPromise(Stream.runCollect(stream)) // => [1, 2, 3, "a", "b", "c"]

在这个例子中,s1 可能遇到 defect,但我们没有让应用崩溃,而是用 Stream.catchCause 优雅地切换到另一个 Stream s2。这样即使面对意料之外的问题,应用也能保持稳健并继续处理数据。

从部分错误中恢复

在 Stream 处理中,有些场景需要只从特定类型的失败中恢复。Stream.catchFilterStream.catchCauseFilter 函数正好派上用场,它们让你可以有针对性地处理和缓解错误。

如果你想从某个特定的错误中恢复,可以使用 Stream.catchFilter

import { Stream, Effect, Filter } from "effect"

const s1 = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.fail("Oh! Error!")),
  Stream.concat(Stream.make(4, 5)),
)

const s2 = Stream.make("a", "b", "c")

const stream = Stream.catchFilter(
  s1,
  Filter.fromPredicate((error) => error === "Oh! Error!"),
  () => s2,
)

await Effect.runPromise(Stream.runCollect(stream)) // => [1, 2, 3, "a", "b", "c"]

如果想从某个特定的 cause 中恢复,可以使用 Stream.catchCauseFilter 函数:

import { Stream, Effect, Cause } from "effect"

const s1 = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.die(new Error("Oh! Error!"))),
  Stream.concat(Stream.make(4, 5)),
)

const s2 = Stream.make("a", "b", "c")

const stream = Stream.catchCauseFilter(s1, Cause.findDie, () => s2)

await Effect.runPromise(Stream.runCollect(stream)) // => [1, 2, 3, "a", "b", "c"]

失败时执行清理

Stream.onError 会在 Stream 失败时执行一个 effect,然后保留原本的失败。它适合用于清理或诊断,而不是用于恢复。

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

const stream = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.die(new Error("Oh! Boom!"))),
  Stream.concat(Stream.make(4, 5)),
  Stream.onError(() =>
    Console.log(
      "Stream application closed! We are doing some cleanup jobs.",
    ).pipe(Effect.orDie),
  ),
)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
Stream application closed! We are doing some cleanup jobs.
Error: Oh! Boom!
*/

重试失败的 Stream

有时,Stream 遇到的失败是临时的、可以恢复的。这时 Stream.retry 操作符就派上了用场。它允许你指定一个重试计划,Stream 会按照该计划重试。

示例

import { Stream, Effect, Schedule } from "effect"
import * as NodeReadLine from "node:readline"

const stream = Stream.make(1, 2, 3).pipe(
  Stream.concat(
    Stream.fromEffect(
      Effect.gen(function* () {
        const s = yield* readLine("Enter a number: ")
        const n = parseInt(s)
        if (Number.isNaN(n)) {
          return yield* Effect.fail("NaN")
        }
        return n
      }),
    ).pipe(Stream.retry(Schedule.exponential("1 second"))),
  ),
)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
Enter a number: a
Enter a number: b
Enter a number: c
Enter a number: 4
[ 1, 2, 3, 4 ]
*/

const readLine = (message: string): Effect.Effect<string> =>
  Effect.promise(
    () =>
      new Promise((resolve) => {
        const rl = NodeReadLine.createInterface({
          input: process.stdin,
          output: process.stdout,
        })
        rl.question(message, (answer) => {
          rl.close()
          resolve(answer)
        })
      }),
  )

在这个例子中,Stream 要求用户输入一个数字,但如果输入了无效值(例如 “a”、“b”、“c”),就会以 “NaN” 失败。不过,我们使用了带指数退避计划的 Stream.retry,这意味着它会在逐渐变长的延迟之后重试。这样我们就能处理临时错误,并最终收集到合法的输入。

细化错误

处理 Stream 时,有些场景需要只保留特定的错误,并用其余错误终止 Stream。你可以通过在 Stream.catch 内对错误进行模式匹配来实现:对想保留的错误重新失败,对其余错误则让它 die(Stream.die)。

示例

import { Stream, Option, Effect, Exit } from "effect"

const stream = Stream.fail(new Error())

const res = Stream.catch(stream, (error) => {
  const refined =
    error instanceof SyntaxError ? Option.some(error) : Option.none()
  return Option.isSome(refined) ? Stream.fail(refined.value) : Stream.die(error)
})

await Effect.runPromiseExit(Stream.runCollect(res)) // => Exit.die(new Error())

在这个例子中,stream 最初以一个通用的 Error 失败。不过,res 会过滤并只保留 SyntaxError 类型的错误,让 Stream 用这些错误重新失败。任何其他错误都会通过 Stream.die 变成 defect,从而终止 Stream。

超时

处理 Stream 时,有些场景需要处理超时,例如当 Stream 在一段时长内没有产出一个值时就终止它。本节我们将探讨如何使用各种操作符来管理超时。

timeout

Stream.timeout 操作符允许你为 Stream 设置超时。如果 Stream 在指定时长内没有产出一个值,它就会终止。

import { Stream, Effect } from "effect"

const stream = Stream.fromEffect(Effect.never).pipe(Stream.timeout("2 seconds"))

await Effect.runPromise(Stream.runCollect(stream)) // => []

timeoutFail

Stream.timeoutOrElse 操作符把超时与自定义的失败消息结合起来。如果 Stream 超时,它就会以指定的错误消息失败。

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

const stream = Stream.fromEffect(Effect.never).pipe(
  Stream.timeoutOrElse({
    duration: "2 seconds",
    orElse: () => Stream.fail("timeout"),
  }),
)

await Effect.runPromiseExit(Stream.runCollect(stream)) // => Exit.fail("timeout")

timeoutFailCause

Stream.timeoutOrElse 类似,Stream.timeoutOrElse 把超时与自定义的失败 cause 结合起来。如果 Stream 超时,它就会以指定的 cause 失败。

import { Stream, Effect, Cause, Exit } from "effect"

const stream = Stream.fromEffect(Effect.never).pipe(
  Stream.timeoutOrElse({
    duration: "2 seconds",
    orElse: () => Stream.failCause(Cause.die("timeout")),
  }),
)

await Effect.runPromiseExit(Stream.runCollect(stream)) // => Exit.die("timeout")

timeoutTo

Stream.timeoutOrElse 操作符允许你在第一个 Stream 于指定时长内没有产出值时切换到另一个 Stream。

import { Stream, Effect } from "effect"

const stream = Stream.fromEffect(Effect.never).pipe(
  Stream.timeoutOrElse({
    duration: "2 seconds",
    orElse: () => Stream.make(1, 2, 3),
  }),
)

await Effect.runPromise(Stream.runCollect(stream)) // => [1, 2, 3]