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

创建 Stream

学习创建 Effect stream 的各种方法,涵盖从基础构造函数到异步数据源、分页与调度的处理。

在本节中,我们将探讨创建 Effect Stream 的各种方法。这些方法能帮助你生成契合自身需求的 stream。

常用构造函数

make

你可以使用 Stream.make 构造函数创建一个纯 stream。该构造函数接受一组数量可变的值作为参数。

import { Stream, Effect } from "effect"

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

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

empty

有时你可能需要一个不产生任何值的 stream。这种情况下,可以使用 Stream.empty。这个构造函数创建的 stream 始终保持为空。

import { Stream, Effect } from "effect"

const stream = Stream.empty

await Effect.runPromise(Stream.runCollect(stream)) // => []

void

如果你需要一个只包含单个 void 值的 stream,可以使用 Stream.succeed(void 0)。当你想用一个 stream 表示单个事件或信号时,这很方便。

import { Stream, Effect } from "effect"

const stream = Stream.succeed(void 0)

await Effect.runPromise(Stream.runCollect(stream)) // => [undefined]

range

要创建指定范围 [min, max](包含 minmax 两个端点)内的整数 stream,可以使用 Stream.range。这在生成连续数字的 stream 时特别有用。

import { Stream, Effect } from "effect"

// Creating a stream of numbers from 1 to 5
const stream = Stream.range(1, 5)

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

iterate

使用 Stream.iterate,你可以通过对初始值反复应用一个函数来生成 stream。初始值会成为 stream 产生的第一个元素,随后依次是由 f(init)f(f(init)) 等产生的值。

import { Stream, Effect } from "effect"

// Creating a stream of incrementing numbers
const stream = Stream.iterate(1, (n) => n + 1) // Produces 1, 2, 3, ...

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

scoped

Stream.scoped 用于从作用域资源创建一个只含单个值的 stream。当处理需要显式获取、使用与释放的资源时,它会很有用。

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

// Creating a single-valued stream from a scoped resource
const stream = Stream.scoped(
  Stream.fromEffect(
    Effect.acquireUseRelease(
      Console.log("acquire"),
      () => Console.log("use"),
      () => Console.log("release"),
    ),
  ),
)

await Effect.runPromise(Stream.runCollect(stream)) // => [undefined]
/*
Output:
acquire
use
release
*/

从成功与失败创建

Effect 数据类型很相似,你可以使用 failsucceed 函数生成 Stream

import { Stream, Effect } from "effect"

// Creating a stream that can emit errors
const streamWithError: Stream.Stream<never, string> = Stream.fail("Uh oh!")

Effect.runPromise(Stream.runCollect(streamWithError))
// throws Error: Uh oh!

// Creating a stream that emits a numeric value
const streamWithNumber: Stream.Stream<number> = Stream.succeed(5)

Effect.runPromise(Stream.runCollect(streamWithNumber)).then(console.log)
// [ 5 ]

从数组创建

你可以像这样从数组构造 stream:

import { Stream, Effect } from "effect"

// Creating a stream with values from a single array
const stream = Stream.fromArray([1, 2, 3])

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

此外,你也可以从多个数组创建 stream:

import { Stream, Effect } from "effect"

// Creating a stream with values from multiple arrays
const stream = Stream.fromArrays([1, 2, 3], [4, 5, 6])

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

从 Effect 创建

你可以使用 Stream.fromEffect 构造函数从 Effect 工作流生成 stream。例如下面这个 stream,它生成一个随机数:

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

const stream = Stream.fromEffect(Random.nextInt)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// Example Output: [ 1042302242 ]

// The value is random, but the stream always emits exactly one element
const result = await Effect.runPromise(Stream.runCollect(stream))
result.length // => 1

这个方法让你能够无缝地把 Effect 的输出转换为 stream,为在 stream 中处理异步操作提供了一种直接的方式。

从异步回调创建

假设你有一个依赖回调的异步函数。如果你想把这些回调发出的结果捕获为一个 stream,可以使用 Stream.callback 函数。这个函数专门用于适配那些会多次调用自身回调的函数,并把结果以 stream 的形式发出。

下面通过一个例子来拆解它的用法:

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

const events = [1, 2, 3, 4]

const stream = Stream.callback<number>((queue) =>
  Effect.sync(() => {
    events.forEach((n) => {
      setTimeout(() => {
        if (n === 3) {
          // Terminate the stream
          Queue.endUnsafe(queue)
        } else {
          // Add the current item to the stream
          Queue.offerUnsafe(queue, n)
        }
      }, 100 * n)
    })
  }),
)

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

传给 Stream.callback 的函数会接收一个 Queue,你可以在异步代码中使用它来驱动这个 stream。各种可能的操作含义如下:

  • 在队列上调用 Queue.offerUnsafe(或基于 effect 的 Queue.offer)会把给定的元素作为 stream 的一部分发出。

  • 在队列上调用 Queue.fail(或 Queue.failCauseUnsafe/Queue.failCause)会以指定的错误终止 stream。

  • 在队列上调用 Queue.endUnsafe/Queue.end 会发出 stream 结束的信号,从而成功地终止它。

