Electric 发布托管版 Durable Streams 0.2.0:面向 AI Agent 与多用户协作的可恢复事件流
【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric
Durable Streams 是 Electric 从 Postgres 同步引擎中抽象出的核心原语——一种具备持久性、可恢复性的 HTTP 事件流协议。本文基于 0.2.0 版本的官方发布公告,系统讲解它为何是智能体(Agent)应用时代缺失的协调原语、0.2.0 带来的幂等生产者与 exactly-once 语义,以及它在 Electric Cloud 上的托管形态、性能表现与上手路径,并结合仓库中的协议文档与 Rust 服务端源码,帮助你理解其底层实现原理。
从 Postgres 同步引擎到 AI 时代的协调原语
Electric 团队花了三年时间构建 Postgres 原生同步引擎(见 1.0 发布公告),在这个过程中他们意识到,真正有价值的不只是"Postgres 同步"本身,而是隐藏在其下的传输原语。这一原语在 2025 年 12 月被独立发布为开源的 Durable Streams 协议,并在 2026 年 1 月 22 日随 0.2.0 版本正式托管到 Electric Cloud(官方公告)。
公告将这一转变归结为协调模型(coordination model)的进化:
- 传统的request/response 模型假设双方轮流发言:一个请求、一个响应。
- 而Agentic 应用(多智能体协作应用)中,多个智能体和多个用户同时在行动,他们需要的是共享的日志——每个人都能读、能恢复、能响应。
人们为了拼凑这个模型,不得不在 Redis、WebSockets 和重试逻辑之间反复横跳。而 Durable Streams 恰好是为此而生的原语:持久、可恢复、基于 HTTP 的事件流。公告同时推荐阅读同期发表的 Durable Sessions 协作文章,其中对交互范式的演进有更完整的论述。
缺失的原语:可寻址的 append-only 日志
核心模型
一个 durable stream(持久流)就是一个拥有自己 URL 的可寻址、append-only 日志。客户端可以从任意位置读取日志,也可以订阅(tail)实时更新。与之对比,现有流式基础设施并不适合这个场景:
| 方案 | 定位 | 局限 |
|---|---|---|
| WebSocket / SSE | 客户端流式连接 | 瞬态连接,断线即丢失在途数据,无法从特定位置恢复 |
| Kafka / Redis Streams | 后端原语 | 持久可靠,但客户端协议需要自己从零构建 |
Durable Streams 用纯 HTTP 解决了这个问题。仓库中的 Electric Streams 概览 给出了协议定义的六个操作:
PUT /streams/my-stream # Create(创建) POST /streams/my-stream # Append(追加) GET /streams/my-stream?offset=… # Read(读取) HEAD /streams/my-stream # Metadata(元数据) POST /streams/my-stream # Close(关闭,带 Stream-Closed: true) DELETE /streams/my-stream # Delete(删除)协议不规定具体的 URL 结构,可以是/v1/stream/{id}、/events/{topic}或任何适合应用的形态。
偏移量(Offset):恢复的基石
流中的每个位置由一个offset标识,它有两个关键性质:
- 不透明(Opaque):永远不要解析、构造或假设 offset 的内部格式,把它当作服务器返回的字符串即可。
- 字典序可排序(Lexicographically sortable):同一流中两个 offset 可以用标准字符串比较判断先后。
协议定义了两种哨兵值:
| 值 | 含义 |
|---|---|
"-1" | 流的开始位置(等价于省略 offset) |
"now" | 当前尾部位置——跳过所有已有数据,只读新消息 |
当读取响应返回时,会携带Stream-Next-Offset响应头,告诉客户端下次从哪里继续。客户端只需保存该值,即可在任何时候恢复。
内容类型与 JSON 模式
流的内容类型在创建时设定,决定了服务器如何处理消息边界:
- 字节流:对大多数内容类型(
application/octet-stream、text/plain、application/x-ndjson等),流是字节的原始拼接,消息边界由应用自己处理。 - JSON 模式:以
Content-Type: application/json创建的流会获得特殊处理——每次 POST 的载荷作为独立消息保存;POST 一个 JSON 数组[a, b, c]会存成三条独立消息(一次请求批量写多条);GET 返回请求范围内消息组成的 JSON 数组。
仓库的 JSON 模式文档 给出了完整示例。这条语义对 AI Agent 场景尤为重要:聊天消息、智能体事件、状态更新、日志都可以作为结构化消息流式传输。
0.2.0 的核心新增:幂等生产者与 exactly-once 语义
0.2.0 版本协议层面的核心增强是**幂等生产者(idempotent producers)**与exactly-once 写入语义。任何 HTTP 客户端都可以通过 POST 追加数据,但如果要获得 exactly-once 语义,生产者需要用三个请求头标识自己:
| 请求头 | 作用 |
|---|---|
Producer-Id | 生产者的稳定标识(如"order-service-1") |
Producer-Epoch | 生产者重启时递增,用于建立新会话 |
Producer-Seq | 同一 epoch 内单调递增的请求序号 |
三个头必须同时提供或同时不提供。服务器为每个(stream, producerId, epoch)元组追踪最近接受的序号:如果收到已经见过的序号,就返回一个去重后的成功响应而不是重复写入数据——这使得重试变得安全。
POST /streams/orders Producer-Id: order-service-1 Producer-Epoch: 0 Producer-Seq: 0 {"order": "abc"} < 200 OK重试同一请求时返回204 No Content,表示数据已经写入过。配合基于 epoch 的隔离(fencing):生产者重启后递增 epoch 并重置序号,服务器接受新 epoch,同时将仍使用旧 epoch 的"僵尸生产者"隔离出去——这些请求会收到403 Forbidden,从而防止崩溃重启后重复写入。
源码层面的印证
这一语义在仓库的 Rust 服务端实现中得到了落实。packages/durable-streams-rust/src/handlers.rs中的handle_append首先解析幂等头Producer-Id/Producer-Epoch/Stream-Seq,重复的(producer, epoch, seq)直接确认而不再追加(见 Rust 服务端架构文档 的 Write path 一节)。项目 README(packages/durable-streams-rust/README.md)也明确列出其协议覆盖范围:create / append / read(catch-up、long-poll、SSE)、HEAD、DELETE、JSON 模式、幂等生产者、close、TTL/expiry、cursors、ETag/304、安全头与流分叉(stream forks)。
托管版特性:Electric Cloud 上的 Durable Streams
0.2.0 的第二个重点是托管形态:Hosted Durable Streams 在Electric Cloud(云平台文档)上正式上线。Electric Cloud 同时托管Postgres 同步(sync 文档),因此你可以在同一个应用里把实时流与同步的关系型数据结合起来。官方公告给出的托管版关键能力如下:
- 读取不命中源站:Electric Cloud 的 Sync CDN 服务所有读取,单流已测试到100 万并发连接。
- 快速写入:小消息24 万次写入/秒,持续吞吐15-25 MB/秒。
- 简单定价:读取免费;每月前 500 万次写入免费,之后按量计费。
- 400+ 一致性测试(192 个服务端 + 212 个客户端)确保协议正确性。
- 10 种语言的客户端库:TypeScript、Python、Go、Rust、Java、Swift、PHP、Ruby、Elixir 和 .NET,全部通过完整一致性测试。
需要说明的是,上述数字来自官方发布公告;如果你希望自行验证开源自托管实现的性能基线,可参考仓库 Rust 服务端 README 中记录的基准数据:在wal模式下、8 个固定 CPU 上约 42 万次追加/秒(1 万流)与约 39 万次(10 万流),读取回放聚合吞吐约 2.8 GiB/秒;百万流基数测试文档 则记录了 1M 流基数下 110 万 ops/s 的压测结论——这些是同一套协议在开源实现上的独立证据。
一致性测试的意义
一致性测试是这套协议"跨语言可互操作"的保障。0.1.0 发布时(0.1.0 发布博客)服务端一致性测试为 124 个、客户端为 110 个,覆盖 offset 语义(单调性、字节级精确恢复、无跳过无重复、跨会话持久化)、重试行为(对 500/503/429 自动重试并尊重 Retry-After,对 4xx 不重试)、实时流模式(SSE 与 long-poll 的一致性)、消息顺序(所有读取模式下严格有序)以及生产者操作等。0.2.0 将总数推进到 400+,公告明确表示"协议已经成熟"。在仓库中,Rust 实现通过 conformance 测试目录 接入同一套一致性套件,CI 会对每种运行配置(wal/memory持久化、尾缓存开关、读卸载策略)跑完整套件。
快速开始:在 Electric Cloud 上创建第一个流
官方公告给出了托管版最简上手路径:
- 注册 Electric Cloud 并创建一个服务(service)。
- 用 curl 创建第一个流:
curl -X PUT \ -H "Authorization: Bearer <your-token>" \ -H "Content-Type: application/json" \ "https://api.electric-sql.cloud/v1/stream/<your-service-id>/my-stream"- 之后对它写、读、实时尾随——全部是普通 HTTP。
自托管快速开始
如果你希望先在本地体验协议本身,仓库的 快速开始文档 给出了自托管路径:下载durable-streams-server二进制后运行./durable-streams-server dev,即会在http://localhost:4437启动一个内存服务器,流端点为/v1/stream/*。然后:
# 创建流 curl -X PUT http://localhost:4437/v1/stream/hello \ -H 'Content-Type: text/plain' # 追加数据 curl -X POST http://localhost:4437/v1/stream/hello \ -H 'Content-Type: text/plain' \ -d 'Hello, Durable Streams!' # 读取全部(从 -1 开始) curl "http://localhost:4437/v1/stream/hello?offset=-1" # 实时尾随(终端一) curl -N "http://localhost:4437/v1/stream/hello?offset=-1&live=sse"在另一个终端对同一流追加数据,第一个终端会立刻收到新数据。这里的live=sse就是协议提供的实时模式之一;另一种是live=long-poll(长轮询,服务器挂起连接直到新数据到达或超时返回 204)。两种模式可以针对同一流互换使用。
深入实战:客户端库、CLI 与 Durable Proxy
TypeScript 客户端
仓库的 TypeScript 客户端文档 展示了三种 API 形态:
stream():fetch 风格的只读 API,支持offset、live参数,返回的StreamResponse提供body()/json()/text()/bodyStream()/jsonStream()/textStream()/subscribeJson()等多种消费方式。DurableStream:create / append / read / close / delete 的持久句柄,可在创建时指定ttlSeconds。IdempotentProducer:exactly-once 写入的推荐路径,内置自动批量(batching)与流水线(pipelining),并提供autoClaim与onError回调:
import { 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()live参数支持true(默认实时行为)、false(仅补读)、"sse"(强制 SSE)、"long-poll"(强制长轮询)四种取值。
CLI 工具
CLI 文档 描述了@durable-streams/cli的使用方式:全局安装后得到durable-stream命令,支持create(含--json快捷方式)、write(支持 stdin 管道与--batch-json数组扁平化)、read(先读全部历史再实时尾随)、delete四个子命令;通过STREAM_URL/STREAM_AUTH环境变量或--url/--auth标志配置连接,认证值原样作为Authorization头发送(支持Bearer、Basic等任意 scheme)。
Durable Proxy:让现有流式接口可恢复
公告中"Coming soon"部分提到的 HTTP 代理,在仓库中已有对应实现文档——Durable Proxy。@durable-streams/proxy可以将请求转发给上游 AI 流式接口,把流式响应持久化到 Durable Streams,并向客户端返回一个可以随时重连的持久读取 URL;客户端通过createDurableFetch(带requestId与autoResume)即可在刷新与断线后从原位置继续。这与公告"让你现有的 token 流在零代码改动下变得可恢复"的规划一致。
底层实现:Rust 服务端如何同时做到持久与高性能
自托管/开源实现(durable-streams-server,Rust 编写)展示了这一协议在工程上能达到的高度。架构文档 的核心论点是:把每条流以将要上线的字节原样存储,于是写入就是一次追加,读取就是一次字节区间读取。
- 连续线上字节存储:每条流的数据文件保存的正是读者收到的字节,没有按消息的重构、没有逐消息拷贝。读取是对文件的
pread/字节区间,在 Linux 上用sendfile(2)零拷贝(页缓存 → socket)服务。 - 分片 WAL 组提交:默认
wal模式下,追加在记录进入分片预写日志(WAL)后才确认;每个分片有一个组提交提交者,把大量流的追加合并进单次fdatasync(macOS 上为F_FULLFSYNC),因此吞吐不随流数量下降——对 1 万条流和 10 万条流同样快。 - 逐流串行化、读不加锁:每条流一把异步互斥锁排序追加;读取只做短暂快照和定位读,从不阻塞写者。
- watch 通道唤醒:long-poll 与 SSE 订阅者挂在每流 watch 通道上,追加发布新 tail 时一次唤醒,无轮询循环。
- 持久性门控的可见性:读者可观察的 tail 仅在组提交 fsync 完成后才发布(PROTOCOL.md §4.1),崩溃永远不会回滚读者已经看到的数据。
- 可观测性:通过
--features telemetry开启 OpenTelemetry,关键的ds.append.fsync.batch_size(组提交健康度)与ds.read.offload.wait(冷读池压力)两个指标用于生产监控。
可选能力还包括冷存储分层(--tier s3,把已封存的段卸载到 S3 兼容对象存储,历史补读从对象存储/CDN 提供)与尾缓存(--tail-cache-bytes,让 N 个已追平订阅者共享一次读取)。所有运行配置都是协议等价的,CI 会对每种配置跑完整一致性套件。
下一步与规划
公告明确列出了 0.2.0 之后的规划:Vercel AI SDK 与 TanStack AI 的即插即用 transport、Yjs 协作编辑支持、让现有 token 流可恢复的 HTTP 代理。在仓库的当前状态下,这些方向均已落地为文档与集成:TanStack AI 集成、Vercel AI SDK 集成、Yjs 集成 以及上文的 Durable Proxy,后续还有 流分叉(fork)、StreamDB、Durable State 等更高层能力。
公告最后给出的判断值得开发者参考:如果你已经写完 agent 循环、调试过 WebSocket 重连竞态、怀疑过 RedisPUBLISH是否真的投递了消息——现在可以停下来,把持久化与可恢复交给一个协议成熟、托管可用、且拥有 400+ 一致性测试与 10 种语言客户端的基础设施。协议是生产就绪的,剩下的问题——用它构建复杂、可塑、智能体化的应用时的体感与边界——正是社区需要共同探索的部分。你可以从 Cloud 快速开始 或 自托管快速开始 直接起步,在构建前先通过 协议概览 完整理解 offset、实时模式与生命周期语义。
【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考