剩余元素
学习如何处理 Stream 中未被消费的元素:收集或忽略剩余元素,从而实现高效的数据处理。
在本节中,我们将探讨如何处理 Sink 未消费、被留下的元素。Sink 可能只处理上游源中的一部分元素,而把其余元素留作「剩余元素」(leftovers)。下面介绍如何收集或忽略这些剩余元素。
收集剩余元素
如果 Sink 没有消费上游源中的所有元素,那么剩下的元素就称为剩余元素(leftovers)。若要捕获这些剩余元素,可以使用 Sink.collectLeftover,它返回一个元组,其中包含 Sink 操作的结果以及所有未消费的元素。
示例(收集剩余元素)
import { Stream, Sink, Effect } 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.collectLeftover)
Effect.runPromise(Stream.run(stream, sink1)).then(console.log)
/*
Output:
[
{ _id: 'Chunk', values: [ 1, 2, 3 ] },
{ _id: 'Chunk', values: [ 4, 5 ] }
]
*/
// Take only the first element and collect the rest as leftovers
const sink2 = Sink.head<number>().pipe(Sink.collectLeftover)
Effect.runPromise(Stream.run(stream, sink2)).then(console.log)
/*
Output:
[
{ _id: 'Option', _tag: 'Some', value: 1 },
{ _id: 'Chunk', values: [ 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.collectLeftover,
)
Effect.runPromise(Stream.run(stream, sink)).then(console.log)
/*
Output:
[ { _id: 'Chunk', values: [ 1, 2, 3 ] }, { _id: 'Chunk', values: [] } ]
*/