Electric Streams TypeScript 客户端完全指南:`stream()` 读取、`DurableStream` 读写与 `IdempotentProducer` 精确一次写入
Electric Streams TypeScript 客户端完全指南stream()读取、DurableStream读写与IdempotentProducer精确一次写入【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electricdurable-streams/client是 Electric Streams开放 Durable Streams 协议的托管实现官方 TypeScript 客户端。它直接面向协议提供读取与写入两大能力stream()提供 fetch 风格的流式读取DurableStream提供 create / append / read / close / delete 的完整句柄而IdempotentProducer则通过(producerId, epoch, seq)元组实现带批处理与重试的精确一次exactly-once写入。阅读本文后你将能够在 TypeScript 应用中直接消费 Durable Streams、安全地写入事件并实现断点续读的实时订阅。为什么需要 TypeScript 客户端Durable Streams 是一个URL 可寻址、append-only、可回放的字节序列每个流拥有独立 URL创建时确定内容类型数据一旦写入便不再改变只能向尾部追加。协议只定义了六种操作PUT创建、POST追加 / 关闭、GET读取、HEAD元数据与 DELETE删除全部基于标准 HTTP协议概览。因此任何 HTTP 客户端都可以读写流但直接裸调 HTTP 需要自己处理偏移量、批处理、重试、SSE / long-poll 连接生命周期等大量细节。durable-streams/client把这些协议细节封装成类型安全、fetch 风格的 APIstream()只读消费fetch 风格适合只消费不生产的应用DurableStream持久句柄支持创建、追加、读取、关闭、删除IdempotentProducer精确一次写入内置自动批处理、流水线与重试。在 packages/agents-server/src/stream-client.ts 中Electric 的 Agent 服务端正是基于这些 API 实现了完整的流封装层StreamClient涵盖了 create、fork、append、幂等追加、catch-up 读取、long-poll 等待、subscription 管理等场景是这套客户端在生产中的直接例证。安装npm install durable-streams/client安装后即可从包中导入stream、DurableStream、IdempotentProducer等符号。客户端不依赖浏览器特性Node.js 与浏览器环境均可使用。只读 APIstream()当应用只需要消费一个流时使用stream()。它是对协议 GET 请求的封装参数与 fetch 类似import { stream } from durable-streams/client const res await stream{ message: string }({ url: https://streams.example.com/my-account/chat/room-1, offset: savedOffset, live: true, }) const items await res.json() console.log(items)几个关键参数参数说明url流的完整 URL。协议不规定 URL 结构可以是/v1/stream/{id}或/events/{topic}见 Quickstart 中的/v1/stream/*端点示例offset起始偏移量。省略或传-1表示从头读取传now表示从当前尾部开始、跳过已有数据。偏移量是不透明字符串不要解析其内部格式只应保存服务器返回的值用于续读live实时模式开关取值见下文Live modes一节StreamResponse 的多种消费方式stream()返回的StreamResponse支持多种消费模式可按场景自由组合// 一次性获取 const bytes await res.body() // Uint8Array 原始字节 const items await res.json() // 解析为 JSON适合 JSON mode 流 const text await res.text() // 解码为文本 // 流式获取 const byteStream res.bodyStream() // ReadableStreamUint8Array const jsonStream res.jsonStream() // 逐个 JSON 消息 const textStream res.textStream() // 文本流 // 订阅式回调 const unsubscribe res.subscribeJson(async (batch) { await processBatch(batch.items) })subscribeJson以批次为单位回调批次内的items是结构化消息数组。务必保存订阅批次中返回的 offset下次连接时把它作为offset传入即可从相同位置继续消费实现精确的断点续读resumability。这一点在服务端封装中也有体现StreamClient.read()使用handle.stream({ offset, live: false })后调用response.body()返回原始字节而waitForMessages()则用response.subscribeBytes()逐块累积并用response.closedPromise 判断流结束packages/agents-server/src/stream-client.ts——这说明StreamResponse底层同时暴露字节块、offset 与结束状态供上层灵活组合。精确一次写入IdempotentProducer对于需要可靠、高吞吐且具备精确一次语义的写入场景使用IdempotentProducerimport { DurableStream, IdempotentProducer } from durable-streams/client const stream await DurableStream.create({ url: https://streams.example.com/events, contentType: application/json, }) const producer new IdempotentProducer(stream, event-processor-1, { autoClaim: true, onError: (err) console.error(Batch failed:, err), }) for (const event of events) { producer.append(event) } await producer.flush() await producer.close()这是需要安全重试与去重时推荐的写入路径。底层原理协议级幂等IdempotentProducer的精确一次语义来自协议本身。写入时客户端会携带三个请求头协议文档请求头用途Producer-Id生产者稳定标识例如event-processor-1Producer-Epoch生产者重启时递增用于开启新会话Producer-Seq同一 epoch 内单调递增的序号服务器为每个(stream, producerId, epoch)元组记录最后接受的序号收到已见过的序号时返回去重后的成功响应重复请求得到204 No Content不会重复写入数据——这使得网络错误后重试完全安全生产者重启后递增 epoch 并将 seq 归零服务器接受新 epoch 的同时会fence 掉仍在使用旧 epoch 的僵尸生产者其请求返回403 Forbidden从根源上杜绝崩溃重启后写入重复数据。客户端侧行为客户端在flush()时把批量数据合入少量请求配合流水线pipelining实现高吞吐autoClaim选项让生产者自动接管流onError提供批次失败的回调。close()会结束生产者并释放资源。在 Electric 服务端的封装中幂等追加是这样使用的packages/agents-server/src/stream-client.tsconst producer new IdempotentProducer(stream, opts.producerId, { epoch: opts.epoch ?? 0, }) try { producer.append(data) await producer.flush() } finally { await producer.detach() }注意这里在flush()后调用的是detach()而非close()detach()只解除生产者与流的绑定而不关闭流句柄适合一次写入后立即释放的场景而文档示例中close()用于整个写入会话结束时。读写 APIDurableStream当需要持久的句柄进行 create、append、read 等操作时使用DurableStreamimport { DurableStream } from durable-streams/client const handle await DurableStream.create({ url: https://streams.example.com/my-account/chat/room-1, contentType: application/json, ttlSeconds: 3600, }) await handle.append(JSON.stringify({ type: message, text: Hello })) const res await handle.stream{ type: string; text: string }() res.subscribeJson(async (batch) { for (const item of batch.items) { console.log(item.text) } })参数与能力说明contentType在创建时决定流的消息边界语义。设为application/json即启用 JSON mode——每条 POST 保存为独立消息、POST 数组会自动展平为多条消息、GET 返回 JSON 数组ttlSeconds流的相对存活时间对应协议头Stream-TTL。协议还支持绝对过期时间Stream-Expires-AtRFC 3339 时间戳二者互斥协议文档。超过 TTL 的旧数据可能被服务端清理此时读取已丢弃的 offset 会收到410 Gone应用应按需回退到-1或nowhandle.append()追加数据可传入字符串或字节handle.stream()发起读取同样返回StreamResponse支持上面所有消费模式。流的完整生命周期DurableStream句柄不只覆盖 create 和 append还提供元数据、关闭与删除操作。协议层的对应关系为协议文档创建PUT幂等重复创建同配置流返回200 OK追加POST响应携带Stream-Next-Offset头指明下一读取位置元数据HEAD返回Stream-Next-Offset、内容类型、TTL 与关闭状态不传输数据关闭带Stream-Closed: true的 POST标记流永久 EOF仍可完整读取旧数据仅拒绝新追加且可在同一次请求中原子地追加最后一条消息并关闭删除DELETE移除流及其全部数据。客户端中对应handle.head()返回{ exists, offset }与handle.close({ body })可在关闭时携带最后一条数据并返回finalOffset。StreamClient的exists()即通过DurableStream.head()实现append()在写入后通过head()获取最新 offset 返回给调用方packages/agents-server/src/stream-client.ts批量场景下还有静态形式DurableStream.create()、DurableStream.delete()、DurableStream.head()可直接调用。Live modes实时订阅模式stream()与handle.stream()的live参数有四种取值取值行为true使用该流类型默认的实时行为通常为 SSEfalse仅做 catch-up 读取不等待新数据sse强制使用 Server-Sent Eventslong-poll强制使用 HTTP 长轮询理解两种实时模式需要了解协议的 catch-up live 两段式模型消费方先用 offset 读取历史数据当收到Stream-Up-To-Date: true头表示已追平此时再切换 live 模式订阅新数据协议文档。long-poll服务器保持连接直到新数据到达或超时。超时返回204 No Content客户端用同一 offset 重试。适合二进制内容与简单请求/响应语义SSE服务器持续推送data事件携带流负载与control事件携带streamNextOffset、upToDate、streamClosed等元数据。客户端应保存最后一个control事件中的streamNextOffset用于重连若收到streamClosed: true则不再重连。服务器约每 60 秒主动关闭一次 SSE 连接以支持 CDN 连接折叠connection collapsing客户端需按此节奏自动重连。在 Rust 服务端实现中SSE 走专门的快路径engine_raw.rs在请求到达时判断响应是否可进入live tail fast path将连接移交给sse_reactor注册而长轮询唤醒则保证其尾部数据位于 page cache热路径以便快速响应packages/durable-streams-rust/src/engine_raw.rs、packages/durable-streams-rust/src/sse_reactor.rs。这说明客户端选择的 live 模式会直接影响服务端的 I/O 路径。客户端层waitForMessages()正是用live: long-pollAbortController超时实现等待新消息的语义packages/agents-server/src/stream-client.ts。与 JSON mode 配合使用客户端最典型的用法是搭配 JSON mode 使用。创建时设置contentType: application/json读取时传入泛型类型即可获得类型安全的结构化消息import { DurableStream, stream } from durable-streams/client const events await DurableStream.create({ url: http://localhost:4437/v1/stream/events, contentType: application/json, }) await events.append(JSON.stringify({ type: user.created, id: 123 })) const res await stream{ type: string; id: string }({ url: http://localhost:4437/v1/stream/events, json: true, }) const items await res.json() console.log(items)JSON mode 的核心行为详见 JSON mode 文档每个 POST 保存一条独立消息POST 数组如[a, b, c]会被展平为三条消息便于单请求批量写入GET 返回所请求范围内的 JSON 数组。因此 JSON 流天然适合聊天消息、Agent 事件、状态更新与日志而字节流application/octet-stream、text/plain等则适合原始二进制数据或已有自定帧格式的场景。何时使用直接构建在协议之上时使用 TypeScript 客户端——它是协议层的第一方封装offset、重连、幂等都由客户端处理流负载是结构化消息时使用 JSON mode将contentType设为application/json需要更上层的 AI 集成时可选用 Vercel AI SDK 或 TanStack AI 集成它们构建在 Durable Streams 之上提供可续传的 token 流与持久化 Agent 会话其他官方语言客户端Python、Go、Elixir、.NET、Swift、PHP、Java、Rust、Ruby 共 10 种实现了同一协议并通过一致的行为测试参见 客户端库总览。仓库中的更多线索协议概览与核心概念offset 哨兵值、幂等写入头、epoch fencing、catch-up/live 模型、流生命周期、CDN 缓存等完整规范Quickstart用durable-streams-server dev启动内存服务器并配合curl走通创建 / 追加 / 读取 / 实时尾随全流程服务端封装示例Electric Agent 服务端基于该客户端的生产级用法服务端一致性测试Rust 服务端通过durable-streams/server-conformance-tests跑完整协议一致性套件验证客户端所依赖的协议行为。【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考