创建 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](包含 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)
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 数据类型很相似,你可以使用 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)
// { _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.unfoldChunk 与 Stream.unfoldChunkEffect,它们是专为处理 Chunk 数据类型而设计的。
分页
paginate
Stream.paginate 与 Stream.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.paginateChunk 与 Stream.paginateChunkEffect,它们是专为处理 Chunk 数据类型而设计的。
展开与分页的对比
你可能会好奇 unfold 与 paginate 这两个组合子之间有什么区别,以及何时该用其中一个而不是另一个。让我们通过一个示例来深入探讨。
设想我们有一个分页 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 中,有两种至关重要的异步消息数据类型:Queue 和 PubSub。你可以分别借助 Stream.fromQueue 和 Stream.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
]
}
*/