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

创建 Stream

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

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

常用构造函数

make

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

import { Stream, Effect } from "effect"

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

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

empty

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

import { Stream, Effect } from "effect"

const stream = Stream.empty

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [] }

void

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

import { Stream, Effect } from "effect"

const stream = Stream.void

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 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)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 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, ...

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
// { _id: 'Chunk', values: [ 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(
  Effect.acquireUseRelease(
    Console.log("acquire"),
    () => Console.log("use"),
    () => Console.log("release"),
  ),
)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
acquire
use
release
{ _id: 'Chunk', values: [ undefined ] }
*/

从成功与失败创建

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)
// { _id: 'Chunk', values: [ 5 ] }

从 Chunk 创建

你可以像这样从 Chunk 构造 stream:

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

// Creating a stream with values from a single Chunk
const stream = Stream.fromChunk(Chunk.make(1, 2, 3))

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

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

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

// Creating a stream with values from multiple Chunks
const stream = Stream.fromChunks(Chunk.make(1, 2, 3), Chunk.make(4, 5, 6))

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 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: { _id: 'Chunk', values: [ 1042302242 ] }

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

从异步回调创建

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

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

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

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

const stream = Stream.async(
  (emit: StreamEmit.Emit<never, never, number, void>) => {
    events.forEach((n) => {
      setTimeout(() => {
        if (n === 3) {
          // Terminate the stream
          emit(Effect.fail(Option.none()))
        } else {
          // Add the current item to the stream
          emit(Effect.succeed(Chunk.of(n)))
        }
      }, 100 * n)
    })
  },
)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2 ] }

StreamEmit.Emit<R, E, A, void> 类型表示一个可以被多次调用的异步回调。该回调接收一个类型为 Effect<Chunk<A>, Option<E>, R> 的值。各种可能的结果含义如下:

  • 当传给回调的值在成功时得到 Chunk<A>,表示应当把指定的元素作为 stream 的一部分发出。

  • 如果传给回调的值以 Some<E> 失败,表示以指定的错误终止 stream。

  • 当传给回调的值以 None 失败时,它充当 stream 结束的信号,本质上是终止这个 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)

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

fromIterableEffect

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

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

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

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

const getUsers = Database.pipe(Effect.andThen((_) => _.getUsers))

const stream = Stream.fromIterableEffect(getUsers)

Effect.runPromise(
  Stream.runCollect(
    stream.pipe(
      Stream.provideService(Database, {
        getUsers: Effect.succeed(["user1", "user2"]),
      }),
    ),
  ),
).then(console.log)
// { _id: 'Chunk', values: [ '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
)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 1, 2 ] }

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

从重复创建

重复单个值

你可以使用 Stream.repeatValue 构造函数创建一个无休止重复某个特定值的 stream:

import { Stream, Effect } from "effect"

const stream = Stream.repeatValue(0)

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
// { _id: 'Chunk', values: [ 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)

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
// { _id: 'Chunk', values: [ 1, 1, 1, 1, 1 ] }

重复 Effect 的结果

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

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

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

const stream = Stream.repeatEffect(Random.nextInt)

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

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

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

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

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

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

生成 tick

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

import { Stream, Effect } from "effect"

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

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
/*
Output:
{
  _id: 'Chunk',
  values: [ undefined, undefined, undefined, undefined, undefined ]
}
*/

从展开/分页创建

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

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

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

展开

unfold

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

declare const unfold: <S, A>(
  initialState: S,
  step: (s: S) => Option.Option<readonly [A, S]>,
) => Stream<A>

它的工作方式如下:

  • initialState。这是初始状态值。
  • step。状态函数 step 接收当前状态 s 作为输入。如果该函数的结果是 None,则 stream 结束。如果是 Some<[A, S]>,那么 stream 中的下一个元素就是 A,同时状态 S 会更新,供下一步处理使用。

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

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

const stream = Stream.unfold(1, (n) => Option.some([n, n + 1]))

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
// { _id: 'Chunk', values: [ 1, 2, 3, 4, 5 ] }

unfoldEffect

有时,我们可能需要在展开过程中执行带 effect 的状态变换。这正是 Stream.unfoldEffect 的用武之地,它让我们可以在生成 stream 的同时处理 effect。

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

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

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

Effect.runPromise(Stream.runCollect(stream.pipe(Stream.take(5)))).then(
  console.log,
)
// Example Output: { _id: 'Chunk', values: [ 1, 1, 1, 1, -1 ] }

其他变体

还有一些类似的操作,比如 Stream.unfoldChunkStream.unfoldChunkEffect,它们是专为处理 Chunk 数据类型而设计的。

分页

paginate

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

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

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

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

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
// { _id: 'Chunk', values: [ 0, 1, 2, 3 ] }

它的工作方式如下:

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

其他变体

还有一些类似的操作,比如 Stream.paginateChunkStream.paginateChunkEffect,它们是专为处理 Chunk 数据类型而设计的。

展开与分页的对比

你可能会好奇 unfoldpaginate 这两个组合子之间有什么区别,以及何时该用其中一个而不是另一个。让我们通过一个示例来深入探讨。

设想我们有一个分页 API,它以分页的方式提供大量数据。当我们向这个 API 发起请求时,它会返回一个 ResultPage 对象,其中包含当前页的结果,以及一个标志,用于指示它是否是最后一页、或者下一页是否还有更多数据要获取。下面是我们这个 API 的简化表示:

