StreamDB 实战指南:在 Durable Stream 上构建类型安全的响应式数据库

📅 发布时间:2026/9/16 20:27:02
StreamDB 实战指南:在 Durable Stream 上构建类型安全的响应式数据库
StreamDB 实战指南在 Durable Stream 上构建类型安全的响应式数据库【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electricStreamDB 是 Electric Streams 生态中面向 Agent 会话与实时应用的状态层传入一个 Standard Schema 定义它就能把一条可持久、可回放的 Durable Stream 变成一个带类型集合、响应式查询与乐观写操作optimistic actions的流上数据库。读完本文你将掌握从定义 Schema、连接流、preload物化到用 TanStack DB 编写增量响应式查询、用事务 ID 做可靠写入的完整实战链路。StreamDB 是什么StreamDB 位于 Electric Streams 的分层协议之上与 Durable State 同属durable-streams/state包。它解决的核心问题是把一条 append-only 的字节流变成有结构的、类型安全、且默认响应式的数据库。从 协议分层 看其架构可以拆成三层Durable Streams—— 可靠的、可续传的字节投递层流是URL 可寻址、只追加、持久有序的字节序列协议操作见 Streams 概览State Protocol—— 在流之上定义结构化的insert/update/delete变更事件与快照控制事件StreamDB—— 消费这些事件按type路由进 TanStack DB 集合提供过滤、连接join、聚合与乐观更新能力。典型使用场景是 Agent 会话状态工具调用、消息、在线状态、Agent 注册表等多类实体天然适合复用同一条流。仓库中的 官方博客文章 即描述了如何用一套 Schema 同时承载messages、presence、agents三类实体。安装npm install durable-streams/state tanstack/db其中tanstack/db是 peer dependencyStreamDB 的集合与查询依赖它必须一并安装如果只需要MaterializedState这类纯物化层则可以跳过参见 Durable State 文档。定义 StandardSchema用createStateSchema定义状态结构。每个集合把一种实体类型映射到一个 Standard Schema 校验器和一个主键字段import { createStateSchema, createStreamDB } from durable-streams/state import { z } from zod const userSchema z.object({ id: z.string(), name: z.string(), email: z.string().email(), }) const messageSchema z.object({ id: z.string(), userId: z.string(), text: z.string(), timestamp: z.string(), }) const schema createStateSchema({ users: { schema: userSchema, type: user, primaryKey: id, }, messages: { schema: messageSchema, type: message, primaryKey: id, }, })Standard Schema 是一个跨库的校验协议任何实现了该协议的库都能用Zod、Valibot、ArkType或手写实现。Schema 还会生成类型化的事件辅助函数帮你构造合法的变更事件value是实体数据oldValue可用于冲突检测key是主键值schema.users.insert({ value: { id: 1, name: Alice, email: aliceexample.com }, }) schema.users.update({ value: updatedUser, oldValue: previousUser }) schema.users.delete({ key: 1 })底层事件格式这些辅助函数产出的正是 State Protocol 规定的标准变更事件。每个事件以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是该类型内实体的唯一标识valueinsert/update 时实体数据old_value否旧值用于冲突检测headers.operation是insert/update/delete之一headers.txid否用于确认写入的事务标识headers.timestamp否RFC 3339 时间戳同一条流可以共存多种实体类型——聊天室流可能交织着user、message、reaction、typing事件全部按序处理。此外协议还定义了snapshot-start/snapshot-end/reset等控制事件headers 中带control而非operation用于初始连接或 schema 迁移时的全量快照下发StreamDB 会透明处理这些边界。创建 StreamDBcreateStreamDB把 Schema 接到一条 Durable Stream 上得到一个响应式、由流支撑的数据库const db createStreamDB({ streamOptions: { url: https://api.example.com/streams/my-stream, contentType: application/json, }, state: schema, }) await db.preload()这里contentType必须是application/json——JSON 模式下服务端会保留消息边界、展开数组一次 POST 一个数组即批量多条消息、GET 返回 JSON 数组细节见 JSON mode 文档 与 Streams 协议概览。调用preload()会从流头部开始读取物化当前状态随后保持连接接收实时更新。物化过程就是把流上的变更事件按序应用进各集合——这也是MaterializedState所做工作的响应式增强版Durable State 文档 中MaterializedState是无 Schema、无响应式查询的极简形态StreamDB 则在其之上叠加了 TanStack DB 集合。仓库中的实际用法在 packages/agents-runtime/src/entity-stream-db.ts 中createEntityStreamDB正是用createStreamDB把 Agent 实体流runs、steps、texts 等内建集合加自定义 state 集合物化为带类型的 TanStack DB 集合。从源码可以观察到几个关键细节内建集合与自定义集合通过mergedCollections合并每个集合都映射{ schema, type, primaryKey }三元组type缺省为state:nameprimaryKey缺省为key与本文的createStateSchema定义完全一致集合 ID 通过getStreamDBCollectionId(streamUrl, name)生成用于把事件type反查回集合名collectionNameByEventType对reset控制事件会清空行偏移与时间线排序记录然后重新物化印证了控制事件的处理路径。响应式查询StreamDB 的集合就是 TanStack DB 集合。用useLiveQuery编写数据变化时自动更新的查询import { useLiveQuery } from tanstack/react-db import { eq, count } from tanstack/db const allUsers useLiveQuery((q) q.from({ users: db.collections.users })) const activeUsers useLiveQuery((q) q .from({ users: db.collections.users }) .where(({ users }) eq(users.active, true)) ) const messagesWithAuthors useLiveQuery((q) q .from({ messages: db.collections.messages }) .join({ users: db.collections.users }, ({ messages, users }) eq(messages.userId, users.id) ) .select(({ messages, users }) ({ text: messages.text, userName: users.name, })) ) const messageCount useLiveQuery((q) q .from({ messages: db.collections.messages }) .select(({ messages }) ({ total: count(messages.id) })) )TanStack DB 基于 differential dataflow差分数据流实现查询是增量更新的新事件到达时只重算受影响的数据而非全量重扫。这意味着跨集合 join、聚合、过滤等派生视图都可以组合复用。除 React 外官方还提供 Solid 与 Vue 适配器。派生集合原始流数据往往需要进一步物化——例如把 token 分片聚合成完整消息。派生集合derived collections用createLiveQueryCollection声明式完成且自身也是 TanStack DB 集合可继续查询、过滤、再派生。完整的分组 关联子查询 物化示例见 官方博客文章 的 Derive collections 一节chunks 同步在流上messages 从 chunks 物化approvals 从 messages 派生每一层都是响应式、类型安全、增量更新的。生命周期await db.preload() db.close() await db.utils.awaitTxId(txid-uuid, 5000)preload()从头读取并物化随后保持 live 连接close()释放连接与订阅资源。在 React 组件里务必配合useEffect的清理函数调用见下文 Best practicesawaitTxId(txid, timeoutMs)等待某个事务 ID 回传到本地物化状态第二个参数是超时毫秒数。这是确认写入已生效的关键工具。乐观操作Optimistic actionsStreamDB 通过 TanStack DB 的 action 系统支持乐观更新本地状态立即变更同时把变更异步持久化到流const db createStreamDB({ streamOptions: { url: streamUrl, contentType: application/json }, state: schema, actions: ({ db, stream }) ({ addUser: { onMutate: (user) { db.collections.users.insert(user) }, mutationFn: async (user) { const txid crypto.randomUUID() await stream.append( JSON.stringify( schema.users.insert({ value: user, headers: { txid } }) ) ) await db.utils.awaitTxId(txid) }, }, }), }) await db.actions.addUser({ id: 1, name: Alice, email: aliceexample.com })流程拆解onMutate立即把数据插入本地集合UI 在网络往返之前就已更新mutationFn生成txid把带txid头的事件 append 到 Durable Stream然后awaitTxId(txid)等待它从流回传确认若服务端写入失败TanStack DB 会自动回滚乐观更新。仓库源码在 entity-stream-db.ts 中展示了同一模式的自动化工序为每个自定义集合自动生成name_insert/name_update/name_delete三个 actiononMutate操作本地集合mutationFn统一走共享的持久化事务管线persistMutations并用WRITE_TXID_TIMEOUT_MS 20_000作为awaitTxId的超时上限源码第 107 行。常见模式键值存储Key/value store把主键设为key即可把流当作持久化的配置/键值存储const schema createStateSchema({ config: { schema: configSchema, type: config, primaryKey: key, }, }) await stream.append( JSON.stringify( schema.config.insert({ value: { key: theme, value: dark } }) ) )在线状态追踪Presence tracking用userId作主键每次更新都覆盖同一实体的最新状态const schema createStateSchema({ presence: { schema: presenceSchema, type: presence, primaryKey: userId, }, }) await stream.append( JSON.stringify( schema.presence.update({ value: { userId: alice, status: online, lastSeen: Date.now() }, }) ) )多类型聊天室Multi-type chat room多种实体类型复用同一条流天然按序处理、按类型路由const schema createStateSchema({ users: { schema: userSchema, type: user, primaryKey: id }, messages: { schema: messageSchema, type: message, primaryKey: id }, reactions: { schema: reactionSchema, type: reaction, primaryKey: id }, typing: { schema: typingSchema, type: typing, primaryKey: userId }, }) await stream.append(JSON.stringify(schema.users.insert({ value: user }))) await stream.append(JSON.stringify(schema.messages.insert({ value: message }))) await stream.append( JSON.stringify(schema.reactions.insert({ value: reaction })) )最佳实践使用对象值object values。StreamDB 的主键模式要求 value 是对象而非原始类型// Wont work { type: count, key: views, value: 42 } // Works { type: count, key: views, value: { id: views, count: 42 } }始终调用close()。在组件卸载时释放连接避免订阅泄漏useEffect(() { const db createStreamDB({ streamOptions, state: schema }) return () db.close() }, [])关键操作用事务 ID。为写操作附加txid并用awaitTxId确认生效超时按场景调整const txid crypto.randomUUID() await stream.append( JSON.stringify(schema.users.insert({ value: user, headers: { txid } })) ) await db.utils.awaitTxId(txid, 10000)在边界做校验。用 Standard Schema 对入口数据做严格约束失败发生在写入之前const userSchema z.object({ id: z.string().uuid(), email: z.string().email(), age: z.number().min(0).max(150), })快速上手一条流如果想在本地完整跑通建流 → 写 StreamDB → 实时读的链路可以先用仓库中的 Rust 参考服务packages/durable-streams-rust/README.md起一个本地流服务器./durable-streams-server --port 4438 --data-dir ./data然后按 Quickstart 用 curl 建流、追加、?offset-1livesse实时尾随。之后把createStreamDB的url指向本地流即可用文中代码开始验证集合物化、响应式查询与乐观写入。了解更多Durable State——底层 State Protocol 的变更事件、控制事件与MaterializedState详解Streams 协议概览——offset、消息边界、live 模式与流生命周期Quickstart——用 curl 快速体验 Durable StreamsJSON mode——StreamDB 依赖的 JSON 消息语义StreamDB 官方博客文章——派生集合与 Agent 会话状态实战createEntityStreamDB 源码——createStreamDB在 Electric Agents 运行时中的真实集成【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考