尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

深入解析 @electric-sql/experimental:ElectricSQL 实验性 TypeScript 同步能力的安装、开发与测试

深入解析 @electric-sql/experimental:ElectricSQL 实验性 TypeScript 同步能力的安装、开发与测试 深入解析 electric-sql/experimentalElectricSQL 实验性 TypeScript 同步能力的安装、开发与测试【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electricelectric-sql/experimental是 ElectricSQL 官方在packages/experimental目录下发布的实验性 TypeScript 客户端包以electric-sql/client为基础提供matchStream/matchBy消息匹配工具与MultiShapeStream/TransactionalMultiShapeStream多 Shape 流等尚处于实验阶段的同步能力。本文基于该包源码packages/experimental与官方 README完整讲解其安装方式、核心 API 的用法与实现原理并给出从仓库构建、联调后端到运行测试套件的完整本地开发流程帮助你在真实项目中安全地试用这些前沿特性。一、这个包是什么实验性能力的定位与导出结构electric-sql/experimental定位为Experimental TypeScript features for ElectricSQL即 ElectricSQL 生态中尚未进入稳定客户端 API 的同步特性集合。它把electric-sql/clientpackages/typescript-client作为必选的 peer dependency 与唯一运行时依赖所有实验特性都构建在官方客户端之上因此需要先安装electric-sql/client才能使用。从 src/index.ts 可以看到包的全部公共 API 只有两个导出export * from ./match export * from ./multi-shape-stream即整个包由两大功能模块组成模块源码文件核心导出消息匹配工具src/match.tsmatchStream、matchBy多 Shape 流src/multi-shape-stream.tsMultiShapeStream、TransactionalMultiShapeStream、MultiShapeStreamInterface、MultiShapeMessagesBigInt 工具src/bigint-utils.tsbigIntMax、bigIntMin、bigIntCompare其中bigint-utils.ts不直接对外导出而是被multi-shape-stream.ts内部使用用于处理 Postgres LSN 这类以字符串形式传输的大整数比较与排序。二、安装一条命令接入实验特性官方 README 给出的一行安装命令即可完成接入npm i electric-sql/experimental根据 package.json 中的声明需要注意几个关键点peer dependencyelectric-sql/client是必选peerDependenciesMeta中optional: false使用前必须自行安装同版本客户端双模块格式exports同时提供了importdist/index.mjs 类型dist/index.d.ts与requiredist/cjs/index.cjs 类型dist/cjs/index.d.ctsESM 与 CJS 项目均可直接引用体积友好sideEffects: false支持 tree-shaking未使用的模块会被摇树移除licenseApache-2.0files只发布dist与src两个目录。包当前版本为 6.0.27从其 CHANGELOG.md 可以看出每个版本都会同步跟进electric-sql/client的更新例如 6.0.27 对应 client 1.5.27因此升级时建议保持两者版本同步。三、核心 API 一用 matchStream / matchBy 等待特定变更在同步场景中经常需要等到某条符合条件的数据变更到达再继续执行。src/match.ts提供的matchStream正是为此设计的 Promise 封装。3.1 matchStream 的签名与行为export function matchStreamT extends Rowunknown( stream: ShapeStreamInterfaceT, operations: ArrayOperation, matchFn: (message: ChangeMessageT) boolean, timeout 60000 // ms ): PromiseChangeMessageT从 src/match.ts 的实现可以梳理出完整行为对传入的ShapeStreamInterface调用subscribe订阅消息流在每个消息批次中先用isChangeMessage过滤出真正的数据变更消息再从中查找操作类型命中且matchFn返回 true的第一条命中后clearTimeout、unsubscribe并resolve(message)返回该ChangeMessage若在timeout毫秒默认 60000ms内未命中则打印错误日志并reject。参数含义参数说明stream任意ShapeStreamInterface实例通常是new ShapeStream(...)operations要匹配的操作数组如[insert]、[update]、[delete]或组合matchFn自定义匹配谓词入参为ChangeMessageT返回布尔值timeout等待超时毫秒默认 60 秒超时后 reject3.2 matchBy一键生成列值匹配器matchBy是matchStream最常用的配套工厂函数export function matchByT extends Rowunknown( column: string, value: ValueGetExtensionsT ): (message: ChangeMessageT) boolean { return (message: ChangeMessageT) message.value[column] value }它返回一个闭包谓词判断消息的value[column]是否严格等于目标值非常适合按主键等待某条记录写入。3.3 测试用例印证test/match.test.ts 中的should match用例演示了典型用法先创建一个订阅issues表的ShapeStream用matchBy(id, id)构造匹配器再在 10ms 后向数据库插入一条{ id, title: test title }记录最终断言matchStream(stream, [insert], matchFn, 200)解析出的result.value.title为test title。这条用例同时验证了匹配器按列值命中、operations过滤生效、以及消息在订阅后异步到达时的 Promise 解析机制。实际项目中的典型用法import { ShapeStream } from electric-sql/client import { matchBy, matchStream } from electric-sql/experimental const stream new ShapeStream({ url: http://localhost:3000/v1/shape, params: { table: issues, where: status \open\ }, }) const change await matchStream(stream, [insert, update], matchBy(id, a1b2c3), 30_000) console.log(目标变更已到达, change.value)四、核心 API 二MultiShapeStream 多 Shape 统一订阅单个ShapeStream只能订阅一个 Shape当需要同时监听多张表或同一张表的不同where过滤条件时MultiShapeStream提供了一个订阅、多 Shape 分发的能力并保证每个 Shape 在checkForUpdatesAfterMs时间间隔内至少收到一次来自 Electric 的up-to-date控制消息。4.1 构造选项与两种初始化方式从 src/multi-shape-stream.ts 的MultiShapeStreamOptions可知配置只有三个字段interface MultiShapeStreamOptionsTShapeRows { shapes: { [K in keyof TShapeRows]: | ShapeStreamOptionsTShapeRows[K] // 配置对象自动被包装成 ShapeStream | ShapeStreamTShapeRows[K] // 或直接传入已创建的 ShapeStream 实例 } start?: boolean // 是否立即启动默认 true checkForUpdatesAfterMs?: number // 强制检查更新的间隔默认 100ms }官方 JSDoc 给出的两种写法都受支持// 方式一直接传配置对象内部会自动创建 ShapeStream并强制 start: false 以便整体统一启动 const multiShapeStream new MultiShapeStream({ shapes: { shape1: { url: http://localhost:3000/v1/shape1 }, shape2: { url: http://localhost:3000/v1/shape2 }, }, }) // 方式二传入已构造的 ShapeStream 实例 const multiShapeStream new MultiShapeStream({ shapes: { shape1: new ShapeStream({ url: http://localhost:3000/v1/shape1 }), shape2: new ShapeStream({ url: http://localhost:3000/v1/shape2 }), }, }) multiShapeStream.subscribe((msgs) { console.log(msgs) })构造函数内部的关键逻辑对应源码第 140-165 行通过Object.fromEntries把shapes规范化为{ name: ShapeStream }映射若传入的是配置对象则用start: false创建ShapeStream保证所有 Shape 由多 Shape 流整体统一启动#start()中若发现某个 Shape 已启动会直接抛错维护#lastDataLsns与#lastUpToDateLsns两张 LSN 记账表初始值均为BigInt(-1)默认start true构造后立即#start()也可以传start: false延迟到首次subscribe时启动subscribe内部有if (!this.#started) this.#start()。4.2 消息增强每条消息带 shape 标识订阅回调收到的消息类型为MultiShapeMessagesTShapeRows它在原始ChangeMessage/ControlMessage基础上增加了一个shape字段标识消息来自哪个 Shapeinterface MultiShapeChangeMessageT, ShapeNames extends ChangeMessageT { shape: ShapeNames } interface MultiShapeControlMessageShapeNames extends ControlMessage { shape: ShapeNames }#start()中每个内部 Shape 的订阅回调会把收到的每条消息展开并附加shape: key后_publish给多 Shape 流的订阅者因此消费方可以据此区分消息来源。4.3 核心机制基于 LSN 的 checkForUpdates 协调这是MultiShapeStream最精巧的部分源码第 167-246 行核心思路是记账对每个 Shape 分别跟踪收到的数据消息最大 LSN#lastDataLsns与up-to-date控制消息的global_last_seen_lsn#lastUpToDateLsns触发只要任一 Shape 收到了新的数据消息maxDataLsn增大就通过#scheduleCheckForUpdates()在checkForUpdatesAfterMs默认 100ms后执行一次检查追赶刷新检查时计算所有 Shape 中最大的数据 LSNmaxDataLsn凡#lastUpToDateLsns[key] maxDataLsn的 Shape说明它还没追平最新进度调用shape.forceDisconnectAndRefresh()强制其重新拉取从而保证慢 Shape 不会被数据更快的新写入永久饿死。用源码中的注释概括多 Shape 流需要作为一个整体一起启动并确保所有 Shape 在checkForUpdatesAfterMs间隔内至少收到一次up-to-date消息这正是实时性一致性的保证。4.4 对外接口与状态查询MultiShapeStreamInterface源码第 56-78 行定义了完整的对外 API成员类型说明shapes{ [K]: ShapeStreamT }内部 ShapeStream 映射可用于访问shapeHandle等底层能力subscribe(cb, onError?)() void订阅消息返回取消订阅函数首个订阅触发启动unsubscribeAll()void清空所有订阅lastSyncedAt()number \| undefined最近一次同步的 Unix 时间戳取所有 Shape 的最小值isLoading时为undefinedlastSynced()number距上次同步的毫秒数尚未同步过返回InfinityisConnected()boolean所有 Shape 都已连接everyisLoading()boolean任一 Shape 仍在初始拉取中someisUpToDateboolean所有 Shape 都已up-to-dateevery错误处理方面#onError目前会把首个错误广播给所有订阅者的onError回调源码留有 TODO 注释提示未来可能改为在首个错误时断开全部 Shape。五、核心 API 三TransactionalMultiShapeStream 事务级批次MultiShapeStream的消息是来一批推一批但跨多个 Shape 的并发消息无法保证事务边界。TransactionalMultiShapeStream继承自MultiShapeStream并重写_publish实现按事务LSN分组、按op_position排序的批量发布确保同一次数据库事务产生的跨 Shape 变更作为一个整体一次性推送给订阅者。5.1 分组与发布算法源码第 379-467 行实现了完整逻辑累积#accumulate把收到的变更消息按headers.lsn快照消息没有 LSN用0代替分组暂存到#changeMessages映射中同时维护每个 Shape 的完整 LSN#completeLsns——当该 Shape 处于up-to-date且收到headers.last true的末条消息或收到up-to-date控制消息时用其global_last_seen_lsn更新取最低完整 LSN#getLowestCompleteLsnbigIntMin求出所有 Shape 完整 LSN 的最小值只有小于等于该值的 LSN 分组才被认为是所有 Shape 都已完整接收的事务可以安全发布发布按 LSN 升序bigIntCompare取出可发布分组组内按op_position数值升序排序快照消息没有op_position排序返回 0flat()后交给父类的_publish一次推给订阅者随后清理已发布分组。其设计目标可由源码注释直接印证利用lsn元数据推断事务边界利用op_position元数据对事务内消息排序。5.2 测试验证事务分组与跨事务隔离test/multi-shape-stream.test.ts 用两个where条件priority 10与priority 10把同一张issues表拆成lowPriority/highPriority两个 Shape做了系统性验证初始同步应作为一个分组断言messageGroups.length 1且组内恰好两条消息同一事务跨 Shape 变更合并在一个BEGIN/COMMIT内执行移动优先级的两次 update 和一次 insert最终得到一个含 5 条消息的分组2 delete 2 insert 1 insert并逐条断言 delete/insert 与 shape、id 的对应关系事务内按 op_position 排序断言组内所有headers.op_position升序跨事务边界隔离连续两个独立事务各含一次 update最终得到两个独立分组每组 2 条1 delete 1 insert且两组的 LSN 不同组内均按op_position排序。这些用例强有力地验证了事务边界被严格保持这一核心承诺是理解该特性行为的最佳参考资料。六、从源码构建与开发Developelectric-sql/experimental与仓库其他包一样使用 pnpm workspace 管理。官方 README 的完整开发流程如下。6.1 安装 workspace 依赖在仓库根目录执行pnpm install这会按 pnpm-workspace.yaml 声明的 workspace 关系安装全部包其中electric-sql/experimental以workspace:*协议引用本地electric-sql/client见其 devDependencies 与 peerDependencies。6.2 构建由于本包依赖electric-sql/client的产物需要先构建客户端再构建本包cd packages/typescript-client pnpm build构建脚本来自 package.jsonbuild: shx rm -rf dist tsup tsc -p tsconfig.build.json即先用 tsup配置见 tsup.config.ts打包出 ESM/CJS 产物再用tsc -p tsconfig.build.json生成类型声明。七、运行测试套件Test该包的测试是集成测试需要真实运行 Electric sync-service 后端与 Postgres而非纯单元测试。官方 README 给出了完整的两步流程。7.1 第一步启动后端在另一个终端中先进入sync-service目录安装依赖并启动cd ../sync-service mix deps.get mix stop_dev mix compile mix start_dev ies -S mix这条命令链依次完成拉取 Elixir 依赖 → 停止旧的 dev 实例 → 编译 → 启动 dev 实例 → 通过 IEx交互式 Shell以mix运行。需要说明的是该命令面向的是仓库本地的 Elixir 开发环境要求本机已具备 Elixir/Erlang 与 Postgres 环境sync-service 的具体配置可参考 packages/sync-service 的 AGENTS.md 与 README.md。7.2 第二步运行测试后端就绪后回到本包目录执行pnpm test测试基于 Vitestpnpm exec vitest配置见 vitest.config.tsglobalSetup: test/support/global-setup.ts启动前准备测试所需的数据库 schema 与 baseUrl 等全局注入项typecheck.enabled: true同时运行类型级测试对应 multi-shape-stream.test-d.tsfileParallelism: false串行执行避免共享后端资源冲突输出覆盖报告istanbul与 JUnit 报告./junit/test-report.junit.xml。测试基础设施test/support/test-context.ts为每个用例动态创建独立的issues表含id UUID PRIMARY KEY, title TEXT, priority INTEGER并通过 fixture 提供insertIssues、updateIssue、deleteIssue、beginTransaction/commitTransaction、clearIssuesShape、waitForIssues等数据库操作助手以及clearShape对/v1/shape发DELETE请求清理 Shape能力。这使得上面的match.test.ts与multi-shape-stream.test.ts可以稳定地验证消息匹配、跨 Shape 持续同步、取消订阅后不再收到回调、事务分组等行为。八、使用建议与风险提示作为experimental命名空间下的包使用时应注意API 不稳定matchStream、MultiShapeStream、TransactionalMultiShapeStream的签名与行为可能随版本调整参考 CHANGELOG.md 跟踪变更与客户端版本强绑定本包与electric-sql/client同步发版升级时两者需一起升级生产接入需谨慎实验特性主要面向验证与预研正式业务建议以稳定客户端 API 为主实验 API 作为补充了解底层协议理解lsn、global_last_seen_lsn、op_position、up-to-date控制消息等 Shape 协议概念详见 packages/typescript-client/SPEC.md有助于正确使用事务分组特性。综合来看electric-sql/experimental是探索 ElectricSQL 同步能力前沿的最佳入口matchStream/matchBy解决等待特定变更的编排问题MultiShapeStream解决多 Shape 统一订阅与进度追赶问题TransactionalMultiShapeStream则把同步推送到事务级一致性。结合本文给出的安装、构建与测试流程你可以在本地完整跑通这套实验特性并评估其是否适合你的业务。【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表