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

剩余元素

学习如何处理 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: [] } ]
*/