Electric Durable Streams:为 Agent Loop 而生的数据原语 📅 发布时间:2026/9/16 14:07:03 👁 浏览次数: Electric Durable Streams为 Agent Loop 而生的数据原语【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electricAgent 是有状态的agent loop 每转一圈都会累积消息、token 流、工具调用与执行结果这类新形态的数据。Electric 团队为它构建了 Durable Streams——一种持久、可寻址、实时的流式数据原语。本文基于 Electric 官方博客文章与仓库中的协议文档、Rust 参考服务实现完整还原这一原语的设计动机、协议语义与源码级实现细节读完你能掌握如何在本地跑起一个 Durable Streams 服务、理解其核心操作与扩展层Durable State / StreamDB / StreamFS以及它在 agent 会话韧性、多用户协作与审计合规场景中的落地方式。一、Agent Loop状态随迭代持续累积Agent 正在被大规模、加速地部署。其背后共同的执行模式就是agent loop一个 observe观察→ think推理→ act行动→ repeat重复的循环。Agent 接到任务后先推理该做什么再决定一个动作执行后把结果重新喂回自己的上下文作为新的观察。这个循环不断重复每一轮迭代都是一次完整的模型推理调用由模型决定下一步动作每一轮都会累积新的状态消息、工具调用、工具调用结果、观察、产物artifacts等。如果把 agent loop 看作一个工作周期那么累积下来的状态就是这份工作的产出循环运行得越久创造的价值越多。问题随之而来——这份持续累积的状态用什么来承载二、为什么 Postgres 同步不够一种新类型的数据Electric 的起点是为应用开发中的状态转移构建 Postgres Sync可参见 Postgres 同步文档。随着构建在其上的软件形态演变为 AI 应用与智能体系统团队得出一个判断AI 应用同样应该构建在同步之上。但通过 Postgres 来同步这条路被明确否定了。官方给出的算账逻辑很直接每秒 50 个 token × 1000 个并发会话 每秒 5 万次写入而 Postgres 的上限大约在每秒 2 万次写入。再加上中心化的延迟这笔账算不过来。团队的解法是把 agent 的写入路径直接接到一条日志尾部。Electric 中的 shape 本质上是可寻址、可订阅的日志——如果把数据库这一环拿掉、绕开中心化延迟让 agent 直接写进日志会怎样这就是 Durable Streams 的出发点。它建立在久经考验的 Electric 同步协议的泛化之上但该协议本身在博客原文中的定位是每日投递数十亿次状态变更这一表述属于官方文档口径。三、核心定义持久、可寻址、实时的流Durable Stream 是一条拥有自己 URL 的持久、可寻址、仅追加append-only日志你可以直接写入它、实时订阅它、从任意位置回放它。从源码结构看它刻意保持极简——本质上就是追加式的二进制日志负载可以是任何内容投递协议是标准 HTTP因此在哪里都能工作可被缓存并可通过现有 CDN 基础设施水平扩展。结合仓库中的协议概览文档 Electric Streams 总览一个流就是一个可以 POST、可以 GET 的 URL。协议定义了六个操作PUT /streams/my-stream # 创建 POST /streams/my-stream # 追加 GET /streams/my-stream?offset… # 读取 HEAD /streams/my-stream # 元数据 POST /streams/my-stream # 关闭带 Stream-Closed: true DELETE /streams/my-stream # 删除协议不规定任何特定的 URL 结构——/v1/stream/{id}、/events/{topic}都可以。偏移量Offset可回放性的根基流中每个位置由一个offset标识它是不可解析的不透明字符串令牌具有两个关键性质不透明Opaque永远不要解析或构造 offset把它当作从服务端收到的字符串使用字典序可排序同一流的两个 offset 可以用普通字符串比较判断先后。协议定义了两个哨兵值值含义-1流的起始位置等价于省略 offsetnow当前尾部——跳过全部已有数据只读新消息读取响应中的Stream-Next-Offset头告诉你下次从哪里读保存它即可随时从断点恢复——这正是从任意点回放Replayable的协议基础。三种读取模式追平 长轮询 SSECatch-up追平读从 offset 读取已有数据响应头中的Stream-Next-Offset指示下一位置Stream-Up-To-Date: true表示已追平当前全部数据Stream-Closed: true表示流永久关闭EOFlivelong-poll长轮询服务端挂起连接直到新数据到达或超时超时返回 204适合二进制内容livesseServer-Sent Events持续推送data事件负载与control事件元数据如 next offset、closed 状态。服务端会周期性约每 60 秒主动断开 SSE 连接以支持 CDN 连接折叠客户端用最后一个control事件中的streamNextOffset重连。一个典型的读取循环来自 协议总览文档offset -1 loop: GET /stream?offset{offset} process(response.body) offset response.headers[Stream-Next-Offset] if response.headers[Stream-Closed]: break # EOF if response.headers[Stream-Up-To-Date]: # 已追平——切换 live 模式或稍后再轮询幂等生产者让重试安全对 exactly-once 写入语义生产者可以用三个头标识自己头用途Producer-Id生产者的稳定标识如order-service-1Producer-Epoch生产者重启时递增建立新会话Producer-Seqepoch 内按请求单调递增的序号服务端为每个(stream, producerId, epoch)元组记录已接受的最后序号重复序号直接返回去重后的成功响应而不重复写入。重启后递增 epoch 还能**围栏fence out**仍持有旧 epoch 的僵尸生产者——它们的请求会被403 Forbidden拒绝。JSON 模式结构化负载的批量化以Content-Type: application/json创建的流有特殊待遇每次 POST 的负载作为独立消息存储边界保留POST 一个 JSON 数组会被展平为多条消息一次 HTTP 请求即可批量写入GET 则返回包含范围内所有消息的 JSON 数组。这正是 agent 会话中结构化事件工具结果、状态变更的承载方式。四、为 Agent Loop 设计的九个属性博客原文给出的核心属性表完整保留属性意义Persistent持久agent 会话是持久的断线、重启后仍在Addressable可寻址找得到它们每个流有 URL每个位置有 offsetReactive响应式可以实时协作在同一个会话上Replayable可回放可以从任意点加入、审计或重启Forkable可分叉可以分叉会话去探索备选路径Lightweight轻量为每个 agent 即开即用创建成本极低Low-latency低延迟CDN 边缘可达个位数毫秒级延迟Schema-aware感知 schema复用结构化与多模态数据Extensible可扩展通过封装协议与集成扩展五、可扩展的分层栈从原始字节到结构化状态在核心开放协议之上Durable Streams 被设计为可组合的分层栈底层的原始二进制流负责字节投递上层协议负责结构化与多模态数据从而轻松接入智能体系统。数据层Durable State / StreamDB / StreamFSDurable State在流之上叠加结构化的状态变更协议。流改为承载带类型的insert/update/delete事件客户端按序应用即可物化状态。StreamDB一条流中的类型安全响应式数据库构建在 TanStack DB 之上。StreamFSagent 之间的共享文件系统。Durable State 文档给出了具体的事件格式——每个事件以typekey定位一个实体headers.operation携带操作类型{ type: user, key: user:123, value: { name: Alice, email: aliceexample.com }, headers: { operation: insert, txid: abc-123, timestamp: 2025-12-23T10:30:00Z } }字段约定type实体类型用于路由到正确集合与key必填value在 insert/update 时必填old_value可选用于冲突检测headers.operation必为insert/update/delete之一headers.txid、headers.timestamp可选。多个实体类型可以在同一条流中共存——一个聊天室流里可以交织user、message、reaction、typing事件全部按序处理。协议还定义了控制事件用于流管理而非数据变更控制类型用途snapshot-start标记快照开始——当前完整状态的一次全量倾倒snapshot-end标记快照边界结束reset通知客户端清空已物化状态并重启最小消费端是MaterializedState——一个内存键值存储把事件按(type, key)物化为最新值import { MaterializedState } from durable-streams/state const state new MaterializedState() state.apply({ type: user, key: 1, value: { name: Alice }, headers: { operation: insert }, }) state.apply({ type: user, key: 1, value: { name: Alice Smith }, headers: { operation: update }, }) const user state.get(user, 1) // { name: Alice Smith }集成层TanStack AI / Vercel AI SDK / Yjs博客原文列出的集成包括TanStack AI为 TanStack AI 应用添加 durable sessions 支持见 集成文档Vercel AI SDKdurable transport 适配器见 集成文档Yjs实时协作与 CRDT 支持附带快照发现、压缩、游标与用户状态见 集成文档。这些分层组合成的就是Durable Session模式一条持久、共享的会话同时复多路 AI token 流与结构化状态多个用户和 agent 可以随时订阅加入。协议分层顺序是Durable Streams——可靠、可恢复的字节投递State Protocol——流之上的结构化 CRUD 操作应用协议——AI SDK transport、presence、CRDT 等。这样AI agent 向会话中流式输出 token 的同时工具结果、用户在线状态、共享文档等结构化状态走的是同一套基础设施。六、源码级实现Rust 参考服务如何兑现这些属性仓库中的 durable-streams-rust 包 是协议的 Rust 参考实现它把博客中持久、低延迟、可扩展的抽象属性兑现为具体的工程决策。单二进制、零依赖的部署形态服务是一个自包含二进制没有数据库、没有 broker只有一个进程和一个数据目录。安装方式有三种Linux/macOSx64/arm64# 1. cargo构建二进制 durable-streams-server cargo install durable-streams # 2. npm下载预构建二进制 npm install -g electric-ax/durable-streams-server-rust # 3. Docker多架构镜像 docker run -p 4437:4437 electricax/durable-streams-server-rust关键标志与默认值所有标志都是可选的默认得到一个跑在127.0.0.1:4437的、数据目录位于$TMPDIR的持久单节点服务网络与存储标志默认值说明--host127.0.0.1监听地址0.0.0.0接受远程连接--port4437监听端口协议默认值--data-dir$TMPDIR/durable-streams-rust存储目录重启后数据保留--long-poll-timeout-ms30000livelong-poll请求最长阻塞毫秒数持久性Durability标志默认值说明--durabilitywalwal默认分组提交 fdatasync追加在记录落入分片 WAL 后才确认memory无 WAL 无 fsync页缓存写入即确认本地不抗崩溃--wal-shardsCPU 核数WAL 分片/分组提交器数量首次运行后持久化--wal-segment-bytes128 MiB每分片 WAL 段大小读取路径性能旋钮不改变协议行为标志默认值说明--tail-cache-bytes0Linux/65536macOS常驻尾部缓存字节上限Linux 上sendfile本身很快故默认关闭--read-offloadtailsendfile读取的调度位置inline/tail/always冷存储分层默认关闭--tier off|s3、--tier-endpoint、--tier-region、--tier-bucket、--tier-key-prefix、--tier-segment-bytes默认 8 MiB 固定段大小CDN 友好、--tier-compact-bytes默认 64 MiB 小段合并阈值等。S3 凭据只从环境变量读取DS_S3_ACCESS_KEY_ID/DS_S3_SECRET_ACCESS_KEY或标准AWS_*变量不走命令行。为什么写快分片 WAL 的分组提交从 架构文档 看核心论点是把每个流存成线上字节的原样写入就是追加读取就是字节区间。写路径的关键链路在src/handlers.rs的handle_append中解析幂等头Producer-Id/Producer-Epoch/Stream-Seq重复的(producer, epoch, seq)直接确认而不重复追加encode_wire把请求体编码为连续线上表示JSON 模式展平数组并追加分隔符获取每流追加互斥锁——这是唯一的串行化点不同流互不竞争write_all写入数据文件落入页缓存并推进写入者尾部先持久、后可见wal模式下追加被送入该流所属的 WAL 分片处理器等待该分片的分组提交fdatasync完成才推进读者可见的durable_tail、填充常驻尾部缓存、并通过 watch 通道唤醒实时订阅者然后才返回 2xx。不变式是读者永远只看到持久的字节一个并发写入批次共享同一次屏障 fsync因此吞吐量随每次 fsync 折叠的批大小扩展而不是每条消息一次 fsync。这也是实现上10,000 个流和 10 个流一样快的原因——单个提交器把多个流的追加合并进一次胖 WAL fsync。仓库 README 中的基准测试数据wal模式、8 核钉住的单节点 NVMe10 万流下约388k 次追加/秒1 万流约 422k/s无基数悬崖1 万→10 万仅降 8%p50 追加延迟约 2.5 ms追加路径是 fsync 受限的16 核达到约 537k/s 的真实平台期追平回放聚合约2.8 GiB/s512 并发连接内存随流数量而非字节数增长10 万流时容器工作集仍在数百 MiB。macOS 用户需注意wal模式在 macOS 上的持久性屏障是fcntl(F_FULLFSYNC)真正的磁盘写缓存刷新每组提交约 2–10 ms。若本机演示不需要断电持久性可用--durability memory写确认降到约 0.2–0.5 msDS_UNSAFE_FAST_FSYNC1仅用于基准测试生产环境切勿启用。为什么读快零拷贝与唤醒机制连续线上字节存储数据文件就是响应本身catch-up 读是字面意义的字节区间Linux 上用sendfile(2)从内核页缓存直达 socket零拷贝无重新分帧、无逐消息拷贝watch 通道唤醒每次追加在每流的 watch 通道上发一次通知即可唤醒所有长轮询与 SSE 订阅者——没有轮询循环、没有定时器空转常驻尾部缓存N 个已追平的订阅者共享同一段刚追加字节的一份读取与编码fan-out 去重有界内存大读按固定窗口流式下发读多 GB 回填的内存代价是一个窗口不是数据量。HTTP 层是一个手写 HTTP/1.1 引擎src/engine_raw.rs不依赖框架因此能直接持有 socket 做零拷贝发送。崩溃恢复与 Fork 支持每条 WAL 记录带头 CRC32C撕裂头检测器部分写入的头立即失败与负载 CRC32C恢复时校验撕裂或被清零的记录绝不重放启动时从数据文件 .meta边车重建每个流、重连 fork 链、重放 WAL 对齐未检查点的尾部实现清单中明确包含stream forks——这对应属性表中Forkable分叉会话以探索备选路径分叉后的流读路径会沿父链fork parent chain解析逻辑区间。冷存储分层历史数据交给对象存储与 CDN--tier s3可选特性cargo build --release --features tier利用追加式、位置不可变模型数据离开热尾部后即被**封存seal**为固定大小段并卸载到 S3 兼容对象存储R2、MinIO、B2 等最近的热尾部留在本地零拷贝服务。清单manifest是权威读取先对照每流清单sealed_offset之上走本地sendfile之下走远端 range-GET 拼接。持久性从不被削弱——追加仍在本地分组提交 fsync 后才确认卸载严格发生在持久化之后完全封存的区间不可变可直接带Cache-Control: immutable交给 CDN 吸收重复冷读。七、运行起来从零到实时订阅以下操作基于仓库中的 Rust 参考服务构建需要 Rust stable ≥ 1.75# 构建在 packages/durable-streams-rust 下 cargo build --release # → ./target/release/durable-streams-server cargo test --release # 单元 集成测试含协议一致性套件 # 启动 ./target/release/durable-streams-server --port 4438 --data-dir ./data创建流、追加、读回、检查元数据BASEhttp://localhost:4438/my-stream curl -X PUT $BASE -H Content-Type: application/octet-stream # 创建 curl -X POST $BASE -H Content-Type: application/octet-stream \ --data hello; # 追加 curl $BASE # 读取 → hello; curl -I $BASE # HEADoffset、length实时读取数据到达即读到curl $BASE?offsetnowlivelong-poll # 阻塞到下一次追加或超时 curl -N $BASE?offset0livesse # Server-Sent Events 流仓库还内置了协议一致性conformance测试可对接任何实现的服务器验证# 用较短的长轮询超时启动服务器后 RUST_SERVER_URLhttp://localhost:4562 pnpm exec vitest run \ --config packages/durable-streams-rust/conformance/vitest.config.ts托管侧则对应 Electric Cloud 中的 Electric Streams 服务——一个完全托管的 Durable Streams 协议实现开箱即带客户端 SDKTypeScript / Python、CLI、JSON 模式、Durable Proxy、Durable State、StreamDB、StreamFS 以及上文各集成。八、解锁韧性与协作回到博客原文的落点。Durable Streams 解锁韧性与协作的 agent 会话用户可以断线、重连、续跑无需重跑昂贵的推理工作既支持多用户同一会话的实时协作也支持跨时间接入、延续会话的异步协作Agent可以订阅并构建在其他 agent 的工作之上用户和 agent 可以派生、分叉子 agent组建 team、swarm 与层级化多 agent 系统且每一层都带有持久状态审计与合规系统可以记录每个 agent 动作的完整历史接入既有的审计与合规体系弹性扇出由于使用 Electric 投递协议借助现有 CDN 基础设施支持大规模弹性 fan-out 与并发可从零扩到百万级并发实时订阅者。九、小结Agent 是有状态的agent loop 每转一圈都在累积状态这份状态需要一种新的数据原语——这就是 Durable Streams一条拥有 URL 的追加式持久日志标准 HTTP 投递可订阅、可回放、可分叉并在其上生长出 Durable State、StreamDB、StreamFS 与 AI SDK 集成层。仓库内的 Rust 参考服务展示了这套原语在工程上的完整落地分片 WAL 分组提交兑现持久sendfile零拷贝与 watch 通道兑现低延迟fork 与 offset 语义兑现可回放、可分叉。如果你正在为 agent 会话寻找状态存哪、断了怎么办、怎么多人多 agent 共享的答案这套从协议到实现的路径协议总览 → 快速上手 → Durable State → Rust 参考服务值得逐层走一遍。【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考