剩余元素
学习如何处理 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], []]