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

剩余元素

学习如何处理 Stream 中未被消费的元素:收集或忽略剩余元素,从而实现高效的数据处理。

在本节中,我们将探讨如何处理 Sink 未消费的元素。Sink 可能只处理上游源中的一部分元素,而把其余元素留作「剩余元素」(leftovers)。下面介绍如何收集或忽略这些剩余元素。

收集剩余元素

如果 Sink 没有消费上游源中的所有元素,那么剩下的元素就称为剩余元素(leftovers)。Sink.mapEnd 会同时变换 Sink 的结果和它可选的剩余元素,因此你可以把剩余元素移到 Stream.run 返回的结果中。

示例(收集剩余元素)

import { Stream, Sink, Effect, Option } from "effect"

const stream = Stream.make(1, 2, 3, 4, 5)

// Take the first 3 elements and collect any leftovers
const sink1 = Sink.take<number>(3).pipe(
  Sink.mapEnd(([a, leftover]) => [[a, leftover ?? []] as const]),
)

await Effect.runPromise(Stream.run(stream, sink1)) // => [[1, 2, 3], [4, 5]]

// Take only the first element and collect the rest as leftovers
const sink2 = Sink.head<number>().pipe(
  Sink.mapEnd(([a, leftover]) => [[a, leftover ?? []] as const]),
)

await Effect.runPromise(Stream.run(stream, sink2)) // => [Option.some(1), [2, 3, 4, 5]]

忽略剩余元素

如果不需要这些剩余元素,可以用 Sink.ignoreLeftover 忽略它们。这种做法会丢弃所有未消费的元素,让 Sink 操作只关注它需要的元素。

示例(忽略剩余元素)

import { Stream, Sink, Effect } from "effect"

const stream = Stream.make(1, 2, 3, 4, 5)

// Take the first 3 elements and ignore any remaining elements
const sink = Sink.take<number>(3).pipe(
  Sink.ignoreLeftover,
  Sink.mapEnd(([a, leftover]) => [[a, leftover ?? []] as const]),
)

await Effect.runPromise(Stream.run(stream, sink)) // => [[1, 2, 3], []]