资源管理型 Stream
学习如何在 Stream 中管理资源:安全地获取与释放、用于清理任务的终结处理(finalization),以及确保在终结之后执行清理动作,从而在流式应用中稳健地处理资源。
在 Stream 模块中,你会发现大多数构造器都提供了一个特殊变体,用于把作用域资源(scoped resource)提升到 Stream 中。使用这些特定的构造器时,你本质上是在创建对资源管理天然安全的 stream。这些构造器会在创建 stream 之前完成资源的获取,并在 stream 使用完毕后确保它被正确关闭。
Stream 还提供了 Stream.acquireRelease 与 Stream.finalizer 构造器,它们与 Effect.acquireRelease 和 Effect.addFinalizer 有相似之处。这些工具让我们能够在 stream 结束运行之前执行清理或终结处理任务。
获取与释放
本节通过一个示例演示在文件操作中使用 Stream.acquireRelease。
import { Stream, Console, Effect } from "effect"
// Simulating File operations
const open = (filename: string) =>
Effect.gen(function* () {
yield* Console.log(`Opening ${filename}`)
return {
getLines: Effect.succeed(["Line 1", "Line 2", "Line 3"]),
close: Console.log(`Closing ${filename}`),
}
})
const stream = Stream.acquireRelease(
open("file.txt"),
(file) => file.close,
).pipe(Stream.flatMap((file) => file.getLines))
Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
Opening file.txt
Closing file.txt
{
_id: "Chunk",
values: [
[ "Line 1", "Line 2", "Line 3" ]
]
}
*/
在这段代码中,我们用 open 函数模拟文件操作。Stream.acquireRelease 用于确保文件被正确打开与关闭,随后我们使用所获取的资源处理文件中的各行。
终结处理
本节将探讨 stream 中的终结处理(finalization)概念。终结处理让我们能在 stream 结束之前执行某个特定动作。当我们想执行清理任务,或对 stream 做最后的收尾时,它会特别有用。
设想这样一个场景:我们的流式应用需要在执行完成时清理一个临时目录。这可以用 Stream.finalizer 函数来实现:
import { Stream, Console, Effect } from "effect"
const application = Stream.fromEffect(Console.log("Application Logic."))
const deleteDir = (dir: string) => Console.log(`Deleting dir: ${dir}`)
const program = application.pipe(
Stream.concat(
Stream.finalizer(
deleteDir("tmp").pipe(
Effect.andThen(Console.log("Temporary directory was deleted.")),
),
),
),
)
Effect.runPromise(Stream.runCollect(program)).then(console.log)
/*
Output:
Application Logic.
Deleting dir: tmp
Temporary directory was deleted.
{
_id: "Chunk",
values: [ undefined, undefined ]
}
*/
在这个代码示例中,我们首先用 application stream 表示应用逻辑。接着用 Stream.finalizer 定义一个终结处理步骤,它会删除临时目录并输出一条消息。这样就能确保应用执行完毕时临时目录被妥善清理。
确保执行
本节将探讨这样一个场景:我们需要在 stream 终结处理之后执行一些动作。为此,我们可以使用 Stream.ensuring 操作符。
设想这样一种情况:应用已经完成主要逻辑,并终结处理了一些资源,但之后我们还需要执行额外的动作。为此可以使用 Stream.ensuring:
import { Stream, Console, Effect } from "effect"
const program = Stream.fromEffect(Console.log("Application Logic.")).pipe(
Stream.concat(Stream.finalizer(Console.log("Finalizing the stream"))),
Stream.ensuring(
Console.log("Doing some other works after stream's finalization"),
),
)
Effect.runPromise(Stream.runCollect(program)).then(console.log)
/*
Output:
Application Logic.
Finalizing the stream
Doing some other works after stream's finalization
{
_id: "Chunk",
values: [ undefined, undefined ]
}
*/
在这个代码示例中,我们首先用 Application Logic. 这条消息表示应用逻辑。接着用 Stream.finalizer 指定终结处理步骤,它会输出 Finalizing the stream。之后,我们用 Stream.ensuring 表明希望在 stream 终结处理之后再执行一些额外任务,从而产生 Performing additional tasks after stream's finalization 这条消息。这样就能确保终结处理之后的动作按预期执行。