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

创建 Sink

了解如何创建和使用各种用于处理 Stream 的 Sink,包括计数、求和、收集、折叠,以及处理成功与失败。

在 Stream 处理中,Sink 用于消费和处理来自 stream 的元素。本节将探索各种 Sink 构造函数,它们让你可以为特定任务创建 Sink

常用构造函数

Sink.head 只取 stream 的第一个元素,并用 Some 包装它。如果 stream 没有任何元素,则返回 None

示例(获取第一个元素)

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

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

Effect.runPromise(Stream.run(nonEmptyStream, Sink.head())).then(console.log)
/*
Output:
{ _id: 'Option', _tag: 'Some', value: 1 }
*/

const emptyStream = Stream.empty

Effect.runPromise(Stream.run(emptyStream, Sink.head())).then(console.log)
/*
Output:
{ _id: 'Option', _tag: 'None' }
*/

last

Sink.last 只取 stream 的最后一个元素,并用 Some 包装它。如果 stream 没有任何元素,则返回 None

示例(获取最后一个元素)

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

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

Effect.runPromise(Stream.run(nonEmptyStream, Sink.last())).then(console.log)
/*
Output:
{ _id: 'Option', _tag: 'Some', value: 4 }
*/

const emptyStream = Stream.empty

Effect.runPromise(Stream.run(emptyStream, Sink.last())).then(console.log)
/*
Output:
{ _id: 'Option', _tag: 'None' }
*/

count

Sink.count 会消费 stream 的所有元素,并统计传给它的元素数量。

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

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

Effect.runPromise(Stream.run(stream, Sink.count)).then(console.log)
// Output: 4

sum

Sink.sum 会消费 stream 的所有元素,并对传入的数值求和。

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

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

Effect.runPromise(Stream.run(stream, Sink.sum)).then(console.log)
// Output: 10

take

Sink.take 会从 stream 中取出指定数量的值,结果是一个 Chunk 数据类型。

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

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

Effect.runPromise(Stream.run(stream, Sink.take(3))).then(console.log)
/*
Output:
{ _id: 'Chunk', values: [ 1, 2, 3 ] }
*/

drain

Sink.drain 会忽略它的输入,实际上就是把它们丢弃。

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

const stream = Stream.make(1, 2, 3, 4).pipe(Stream.tap(Console.log))

Effect.runPromise(Stream.run(stream, Sink.drain)).then(console.log)
/*
Output:
1
2
3
4
undefined
*/

timed

Sink.timed 会执行 stream 并测量其执行时间,返回一个 Duration

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

const stream = Stream.make(1, 2, 3, 4).pipe(
  Stream.schedule(Schedule.spaced("100 millis")),
)

Effect.runPromise(Stream.run(stream, Sink.timed)).then(console.log)
/*
Output:
{ _id: 'Duration', _tag: 'Millis', millis: 408 }
*/

forEach

Sink.forEach 会针对传给它的每个元素执行所提供的 effect 函数。

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

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

Effect.runPromise(Stream.run(stream, Sink.forEach(Console.log))).then(
  console.log,
)
/*
Output:
1
2
3
4
undefined
*/

从成功与失败创建 Sink

正如你可以定义 stream 来保存或操作数据,你也可以使用 Sink.failSink.succeed 函数创建具有特定成功或失败结果的 Sink

成功的 Sink

下面的示例创建了一个 Sink:它不消费上游源中的任何元素,而是立即以一个指定的数值成功结束:

示例(总是以某个值成功的 Sink)

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

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

Effect.runPromise(Stream.run(stream, Sink.succeed(0))).then(console.log)
// Output: 0

失败的 Sink

在这个示例中,这个 Sink 同样不消费上游源中的任何元素。相反,它以 string 类型的指定错误消息失败:

示例(总是以错误消息失败的 Sink)

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

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

Effect.runPromiseExit(Stream.run(stream, Sink.fail("fail!"))).then(console.log)
/*
Output:
{
  _id: 'Exit',
  _tag: 'Failure',
  cause: { _id: 'Cause', _tag: 'Fail', failure: 'fail!' }
}
*/

收集

收集所有元素

要把数据流中的所有元素汇总到一个 Chunk 中,可以使用 Sink.collectAll

最终输出是一个 Chunk,按元素被发出的顺序包含 stream 中的所有元素。

示例(收集 Stream 中的所有元素)

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

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

Effect.runPromise(Stream.run(stream, Sink.collectAll())).then(console.log)
/*
Output:
{ _id: 'Chunk', values: [ 1, 2, 3, 4 ] }
*/

收集指定数量

要把 stream 中固定数量的元素收集到一个 Chunk 中,可以使用 Sink.collectAllN。这个 Sink 在达到指定上限后就停止收集。

示例(收集有限数量的元素)

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

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

Effect.runPromise(
  Stream.run(
    stream,
    // Collect the first 3 elements into a Chunk
    Sink.collectAllN(3),
  ),
).then(console.log)
/*
Output:
{ _id: 'Chunk', values: [ 1, 2, 3 ] }
*/

在满足条件时收集

要在元素满足特定条件时从 stream 中收集它们,可以使用 Sink.collectAllWhile。这个 Sink 会持续收集元素,直到给定的谓词返回 false

示例(收集元素直到条件不再满足)

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

const stream = Stream.make(1, 2, 0, 4, 0, 6, 7)

Effect.runPromise(
  Stream.run(
    stream,
    // Collect elements while they are not equal to 0
    Sink.collectAllWhile((n) => n !== 0),
  ),
).then(console.log)
/*
Output:
{ _id: 'Chunk', values: [ 1, 2 ] }
*/

收集到 HashSet

