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.catchFilter 和 Stream.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]