StreamDB 实战指南:在 Durable Stream 上构建类型安全的响应式数据库
2026/9/16 20:26:50 网站建设 项目流程

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

【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric

StreamDB 是 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 的字节流变成有结构的、类型安全、且默认响应式的数据库

从 协议分层 看,其架构可以拆成三层:

  1. Durable Streams—— 可靠的、可续传的字节投递层,流是"URL 可寻址、只追加、持久有序"的字节序列(协议操作见 Streams 概览);
  2. State Protocol—— 在流之上定义结构化的insert/update/delete变更事件与快照控制事件;
  3. StreamDB—— 消费这些事件,按type路由进 TanStack DB 集合,提供过滤、连接(join)、聚合与乐观更新能力。

典型使用场景是 Agent 会话状态:工具调用、消息、在线状态、Agent 注册表等多类实体天然适合复用同一条流。仓库中的 官方博客文章 即描述了如何用一套 Schema 同时承载messagespresenceagents三类实体。

安装

npm install @durable-streams/state @tanstack/db

其中@tanstack/db是 peer dependency,StreamDB 的集合与查询依赖它,必须一并安装(如果只需要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: "alice@example.com" }, }) schema.users.update({ value: updatedUser, oldValue: previousUser }) schema.users.delete({ key: "1" })

底层事件格式

这些辅助函数产出的正是 State Protocol 规定的标准变更事件。每个事件以type+key定位实体,以headers.operation表示操作:

{ "type": "user", "key": "user:123", "value": { "name": "Alice", "email": "alice@example.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.timestampRFC 3339 时间戳

同一条流可以共存多种实体类型——聊天室流可能交织着usermessagereactiontyping事件,全部按序处理。此外协议还定义了snapshot-start/snapshot-end/reset等控制事件(headers 中带control而非operation),用于初始连接或 schema 迁移时的全量快照下发,StreamDB 会透明处理这些边界。

创建 StreamDB

createStreamDB把 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:<name>primaryKey缺省为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 practices);
  • awaitTxId(txid, timeoutMs):等待某个事务 ID 回传到本地物化状态,第二个参数是超时毫秒数。这是确认"写入已生效"的关键工具。

乐观操作(Optimistic actions)

StreamDB 通过 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: "alice@example.com" })

流程拆解:

  1. onMutate立即把数据插入本地集合,UI 在网络往返之前就已更新;
  2. mutationFn生成txid,把带txid头的事件 append 到 Durable Stream,然后awaitTxId(txid)等待它从流回传确认;
  3. 若服务端写入失败,TanStack DB 会自动回滚乐观更新。

仓库源码在 entity-stream-db.ts 中展示了同一模式的自动化工序:为每个自定义集合自动生成<name>_insert/<name>_update/<name>_delete三个 action,onMutate操作本地集合,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 是对象而非原始类型:

// Won't 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=-1&live=sse实时尾随。之后把createStreamDBurl指向本地流,即可用文中代码开始验证集合物化、响应式查询与乐观写入。

了解更多

  • Durable State——底层 State Protocol 的变更事件、控制事件与MaterializedState详解
  • Streams 协议概览——offset、消息边界、live 模式与流生命周期
  • Quickstart——用 curl 快速体验 Durable Streams
  • JSON mode——StreamDB 依赖的 JSON 消息语义
  • StreamDB 官方博客文章——派生集合与 Agent 会话状态实战
  • createEntityStreamDB 源码——createStreamDB在 Electric Agents 运行时中的真实集成

【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询