简单来说,这让你完全掌控异步回调与 stream 的交互方式:决定何时发出元素、何时以错误终止,以及何时发出 stream 结束的信号。

从 Iterable 创建

fromIterable

你可以使用 Stream.fromIterable 构造函数从值的 Iterable 创建一个纯 stream。这是把一组值转换为 stream 的直接方式。

import { Stream, Effect } from "effect"

const numbers = [1, 2, 3]

const stream = Stream.fromIterable(numbers)

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

fromIterableEffect

当你有一个产生 Iterable 类型值的 effect 时,可以使用 Stream.fromIterableEffect 构造函数从该 effect 生成 stream。

例如,假设你有一个获取用户列表的数据库操作。由于该操作涉及 effect,你可以利用 Stream.fromIterableEffect 把结果转换为 Stream

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

class Database extends Context.Service<
  Database,
  { readonly getUsers: Effect.Effect<Array<string>> }
>()("Database") {}

const getUsers = Database.use((_) => _.getUsers)

const stream = Stream.fromIterableEffect(getUsers)

await Effect.runPromise(
  Stream.runCollect(
    stream.pipe(
      Stream.provideService(Database, {
        getUsers: Effect.succeed(["user1", "user2"]),
      }),
    ),
  ),
) // => ["user1", "user2"]

这让你能够无缝地处理 effect,并把它们的结果转换为 stream 以便进一步处理。

fromAsyncIterable

异步可迭代对象(async iterable)是另一类可以转换为 stream 的数据源。借助 Stream.fromAsyncIterable 构造函数,你可以处理异步数据源并优雅地处理潜在错误。

import { Stream, Effect } from "effect"

const myAsyncIterable = async function* () {
  yield 1
  yield 2
}

const stream = Stream.fromAsyncIterable(
  myAsyncIterable(),
  (e) => new Error(String(e)), // Error Handling
)

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

在这段代码中,我们定义了一个异步可迭代对象,然后由它创建了一个名为 stream 的 stream。此外,我们还提供了一个错误处理函数,用来管理转换过程中可能出现的任何错误。

从重复创建

重复单个值

你可以使用 Stream.forever(Stream.succeed(value)) 创建一个无休止重复某个特定值的 stream:

import { Stream, Effect } from "effect"

const stream = Stream.forever(Stream.succeed(0))

await Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))) // => [0, 0, 0, 0, 0]

重复 Stream 的内容

Stream.repeat 让你可以按照给定的调度重复指定 stream 的内容。这在生成周期性的事件或值时很有用。

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

// Creating a stream that repeats a value indefinitely
const stream = Stream.repeat(Stream.succeed(1), Schedule.forever)

await Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))) // => [1, 1, 1, 1, 1]

重复 Effect 的结果

假设你有一个 effectful 的 API 调用,并且想用该调用的结果来创建 stream。你可以通过从该 effect 创建 stream 并无限重复它来实现。

下面是一个生成随机数 stream 的例子:

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

const stream = Stream.fromEffectRepeat(Random.nextInt)

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
/*
Example Output:
[ 1666935266, 604851965, 2194299958, 3393707011, 4090317618 ]
*/

// The values are random, but the stream always emits exactly 5 elements
const result = await Effect.runPromise(
  Stream.runCollect(stream.pipe(Stream.take(5))),
)
result.length // => 5

重复 Effect 并在特定条件下终止

你可以重复求值一个给定的 effect,并根据特定条件终止 stream。

在这个例子中,我们通过耗尽(drain)一个 Iterator 来由它创建 stream:

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

const drainIterator = <A>(it: Iterator<A>): Stream.Stream<A> =>
  Stream.fromEffectRepeat(
    Effect.sync(() => it.next()).pipe(
      Effect.andThen((res) => {
        if (res.done) {
          return Cause.done()
        }
        return Effect.succeed(res.value)
      }),
    ),
  )

const numbers = [10, 20, 30]

await Effect.runPromise(
  Stream.runCollect(drainIterator(numbers[Symbol.iterator]())),
) // => [10, 20, 30]

生成 tick

你可以使用 Stream.tick 构造函数创建一个按指定间隔发出 void 值的 stream。这对创建周期性事件很有用。

import { Stream, Effect } from "effect"

const stream = Stream.tick("100 millis")

await Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))) // => [undefined, undefined, undefined, undefined, undefined]

从展开/分页创建

在函数式编程中,unfold 这一概念可以看作 fold 的对偶。

使用 fold 时,我们处理一个数据结构并产出一个返回值。例如,我们可以接收一个 Array<number> 并计算其所有元素之和。

另一方面,unfold 表示这样的一种操作:从一个初始值开始,使用指定的状态函数一次添加一个元素,从而生成一个递归的数据结构。例如,我们可以从 1 开始、以 increment 函数作为状态函数,创建一段自然数序列。

