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

流中的错误处理

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

从失败中恢复

处理可能出错的流时,知道如何优雅地应对这些错误至关重要。Stream.orElse 函数是一个强大的工具,它能在出错时从失败中恢复并切换到备用的流。

示例

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.orElse(s1, () => s2)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: "Chunk",
  values: [ 1, 2, 3, "a", "b", "c" ]
}
*/

在这个例子中,s1 遇到了错误,但我们没有终止整个流,而是通过 Stream.orElse 优雅地切换到了 s2。这保证了即使其中一个流出错,我们也能继续处理数据。

还有一个名为 Stream.orElseEither 的变体,它使用 Either 数据类型,根据成功与失败来区分两个流中的元素:

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.orElseEither(s1, () => s2)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: "Chunk",
  values: [
    {
      _id: "Either",
      _tag: "Left",
      left: 1
    }, {
      _id: "Either",
      _tag: "Left",
      left: 2
    }, {
      _id: "Either",
      _tag: "Left",
      left: 3
    }, {
      _id: "Either",
      _tag: "Right",
      right: "a"
    }, {
      _id: "Either",
      _tag: "Right",
      right: "b"
    }, {
      _id: "Either",
      _tag: "Right",
      right: "c"
    }
  ]
}
*/

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

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.catchAll(s1, (error): Stream.Stream<string | boolean> => {
  switch (error) {
    case "Uh Oh!":
      return s2
    case "Ouch":
      return s3
  }
})

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: "Chunk",
  values: [ 1, 2, 3, "a", "b", "c" ]
}
*/

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

从 Defect 中恢复

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

示例

import { Stream, Effect } from "effect"

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

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

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

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: "Chunk",
  values: [ 1, 2, 3, "a", "b", "c" ]
}
*/

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

从部分错误中恢复

在流处理中,有时你可能只需要从特定类型的失败中恢复。Stream.catchSomeStream.catchSomeCause 函数正是为此而生,它们允许你有选择地处理和缓解错误。

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

import { Stream, Effect, Option } 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.catchSome(s1, (error) => {
  if (error === "Oh! Error!") {
    return Option.some(s2)
  }
  return Option.none()
})

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: "Chunk",
  values: [ 1, 2, 3, "a", "b", "c" ]
}
*/

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

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

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

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

const stream = Stream.catchSomeCause(s1, (cause) => {
  if (Cause.isDie(cause)) {
    return Option.some(s2)
  }
  return Option.none()
})

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: "Chunk",
  values: [ 1, 2, 3, "a", "b", "c" ]
}
*/

恢复为一个 Effect

在流处理中,优雅地处理错误并在需要时执行清理任务非常关键。Stream.onError 函数正好能让我们做到这一点。如果我们的流遇到了错误,我们可以指定一个要执行的清理任务。

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

const stream = Stream.make(1, 2, 3).pipe(
  Stream.concat(Stream.dieMessage("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: RuntimeException: Oh! Boom!
*/

重试失败的流

有时,流可能遇到临时的、可恢复的失败。在这种情况下,Stream.retry 操作符就派上了用场。它允许你指定一个重试计划(schedule),流会按照该计划进行重试。

示例

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
{
  _id: "Chunk",
  values: [ 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)
        })
      }),
  )

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

细化错误

在处理流时,可能会出现这样的情况:你想有选择地保留某些错误,并以其余的错误终止流。你可以使用 Stream.refineOrDie 函数来实现这一点。

示例

import { Stream, Option } from "effect"

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

const res = Stream.refineOrDie(stream, (error) => {
  if (error instanceof SyntaxError) {
    return Option.some(error)
  }
  return Option.none()
})

在这个例子中,stream 最初以一个通用的 Error 失败。不过,我们使用 Stream.refineOrDie 来过滤并只保留类型为 SyntaxError 的错误。任何其他错误都会终止流,而 SyntaxError 会被保留在 refinedStream 中。

超时

在处理流时,可能会遇到需要处理超时的场景,比如当流在一定时长内没有产生值时将其终止。本节我们将探讨如何使用各种操作符来管理超时。

timeout

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

import { Stream, Effect } from "effect"

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

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
{
  _id: "Chunk",
  values: []
}
*/

timeoutFail

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

import { Stream, Effect } from "effect"

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

Effect.runPromiseExit(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: 'Exit',
  _tag: 'Failure',
  cause: { _id: 'Cause', _tag: 'Fail', failure: 'timeout' }
}
*/

timeoutFailCause

Stream.timeoutFail 类似,Stream.timeoutFailCause 把超时与自定义的失败 cause 结合在一起。如果流超时了,它就会以指定的 cause 失败。

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

const stream = Stream.fromEffect(Effect.never).pipe(
  Stream.timeoutFailCause(() => Cause.die("timeout"), "2 seconds"),
)

Effect.runPromiseExit(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: 'Exit',
  _tag: 'Failure',
  cause: { _id: 'Cause', _tag: 'Die', defect: 'timeout' }
}
*/

timeoutTo

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

import { Stream, Effect } from "effect"

const stream = Stream.fromEffect(Effect.never).pipe(
  Stream.timeoutTo("2 seconds", Stream.make(1, 2, 3)),
)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
{
  _id: "Chunk",
  values: [ 1, 2, 3 ]
}
*/