创建 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](包含 min 和 max 两个端点)内的整数 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 数据类型很相似,你可以使用 fail 和 succeed 函数生成 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.unfold 的 step 函数本身就返回一个 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.paginate 与 Stream.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())。
展开与分页的对比
你可能会好奇 unfold 与 paginate 这两个组合子有何区别,以及何时该用哪一个。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,并在 nextCursor 为 undefined 时停止。若用 Stream.unfold 来建模,就需要每次从当前页中剥出一个元素,并把剩余部分作为额外的状态保存起来。Stream.paginate 通过让每一步接收一个数组,已经处理好了这套逻辑。
从 Queue 和 PubSub 创建
在 Effect 中,有两种至关重要的异步消息数据类型:Queue 和 PubSub。你可以分别借助 Stream.fromQueue 和 Stream.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]