消费 Stream
学习消费 Stream 的各种技巧,包括收集元素、使用回调处理,以及使用 fold 与 Sink。
使用 Stream 时,理解如何消费它们产出的数据至关重要。在本指南中,我们将逐一介绍几种常见的 Stream 消费方法。
使用 runCollect
要把 Stream 中的所有元素收集到一个数组里,可以使用 Stream.runCollect 函数。
import { Stream, Effect } from "effect"
const stream = Stream.make(1, 2, 3, 4, 5)
const collectedData = Stream.runCollect(stream)
await Effect.runPromise(collectedData) // => [1, 2, 3, 4, 5]
使用 runForEach
另一种消费 Stream 元素的方式是使用 Stream.runForEach。它接收一个回调函数,该函数会接收 Stream 的每个元素。示例如下:
import { Stream, Effect, Console } from "effect"
const effect = Stream.make(1, 2, 3).pipe(
Stream.runForEach((n) => Console.log(n)),
)
await Effect.runPromise(effect) // => undefined
在这个示例中,我们使用 Stream.runForEach 把每个元素打印到控制台。
使用 runFold
Stream.runFold 通过归约 Stream 的值来消费它,并返回一个包含结果的 Effect。若需要提前终止,可以使用 Stream.runForEachWhile,并把累加器保持为 Effect.suspend 块内的局部变量。
import { Stream, Effect } from "effect"
const foldedStream = Stream.make(1, 2, 3, 4, 5).pipe(
Stream.runFold(
() => 0,
(a, b) => a + b,
),
)
await Effect.runPromise(foldedStream) // => 15
const foldedWhileStream = Effect.suspend(() => {
let acc = 0
return Stream.make(1, 2, 3, 4, 5)
.pipe(
Stream.runForEachWhile((n) => {
acc = acc + n
return Effect.succeed(acc <= 3)
}),
)
.pipe(Effect.map(() => acc))
})
await Effect.runPromise(foldedWhileStream) // => 6
在第一个示例中,Stream.runFold 计算所有元素之和。在第二个示例中,Stream.runForEachWhile 在累加器超过 3 后停止;使谓词返回 false 的那个元素已经被消费,因此结果是 6。
使用 Sink
要使用 Sink 消费 Stream,可以把 Sink 传给 Stream.run 函数。示例如下:
import { Stream, Sink, Effect } from "effect"
const effect = Stream.make(1, 2, 3).pipe(Stream.run(Sink.sum))
await Effect.runPromise(effect) // => 6
在这个示例中,我们使用 Sink 计算 Stream 中所有元素之和。