import { Chunk, Effect } from "effect"

type RawData = string

class PageResult {
  constructor(
    readonly results: Chunk.Chunk<RawData>,
    readonly isLast: boolean,
  ) {}
}

const pageSize = 2

const listPaginated = (
  pageNumber: number,
): Effect.Effect<PageResult, Error> => {
  return Effect.succeed(
    new PageResult(
      Chunk.map(
        Chunk.range(1, pageSize),
        (index) => `Result ${pageNumber}-${index}`,
      ),
      pageNumber === 2, // Return 3 pages
    ),
  )
}

我们的目标是把这样一个分页 API 转换成一个由 RowData 事件组成的 Stream。在最初的尝试中,我们可能会认为使用 Stream.unfold 操作就是可行之道:

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

type RawData = string

class PageResult {
  constructor(
    readonly results: Chunk.Chunk<RawData>,
    readonly isLast: boolean,
  ) {}
}

const pageSize = 2

const listPaginated = (
  pageNumber: number,
): Effect.Effect<PageResult, Error> => {
  return Effect.succeed(
    new PageResult(
      Chunk.map(
        Chunk.range(1, pageSize),
        (index) => `Result ${pageNumber}-${index}`,
      ),
      pageNumber === 2, // Return 3 pages
    ),
  )
}

const firstAttempt = Stream.unfoldChunkEffect(0, (pageNumber) =>
  listPaginated(pageNumber).pipe(
    Effect.map((page) => {
      if (page.isLast) {
        return Option.none()
      }
      return Option.some([page.results, pageNumber + 1] as const)
    }),
  ),
)

Effect.runPromise(Stream.runCollect(firstAttempt)).then(console.log)
/*
Output:
{
  _id: "Chunk",
  values: [ "Result 0-1", "Result 0-2", "Result 1-1", "Result 1-2" ]
}
*/

然而,这种做法有一个缺点:它没有包含最后一页的结果。为了绕开这个问题,我们额外发起一次 API 调用,把这些缺失的结果也包含进来:

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

type RawData = string

class PageResult {
  constructor(
    readonly results: Chunk.Chunk<RawData>,
    readonly isLast: boolean,
  ) {}
}

const pageSize = 2

const listPaginated = (
  pageNumber: number,
): Effect.Effect<PageResult, Error> => {
  return Effect.succeed(
    new PageResult(
      Chunk.map(
        Chunk.range(1, pageSize),
        (index) => `Result ${pageNumber}-${index}`,
      ),
      pageNumber === 2, // Return 3 pages
    ),
  )
}

const secondAttempt = Stream.unfoldChunkEffect(Option.some(0), (pageNumber) =>
  Option.match(pageNumber, {
    // We already hit the last page
    onNone: () => Effect.succeed(Option.none()),
    // We did not hit the last page yet
    onSome: (pageNumber) =>
      listPaginated(pageNumber).pipe(
        Effect.map((page) =>
          Option.some([
            page.results,
            page.isLast ? Option.none() : Option.some(pageNumber + 1),
          ]),
        ),
      ),
  }),
)

Effect.runPromise(Stream.runCollect(secondAttempt)).then(console.log)
/*
Output:
{
  _id: 'Chunk',
  values: [
    'Result 0-1',
    'Result 0-2',
    'Result 1-1',
    'Result 1-2',
    'Result 2-1',
    'Result 2-2'
  ]
}
*/

虽然这种做法可行,但显然 Stream.unfold 并不是从分页 API 获取数据的最友好选择。它需要额外的变通手段,才能把最后一页的结果包含进来。

这正是 Stream.paginate 大显身手的地方。它提供了一种更符合人体工程学的方式,把分页 API 转换成 Effect Stream。让我们用 Stream.paginate 重写这个方案:

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

type RawData = string

class PageResult {
  constructor(
    readonly results: Chunk.Chunk<RawData>,
    readonly isLast: boolean,
  ) {}
}

const pageSize = 2

const listPaginated = (
  pageNumber: number,
): Effect.Effect<PageResult, Error> => {
  return Effect.succeed(
    new PageResult(
      Chunk.map(
        Chunk.range(1, pageSize),
        (index) => `Result ${pageNumber}-${index}`,
      ),
      pageNumber === 2, // Return 3 pages
    ),
  )
}

const finalAttempt = Stream.paginateChunkEffect(0, (pageNumber) =>
  listPaginated(pageNumber).pipe(
    Effect.andThen((page) => {
      return [
        page.results,
        page.isLast ? Option.none<number>() : Option.some(pageNumber + 1),
      ]
    }),
  ),
)

Effect.runPromise(Stream.runCollect(finalAttempt)).then(console.log)
/*
Output:
{
  _id: 'Chunk',
  values: [
    'Result 0-1',
    'Result 0-2',
    'Result 1-1',
    'Result 1-2',
    'Result 2-1',
    'Result 2-2'
  ]
}
*/

从 Queue 和 PubSub 创建

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

从 Schedule 创建

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

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

// Emits values every 1 second for a total of 10 emissions
const schedule = Schedule.spaced("1 second").pipe(
  Schedule.compose(Schedule.recurs(10)),
)

const stream = Stream.fromSchedule(schedule)

Effect.runPromise(Stream.runCollect(stream)).then(console.log)
/*
Output:
{
  _id: 'Chunk',
  values: [
    0, 1, 2, 3, 4,
    5, 6, 7, 8, 9
  ]
}
*/