Effect Streams 实战指南:用 t3code 仓库中的 Effect 构建拉取式数据流管道

发布时间:2026/9/15 13:32:00
Effect Streams 实战指南:用 t3code 仓库中的 Effect 构建拉取式数据流管道
Effect Streams 实战指南用 t3code 仓库中的 Effect 构建拉取式数据流管道【免费下载链接】t3code项目地址: https://gitcode.com/GitHub_Trending/t3/t3code导读Effect Streams 是 Effect 库中用于表达随时间产生值的有效果effectful、拉取式pull-based序列的核心抽象它让你可以用统一的模型处理有限或无限的数据源并在其上组合出类型安全、可并发、可容错的流式管道。本指南以 t3code 仓库随附的 Effect 官方 AI 文档.repos/effect-smol/ai-docs/src/03_stream/index.md为骨架结合其同目录下的三个完整示例10_creating-streams.ts、20_consuming-streams.ts、30_encoding.ts并对照仓库真实源码系统讲解流的创建、转换、消费与编解码。读完本文你将掌握从任意数据源构造 Stream、用算子编排数据管道、以及基于 NDJSON/Msgpack 做流式编解码的完整实战能力。说明t3code 仓库将 Effect 官方文档与其示例代码一并归档在.repos/effect-smol/ai-docs/src/下其中03_stream目录即本节内容来源index.md只提供章节导语完整的技术细节全部沉淀在同目录的.ts示例文件中。一、理解 Effect Stream有效果的拉取式值序列1.1 核心定义根据 章节导语Effect Stream 的定义可以拆解为三个关键词effectful有效果的流中的每个元素都可能在产生时执行副作用I/O、网络请求、计时器等因此流的类型签名是StreamA, E, R——即产生A类型元素、可能以E类型失败、需要R环境依赖pull-based拉取式流是惰性的元素只在被消费时才按需产生拉取而不是预先一次性生成推送给消费者这一点与生成器generator的按需求值模型一致over time随时间元素可以按时间维度持续到来因此流既能建模有限数据源如数组、文件也能建模无限数据源如事件监听器、轮询采样、实时日志。这一拉取 有效果 时序的组合使 Stream 非常适合做指标采样、健康检查、分页 API 拉取、DOM/Node 事件处理、流式日志编解码等场景。1.2 在 t3code 中的实际印证Stream 不是纸上谈兵的抽象t3code 的真实业务代码大量使用了它。例如在 客户端连接层 中通过Stream.runForEach(registry.reconcilePlatform)把平台对账逻辑作为流消费者执行在 连接注册表的测试 中可以看到Stream.filter(...).pipe(Stream.runHead)、Stream.fromEffect(SubscriptionRef.get(...))、Stream.runDrain等算子被用于断言流式行为。这说明掌握本节 API 后你可以直接阅读并修改这类真实代码。二、创建 Stream从常见数据源构造流示例文件.repos/effect-smol/ai-docs/src/03_stream/10_creating-streams.ts展示了从 7 类常见数据源创建流的方法。2.1Stream.fromIterable从数组/可迭代对象创建任何可迭代对象数组、Set、Map、生成器等都可以直接转成流这是最基础的构造方式import { Stream } from effect export const numbers Stream.fromIterablenumber([1, 2, 3, 4, 5])2.2Stream.fromEffectSchedule把单个 Effect 变成轮询采样流把执行一次 Effect和调度器Schedule组合即可生成周期性执行的轮询流。典型的用途是指标采集、健康检查、缓存刷新循环import { Effect, Schedule, Stream } from effect export const samples Stream.fromEffectSchedule( Effect.succeed(3), // 每次轮询执行的 Effect Schedule.spaced(30 seconds) // 每 30 秒触发一次 ).pipe( // Stream.take 限制流最多发射的元素个数 Stream.take(3) )要点Schedule.spaced(30 seconds)也可写作Schedule.spaced(Duration.seconds(30))Stream.take在这里把理论上无限的轮询流截断为有限流是处理无限源时最常用的算子之一。2.3Stream.paginate游标式分页拉取读取一次返回一页的 API如 REST 列表、GitHub 搜索等时Stream.paginate是最佳选择。它接收一个初始游标和一个返回[当页结果, 可选下一个游标]的函数当函数返回Option.none()时流结束import { Array, Effect, Stream } from effect import * as Option from effect/Option export const fetchJobsPage Stream.paginate( 0, // 从第 0 页开始即游标 Effect.fn(function*(page) { // 模拟网络延迟 yield* Effect.sleep(50 millis) const results Array.range(0, 100).map((i) Job ${i 1 page * 100}) // 最多返回 10 页结果 const nextPage page 10 ? Option.some(page 1) : Option.none() return [results, nextPage] as const }) )注意这里使用了Effect.fn(function*(...))生成器语法编写带延迟的有效果分页函数返回元组[values, nextCursor]游标类型完全由你决定数字、字符串、哈希值皆可。2.4Stream.fromAsyncIterable从异步可迭代对象创建并做错误转型async generatorasync function*同样可以转为流。第二个参数是把异步可迭代对象抛出的任何错误转型为类型化错误的映射函数import { Schema, Stream } from effect class LetterError extends Schema.TaggedErrorLetterError()(LetterError, { cause: Schema.Defect() }) {} async function* asyncIterable() { yield a yield b yield c } export const letters Stream.fromAsyncIterable( asyncIterable(), (cause) new LetterError({ cause }) // 统一错误类型便于后续 catchTag )使用Schema.TaggedError定义类型化错误是 Effect 生态的标准做法后续可以直接用Stream.catchTag(LetterError, ...)精确捕获。2.5Stream.fromEventListener把 DOM 事件变为流浏览器端的按钮点击、键盘输入等 DOM 事件可以直接流式化事件对象即为流的元素const button document.getElementById(my-button)! export const events Stream.fromEventListenerPointerEvent(button, click)2.6Stream.callback任何回调式 API 的流式化对于注册回调 注销回调这类 API事件监听器、消息订阅等Stream.callback是万能方案。它接收一个基于Queue的函数在回调里用Queue.offerUnsafe向队列投递值再用Effect.acquireRelease注册监听器并注册流结束时注销的终结器import { Effect, Queue, Stream } from effect export const callbackStream Stream.callbackPointerEvent(Effect.fn(function*(queue) { function onEvent(event: PointerEvent) { Queue.offerUnsafe(queue, event) // 回调中向队列投递值 } // 注册监听器并在流结束时自动注销资源安全 yield* Effect.acquireRelease( Effect.sync(() button.addEventListener(click, onEvent)), () Effect.sync(() button.removeEventListener(click, onEvent)) ) }))这里体现了 Effect 的资源管理思想acquireRelease保证了无论流正常结束还是出错监听器都会被清理避免内存泄漏。2.7NodeStream.fromReadable接入 Node.js Readable 流在 Node 环境中effect/platform-node的NodeStream模块可以把任何Readable文件流、HTTP 响应体等转成 Effect Stream并支持错误转型与惰性求值import { NodeStream } from effect/platform-node import { Schema } from effect import { Readable } from node:stream export class NodeStreamError extends Schema.TaggedErrorNodeStreamError()(NodeStreamError, { cause: Schema.Defect() }) {} export const nodeStream NodeStream.fromReadable({ evaluate: () Readable.from([Hello, , world, !]), // 惰性创建 onError: (cause) new NodeStreamError({ cause }), // 错误转型 closeOnDone: true // 流结束后自动关闭底层 Readable默认为 true })evaluate字段让 Readable 的创建被推迟到流真正被消费时实现惰性求值closeOnDone: true默认值保证资源自动释放。2.8 小结构造方式速查表构造 API适用数据源关键参数Stream.fromIterable数组 / Set / Map / 生成器可迭代对象Stream.fromEffectSchedule周期任务指标、健康检查Effect ScheduleStream.paginate游标式分页 API初始游标 分页函数Stream.fromAsyncIterableasync generator异步可迭代对象 错误映射Stream.fromEventListenerDOM 事件元素 事件名Stream.callback任意回调式 API基于 Queue 的回调注册函数NodeStream.fromReadableNode.js Readableevaluate / onError / closeOnDone三、消费与转换 Stream算子编排数据管道.repos/effect-smol/ai-docs/src/03_stream/20_consuming-streams.ts以订单事件流为业务背景系统展示了转换算子与run*消费方法。文中定义了三层数据模型Order原始订单、NormalizedOrder加总金额、EnrichedOrder加税与优先级贯穿整个示例。3.1 纯转换算子Stream.map逐元素纯变换。把订单加总金额export const normalizedOrders orderEvents.pipe( Stream.map((order): NormalizedOrder ({ ...order, totalCents: order.subtotalCents order.shippingCents })) )Stream.filter按谓词剔除元素。只保留已支付订单export const paidOrders normalizedOrders.pipe( Stream.filter((order) order.status paid) )Stream.flatMap把每个元素展开为一个子流再扁平化合并适合一对多的 fan-out。示例为每个国家生成 1~49 号订单并用第二个参数控制并发度export const allOrders Stream.make(US, CA, NZ).pipe( Stream.flatMap( (country) Stream.range(1, 50).pipe( Stream.map((i): Order ({ ... /* 生成订单 */ })) ), { concurrency: 2 } // 可选控制内部流的并发度 ) )3.2 有效果的转换Stream.mapEffect当逐元素转换需要执行副作用如查税率、风控、调用外部服务时用Stream.mapEffect同样支持并发控制。示例中用concurrency: 4并行计算税额与优先级const enrichOrder Effect.fn(function*(order: NormalizedOrder) { yield* Effect.sleep(5 millis) // 模拟有效果的外部查询 const taxRate order.country US ? 0.08 : 0.13 const taxCents Math.round(order.totalCents * taxRate) return { ...order, taxCents, grandTotalCents: order.totalCents taxCents, priority: order.totalCents 20_000 ? high : normal } }) export const enrichedPaidOrders paidOrders.pipe( Stream.mapEffect(enrichOrder, { concurrency: 4 }) )3.3 消费方法run*家族流只有在被run系列方法消费时才会真正执行。常见的有消费方法行为示例Stream.runCollect收集全部元素为不可变数组Stream.runCollect(enrichedPaidOrders)Stream.runDrain只为副作用运行丢弃所有输出Stream.runDrain(enrichedPaidOrders)Stream.runForEach对每个元素执行一个有效果的消费者enrichedPaidOrders.pipe(Stream.runForEach((o) Effect.logInfo(...)))Stream.runFold归约为单个累积值enrichedPaidOrders.pipe(Stream.runFold(() 0, (acc, o) acc o.grandTotalCents))Stream.run(Sink)通过任意 Sink 消费enrichedPaidOrders.pipe(Stream.map(o o.grandTotalCents), Stream.run(Sink.sum))Stream.runHead取第一个元素返回Option...pipe(Stream.filter(...), Stream.runHead)Stream.runLast取最后一个元素返回Option...pipe(Stream.filter(...), Stream.runLast)其中runHead/runLast返回Option类型而非直接的值优雅地表达了流可能为空的情况——这一模式在 registry.test.ts 中也有体现Stream.runHead后接Effect.map(Option.getOrThrow)。runForEach的日志消费示例export const logOrders enrichedPaidOrders.pipe( Stream.runForEach((order) Effect.logInfo(Order ${order.id} total$${(order.grandTotalCents / 100).toFixed(2)}) ) )3.4 窗口类算子控制下游看到什么配合消费方法可以用窗口算子塑形流的可见部分Stream.take(2)只取前 2 个元素Stream.drop(1)跳过前 1 个元素如丢弃预热数据Stream.takeWhile(predicate)持续取元素直到谓词首次为假如只取priority normal的订单遇到 high 即停。export const firstTwoOrders enrichedPaidOrders.pipe( Stream.take(2), Stream.runCollect ) export const untilLargeOrder enrichedPaidOrders.pipe( Stream.takeWhile((order) order.priority normal), Stream.runCollect )四、流的编解码NDJSON 与 Msgpack.repos/effect-smol/ai-docs/src/03_stream/30_encoding.ts展示了流式数据传输的最后一块拼图用Stream.pipeThroughChannel配合effect/unstable/encoding下的Ndjson与Msgpack模块对流中的结构化数据进行解码/编码。文中明确指出所有 NDJSON 示例把Ndjson换成Msgpack、使用对应的Msgpack.decode()/Msgpack.encode()通道即可。4.1 领域模型与 Schema示例定义了结构化日志条目LogEntry其中Schema.DateTimeUtcFromString负责把 ISO-8601 字符串解码为DateTime.Utcimport { DateTime, Schema } from effect class LogEntry extends Schema.ClassLogEntry(LogEntry)({ timestamp: Schema.DateTimeUtcFromString, level: Schema.Literals([info, warn, error]), message: Schema.String }) {}4.2 解码NDJSON 字符串 → 对象无类型解码Ndjson.decodeString()通道Channel按换行拆分输入字符串并对每行执行JSON.parseexport const decodeUntyped Stream.make( {\timestamp\:\2025-06-01T00:00:00Z\,\level\:\info\,\message\:\start\}\n {\timestamp\:\2025-06-01T00:00:01Z\,\level\:\error\,\message\:\oops\}\n ).pipe( Stream.pipeThroughChannel(Ndjson.decodeString()), Stream.runCollect )带 Schema 校验解码Ndjson.decodeSchemaString(Schema)()在 JSON 解析的基础上进一步对每个值做 Schema 校验一步到位获得类型安全的强类型对象export const decodeTyped Stream.make(/* 同上的 NDJSON 文本 */).pipe( Stream.pipeThroughChannel(Ndjson.decodeSchemaString(LogEntry)()), Stream.runCollect )4.3 编码对象 → NDJSON 字符串Ndjson.encodeString()把每个值序列化为一行 JSONexport const encodeUntyped Stream.make( { timestamp: 2025-06-01T00:00:00Z, level: info, message: start }, { timestamp: 2025-06-01T00:00:01Z, level: error, message: oops } ).pipe( Stream.pipeThroughChannel(Ndjson.encodeString()), Stream.runCollect )Ndjson.encodeSchemaString(Schema)()先让值经过 Schema 编码应用日期格式化等转换再序列化为 NDJSON 行适合把DateTime.Utc等内部类型规范地输出为字符串export const encodeTyped Stream.make( new LogEntry({ timestamp: DateTime.makeUnsafe(2025-06-01T00:00:00Z), level: info, message: start }), new LogEntry({ timestamp: DateTime.makeUnsafe(2025-06-01T00:00:01Z), level: error, message: oops }) ).pipe( Stream.pipeThroughChannel(Ndjson.encodeSchemaString(LogEntry)()), Stream.runCollect )4.4 二进制变体面向 TCP 套接字与文件描述符处理二进制 I/O如 TCP 套接字、文件描述符时应使用非字符串变体Ndjson.decode()接收Uint8Array分块并内部处理文本解码Ndjson.encode()输出Uint8Arrayconst enc new TextEncoder() export const decodeBinary Stream.make( enc.encode({\level\:\info\,\message\:\binary\}\n) ).pipe( Stream.pipeThroughChannel(Ndjson.decode()), Stream.runCollect ) export const encodeBinary Stream.make( { level: info, message: binary } ).pipe( Stream.pipeThroughChannel(Ndjson.encode()), Stream.runCollect )4.5 空行处理ignoreEmptyLinesNDJSON 文件经常包含空行尾部换行、美化输出等。默认情况下空行会引发NdjsonError传入{ ignoreEmptyLines: true }即可跳过export const decodeIgnoringBlanks Stream.make( {\ok\:true}\n\n{\ok\:false}\n ).pipe( Stream.pipeThroughChannel(Ndjson.decodeString({ ignoreEmptyLines: true })), Stream.runCollect )4.6 错误处理NdjsonError编码失败kind: Pack或解码失败kind: Unpack时抛出Ndjson.NdjsonError可用Stream.catchTag或Effect.catchTag精确捕获。错误对象的kind字段指明发生在编码还是解码阶段cause携带底层异常export const handleDecodeErrors Stream.make(not-valid-json\n).pipe( Stream.pipeThroughChannel(Ndjson.decodeString()), Stream.catchTag(NdjsonError, (err) Stream.succeed({ recovered: true, kind: err.kind })), Stream.runCollect )4.7 完整实战解码 → 过滤 → 重新编码最常见的生产模式是读 NDJSON → 转换记录 → 再写出为 NDJSON。示例过滤出 error 级日志后重新编码输出const ndjsonInput {\timestamp\:\2025-06-01T00:00:00Z\,\level\:\info\,\message\:\ok\}\n {\timestamp\:\2025-06-01T00:00:01Z\,\level\:\error\,\message\:\fail\}\n {\timestamp\:\2025-06-01T00:00:02Z\,\level\:\warn\,\message\:\slow\}\n export const filterAndReencode Stream.make(ndjsonInput).pipe( // 1) 把每行解码为经过校验的 LogEntry Stream.pipeThroughChannel(Ndjson.decodeSchemaString(LogEntry)()), // 2) 只保留 error 级条目 Stream.filter((entry) entry.level error), // 3) 把过滤后的条目重新编码回 NDJSON 字符串 Stream.pipeThroughChannel(Ndjson.encodeSchemaString(LogEntry)()), Stream.runCollect )这个三步管道把格式解析、业务过滤、格式输出解耦为三个可独立测试、可替换的阶段是流式 ETL/日志管道的最佳实践模板。五、深入实践建议阅读仓库中的真实用法t3code 的 客户端连接层 与 registry.test.ts 都是 Stream 在生产/测试中的直接范例可作为阅读进阶代码的入口。结合 Effect 生态其他文档本目录隶属于.repos/effect-smol/ai-docs/src/的 Effect 系列教程涉及 Effect 基础、Schema、服务、调度等章节可与本文相互印证例如本文用到的Schema.TaggedError与Effect.fn分别对应 02_schema 与 01_effect 章节的内容。示例运行前提文中示例依赖effect及其平台包如effect/platform-node仓库根目录package.json的engines声明了node: ^24.13.1pnpm版本为11.10.0请在满足这些前提的环境下运行。示例按项目规范写在.ts文件中可通过项目的测试/运行链路直接执行ai-docs的package.json提供了 workspace 级别的effect依赖。命名约定.ts示例文件采用10_、20_、30_数字前缀控制渲染顺序见.repos/effect-smol/ai-docs/README.md理解这一约定有助于你定位不同主题的文档章节。结语Effect Stream 以有效果、拉取式、随时间三个特性把数组、轮询、分页、事件、Node 流等形态各异的数据源统一到同一个算子体系之下用from*系列构造用map/filter/flatMap/mapEffect转换用run*系列消费用pipeThroughChannel NDJSON/Msgpack 完成结构化数据的流式编解码。配合 t3code 仓库中真实的调用证据与完整示例源码这套组合足以支撑你在生产代码中写出类型安全、资源可控、可并发扩展的流式数据管道。【免费下载链接】t3code项目地址: https://gitcode.com/GitHub_Trending/t3/t3code创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考