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

消费 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 中所有元素之和。