流中的错误处理
学习如何在流中处理错误,确保稳健的恢复、重试与优雅的错误管理,从而实现可靠的流式处理。
从失败中恢复
处理可能出错的流时,知道如何优雅地应对这些错误至关重要。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.catchSome 和 Stream.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 ]
}
*/