要把 stream 的元素累积到一个 HashSet 中,可以使用 Sink.collectAllToSet()。这样可以确保每个元素在最终集合中只出现一次。

示例(把去重后的元素收集到 HashSet)

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

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

Effect.runPromise(Stream.run(stream, Sink.collectAllToSet())).then(console.log)
/*
Output:
{ _id: 'HashSet', values: [ 1, 2, 3, 4 ] }
*/

收集到指定大小的 HashSet

如果需要以受控方式把元素收集到有指定最大大小的 HashSet 中,可以使用 Sink.collectAllToSetN。这个 Sink 会收集去重后的元素,直到达到给定的上限。

示例(在限制集合大小的情况下收集去重元素)

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

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

Effect.runPromise(
  Stream.run(
    stream,
    // Collect unique elements, limiting the set size to 3
    Sink.collectAllToSetN(3),
  ),
).then(console.log)
/*
Output:
{ _id: 'HashSet', values: [ 1, 2, 3 ] }
*/

收集到 HashMap

对于更复杂的收集场景,Sink.collectAllToMap 让你可以把元素收集到一个 HashMap<K, A> 中,并指定 key 策略与合并策略。这个 Sink 既需要一个 key 函数来定义每个元素的分组,也需要一个合并函数来合并共享同一个 key 的值。

示例(在 HashMap 中分组并合并 Stream 元素)

在这个示例中,我们用 (n) => n % 3 确定 map 的 key,用 (a, b) => a + b 合并具有相同 key 的元素:

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

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

Effect.runPromise(
  Stream.run(
    stream,
    Sink.collectAllToMap(
      (n) => n % 3, // Key function to group by element value
      (a, b) => a + b, // Merge function to sum values with the same key
    ),
  ),
).then(console.log)
/*
Output:
{ _id: 'HashMap', values: [ [ 0, 6 ], [ 1, 3 ], [ 2, 7 ] ] }
*/

收集到 key 数量受限的 HashMap

要把元素累积到一个 key 数量有上限的 HashMap 中,可以使用 Sink.collectAllToMapN。这个 Sink 会一直收集元素,直到达到指定的 key 上限;它需要一个 key 函数来定义每个元素的分组,以及一个合并函数来合并具有相同 key 的值。

示例(限制 HashMap 中收集的 key 数量)

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

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

Effect.runPromise(
  Stream.run(
    stream,
    Sink.collectAllToMapN(
      3, // Maximum of 3 keys
      (n) => n, // Key function to group by element value
      (a, b) => a + b, // Merge function to sum values with the same key
    ),
  ),
).then(console.log)
/*
Output:
{ _id: 'HashMap', values: [ [ 1, 2 ], [ 2, 2 ], [ 3, 6 ] ] }
*/

折叠

左折叠

如果你想按顺序对每个元素应用一个操作,把 stream 归约为单个累积值,可以使用 Sink.foldLeft 函数。

示例(用 foldLeft 对 Stream 中的元素求和)

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

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

Effect.runPromise(
  Stream.run(
    stream,
    // Use foldLeft to sequentially add each element, starting with 0
    Sink.foldLeft(0, (a, b) => a + b),
  ),
).then(console.log)
// Output: 10

带终止条件的折叠

有时,你可能想折叠 stream 中的元素,但在满足某个特定条件时就停止这一过程。这被称为“短路”(short-circuiting)。你可以用 Sink.fold 函数做到这一点,它允许你定义终止条件。

示例(带提前停止条件的折叠)

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

const stream = Stream.iterate(0, (n) => n + 1)

Effect.runPromise(
  Stream.run(
    stream,
    Sink.fold(
      0, // Initial value
      (sum) => sum <= 10, // Termination condition
      (a, b) => a + b, // Folding operation
    ),
  ),
).then(console.log)
// Output: 15

折叠到某个上限

要累积元素直到达到特定数量,可以使用 Sink.foldUntil。这个 Sink 会一直折叠元素,直到达到指定上限,然后停止。

示例(累积固定数量的元素)

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

const stream = Stream.make(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)

Effect.runPromise(
  Stream.run(
    stream,
    // Fold elements, stopping after accumulating 3 values
    Sink.foldUntil(0, 3, (a, b) => a + b),
  ),
).then(console.log)
// Output: 6

带权重元素的折叠

在某些场景中,你可能希望按定义好的“权重”(weight)或“代价”(cost)来折叠元素,累积元素直到达到指定的最大代价。你可以用 Sink.foldWeighted 做到这一点。

示例(按权重累积元素)

在下面的示例中,每个元素的权重都是 1,当累积权重达到 3 时折叠就会重新开始。

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

const stream = Stream.make(3, 2, 4, 1, 5, 6, 2, 1, 3, 5, 6).pipe(
  Stream.transduce(
    Sink.foldWeighted({
      initial: Chunk.empty<number>(), // Initial empty Chunk
      maxCost: 3, // Maximum accumulated cost
      cost: () => 1, // Each element has a weight of 1
      body: (acc, el) => Chunk.append(acc, el), // Append element to the Chunk
    }),
  ),
)

Effect.runPromise(Stream.runCollect(stream)).then((chunk) =>
  console.log("%o", chunk),
)
/*
Output:
{
  _id: 'Chunk',
  values: [
    { _id: 'Chunk', values: [ 3, 2, 4, [length]: 3 ] },
    { _id: 'Chunk', values: [ 1, 5, 6, [length]: 3 ] },
    { _id: 'Chunk', values: [ 2, 1, 3, [length]: 3 ] },
    { _id: 'Chunk', values: [ 5, 6, [length]: 2 ] },
    [length]: 4
  ]
}
*/