展开

unfold

Stream 模块包含一个 unfold 函数,其定义如下:

declare const unfold: <S, A, E, R>(
  initialState: S,
  step: (s: S) => Effect.Effect<readonly [A, S] | undefined, E, R>,
) => Stream<A, E, R>

它的工作方式如下:

  • initialState。这是初始状态值。
  • step。状态函数 step 接收当前状态 s 作为输入,并返回一个 effect。如果该 effect 解析为 undefined,则 Stream 结束。如果它解析为一个元组 [A, S],那么 Stream 中的下一个元素就是 A,同时状态 S 会更新,供下一步处理使用。

例如,让我们用 Stream.unfold 创建一个自然数 Stream:

import { Stream, Effect } from "effect"

const stream = Stream.unfold(1, (n) => Effect.succeed([n, n + 1] as const))

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

带 Effect 的展开

有时,我们可能需要在展开过程中执行带 effect 的状态变换。由于传给 Stream.unfoldstep 函数本身就返回一个 Effect,它在产出下一个元素和状态时,可以依赖任意带 effect 的计算,比如生成一个随机值。

下面是一个使用 Stream.unfold 创建由随机 1-1 组成的无限 Stream 的示例:

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

const stream = Stream.unfold(1, (n) =>
  Random.nextBoolean.pipe(
    Effect.map((b) => (b ? ([n, -n] as const) : ([n, n] as const))),
  ),
)

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
// Example Output: [ 1, 1, 1, 1, -1 ]

// The sign is random, but the state starts at 1 and only ever flips sign,
// so its absolute value is deterministically always 1
const result = await Effect.runPromise(
  Stream.runCollect(stream.pipe(Stream.take(5))),
)
result.map(Math.abs) // => [1, 1, 1, 1, 1]

分页

paginate

Stream.paginateStream.unfold 类似,但允许一步发出更多的值。

例如,下面的 Stream 会发出 0, 1, 2, 3 这些元素:

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

const stream = Stream.paginate(0, (n) =>
  Effect.succeed([[n], n < 3 ? Option.some(n + 1) : Option.none()] as const),
)

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

它的工作方式如下:

  • 我们从一个初始值 0 开始。
  • 传入的函数接收当前值 n 并返回一个元组。元组的第一个元素是要发出的值(n),第二个元素决定是继续(Option.some(n + 1))还是停止(Option.none())。

展开与分页的对比

你可能会好奇 unfoldpaginate 这两个组合子有何区别,以及何时该用哪一个。Stream.unfold 每一步恰好产出一个值,因此它无法消费那种每一步天然会一次返回一批值的 API。Stream.paginate 正是为这种形态而设计的:每一步返回一个值数组以及下一个状态,因此单次调用可以在决定是否继续之前发出零个、一个或多个元素。

这正是一个分页 API 的形态。设想一个 fetchUsers 请求:给定一个游标,它返回一页用户,如果还有更多数据,还会返回下一页的游标:

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

interface Page<A> {
  readonly items: ReadonlyArray<A>
  readonly nextCursor: number | undefined
}

// A mock paginated API: six users, two per page.
const fetchUsers = (cursor: number): Effect.Effect<Page<string>> => {
  const users = ["Alice", "Bob", "Carol", "Dave", "Erin", "Frank"]
  const pageSize = 2
  const items = users.slice(cursor, cursor + pageSize)
  const nextCursor =
    cursor + pageSize < users.length ? cursor + pageSize : undefined
  return Effect.succeed({ items, nextCursor })
}

const stream = Stream.paginate(0, (cursor) =>
  fetchUsers(cursor).pipe(
    Effect.map(
      (page) => [page.items, Option.fromUndefinedOr(page.nextCursor)] as const,
    ),
  ),
)

await Effect.runPromise(Stream.runCollect(stream)) // => ["Alice", "Bob", "Carol", "Dave", "Erin", "Frank"]

每次调用 fetchUsers 都会返回一整页元素,而 Stream.paginate 会把每一页展平进结果 Stream,并在 nextCursorundefined 时停止。若用 Stream.unfold 来建模,就需要每次从当前页中剥出一个元素,并把剩余部分作为额外的状态保存起来。Stream.paginate 通过让每一步接收一个数组,已经处理好了这套逻辑。

从 Queue 和 PubSub 创建

在 Effect 中,有两种至关重要的异步消息数据类型:QueuePubSub。你可以分别借助 Stream.fromQueueStream.fromPubSub,轻松地把这些数据类型转换为 Stream

从 Schedule 创建

我们可以从一个不需要任何额外输入的 Schedule 创建 Stream。该 Stream 会为 Schedule 输出的每个值发出一个元素,只要 Schedule 继续,它就会一直继续:

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

// Emits values every 100 milliseconds for a total of 10 emissions
const schedule = Schedule.spaced("100 millis").pipe(
  Schedule.upTo({ times: 10 }),
)

const stream = Stream.fromSchedule(schedule)

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