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

Sink 操作

探索用于变换、过滤和适配 Sink 的操作,从而在 Stream 处理中实现自定义的输入输出处理与元素过滤。

在前面几节中,我们学习了如何创建和使用 Sink。现在,让我们来探索一些可以变换或过滤 Sink 行为的操作。

适配 Sink 的输入

有时,你的 Sink 处理的是一种输入类型,而当前的 stream 使用的是另一种类型。Sink.mapInput 函数通过变换输入值,帮助你让 Sink 适配新的输入类型。Sink.map 改变的是 Sink 的输出,而 Sink.mapInput 改变的是它接受的输入。

示例(将字符串输入转换为数值以便求和)

假设你有一个用于计算数字之和的 Sink.sum。如果你的 stream 中包含的是字符串而不是数字,那么 Sink.mapInput 可以把这些字符串转换为数字,从而让 Sink.sum 能与你的 stream 配合工作:

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

// A stream of numeric strings
const stream = Stream.make("1", "2", "3", "4", "5")

// Define a sink for summing numeric values
const numericSum = Sink.sum

// Use mapInput to adapt the sink, converting strings to numbers
const stringSum = numericSum.pipe(
  Sink.mapInput((s: string) => Number.parseFloat(s)),
)

await Effect.runPromise(Stream.run(stream, stringSum)) // => 15

同时变换输入与输出

当你需要同时变换 Sink 的输入和输出时,可以把 Sink.mapInputSink.map 组合起来使用。二者配合,可以让你先变换输入类型、执行操作,再把输出变换为新类型。这在需要在输入类型和输出类型之间做完整转换时很有用。

示例(将输入转换为整数、求和,再把输出转换为字符串)

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

// A stream of numeric strings
const stream = Stream.make("1", "2", "3", "4", "5")

// Convert string inputs to numbers, sum them,
// then convert the result to a string
const sumSink = Sink.sum.pipe(
  // Transform input: string to number
  Sink.mapInput((s: string) => Number.parseFloat(s)),
  // Transform output: number to string
  Sink.map((n) => String(n)),
)

await Effect.runPromise(Stream.run(stream, sumSink)) // => "15"

过滤输入

你可以先过滤为 Sink 提供元素的 stream,再用 Stream.transduce 把它交给 Sink,从而过滤 Sink 最终要处理的元素。这样就能把过滤条件与 Sink.take 这类只关心满足特定条件的元素的 Sink 组合起来。

示例(按每三个一组过滤负数)

在下面的示例中,元素被收集为每三个一组的数组,但只有正数会被包含进来:

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

// Define a stream with positive, negative, and zero values
const stream = Stream.fromIterable([
  1, -2, 0, 1, 3, -3, 4, 2, 0, 1, -3, 1, 1, 6,
]).pipe(
  // Filter out non-positive numbers before grouping
  Stream.filter((n) => n > 0),
  // Collect the remaining elements in groups of 3
  Stream.transduce(Sink.take(3)),
)

await Effect.runPromise(Stream.runCollect(stream)) // => [[1, 1, 3], [4, 2, 1], [1, 1, 6], []]