Electric 同步 Redis 示例用 ShapeStream 把 Postgres 数据实时同步为 Redis 哈希自动完成缓存失效【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric本文围绕 Electric 仓库中的 Redis 演示website/sync/demos/redis.md展开它展示如何用 TypeScript 客户端的ShapeStream订阅 Postgres 表的 shape 日志并把 insert/update/delete 变更实时写入一个 Redis 哈希Hash从而让 Electric 代替你管理缓存失效无需再手工维护 TTL 或写失效逻辑。读完本文你能理解该示例的完整运行流程、逐行掌握其核心代码包括用 Lua 脚本合并部分列更新这一关键技巧并能直接照搬到自己的缓存场景中。示例的核心定位让 Electric 接管缓存失效Redis 最常见的角色是缓存。而缓存最难维护的不是写而是失效数据变了谁负责把缓存里那条过期的记录删掉或更新传统的做法是应用代码在写库后再手动DEL或者给记录设置 TTL——两者分别有漏删导致脏读和无谓重建的问题。该示例的出发点是Redis 里的一份数据与其在 Postgres 中的权威副本之间天然是一条同步链路。Electric 恰好提供这条链路——它通过逻辑复制把 Postgres 的变更以 shape 日志的形式推送给客户端客户端把日志物化成一份与源表一致的数据。示例 examples/redis 就是把这个思路落到 Redis 上Electric 可以同步进 Redis 并自动管理缓存失效cache invalidation。你不需要单独处理缓存失效也不需要为缓存记录设置过期时间TTLElectric 会替你处理。引自 examples/redis/README.md。数据流可以概括为三段Postgres业务表items的增删改Electric 服务在http://localhost:3000/v1/shape端点以 SSE 形式推送 shape 日志初始快照 后续变更TypeScript 进程ShapeStream消费日志按行把变更翻译成 Redis 命令用 Redis 事务MULTI/EXEC批量落盘到一个名为items的哈希中。运行环境与启动步骤该示例是 ElectricSQL monorepo 中 pnpm workspace 的一部分因此所有命令都在仓库根目录的 workspace 语境下执行。完整的操作步骤继承自 examples/redis/README.md1. 在 monorepo 根目录安装并构建所有 workspace 包cd electric # 进入 monorepo 根目录 pnpm install pnpm run -r build2. 进入示例目录启动后端服务Electric Postgres基于 Docker Composecd examples/redis pnpm backend:up这一步会停掉并删除其他示例容器挂载的卷保证示例始终从一个干净的数据库和磁盘启动。从 examples/redis/package.json 可以看到backend:up实际是两段动作的串联backend:up: PROJECT_NAMEredis-example pnpm -C ../../ run example-backend:up pnpm db:migrate, db:migrate: dotenv -e ../../.env.dev -- pnpm exec pg-migrations apply --directory ./db/migrations即先拉起后端容器再对./db/migrations目录执行数据库迁移见下文数据模型一节。3. 启动同步进程pnpm dev # 等价于 tsx src/index.ts4. 用 redis-cli 观察同步结果redis-cli -h 127.0.0.1 -p 6379在 Redis 交互端里redis HKEYS items # 查看哈希中所有字段 redis MONITOR # 实时观看每一条到达 Redis 的命令5. 制造变更并观察实时同步另开一个终端连接 Postgres示例默认凭据psql postgresql://postgres:passwordlocalhost:54321/electricinsert into items (id, title) values (gen_random_uuid(), foo);执行插入后几乎立刻能在MONITOR输出或HKEYS items中看到新字段出现——这就是变更写入 Postgres 即同步进 Redis的效果。6. 结束时清理后端pnpm backend:down数据模型items 表同步的源头是一张极简的表定义在 examples/redis/db/migrations/01-create_items_table.sql-- Create a simple items table. CREATE TABLE IF NOT EXISTS items ( id TEXT PRIMARY KEY NOT NULL, title TEXT NOT NULL ); -- Populate the table with 10 items. -- FIXME: Remove this once writing out of band is implemented WITH generate_series AS ( SELECT gen_random_uuid()::text AS id, foo AS title FROM generate_series(1, 10) ) INSERT INTO items (id, title) SELECT id, title FROM generate_series;id是 TEXT 主键用gen_random_uuid()生成字符串形式迁移同时插入 10 条初始数据这样示例一启动Redis 哈希里就有可见内容。迁移由pg-migrations工具应用databases/pg-migrations在 examples/redis/package.json 的 devDependencies 中。核心代码逐段解析完整源码只有约 85 行位于 examples/redis/src/index.ts。下面按数据流顺序拆解。1. 连接 Redisimport { createClient } from redis const REDIS_HOST localhost const REDIS_PORT 6379 const client createClient({ url: redis://${REDIS_HOST}:${REDIS_PORT}, })使用的是官方的node-redisv4 客户端package.json 中依赖为redis: ^4.6.14连接本地默认端口的 Redis。连接建立后先清理旧数据client.connect().then(async () { console.log(Connected to Redis server) client.del(items) // 清掉上次运行留下的哈希del(items)保证每次运行都从空哈希开始使 Redis 中的数据与本次 shape 日志流一一对应。2. 为什么需要 Lua 脚本Electric 的 update 是部分列更新这是整个示例最关键的一处设计值得展开。Electric 推送的update变更消息里value只包含本次实际被修改的列而不是整行。客户端必须自己把这次的部分更新合并到已有行上。这一点可以从 TypeScript 客户端源码得到印证Shape类在内存中物化 shape 时对 update 消息执行的就是取出旧行、展开合并新值case update: this.#data.set(message.key, { ...this.#data.get(message.key)!, ...message.value, })见 packages/typescript-client/src/shape.ts约 L211-L215。在内存里一个 JS 对象展开合并就够用了。但在 Redis 中读出旧值 → 合并 → 写回是三步直接做会有并发竞态两个客户端或同一客户端的两个批次交错执行时可能互相覆盖对方的更新。示例的解法是把合并逻辑写成一段 Lua 脚本交给 Redis 原子执行const script local current redis.call(HGET, KEYS[1], KEYS[2]) local parsed {} if current then parsed cjson.decode(current) end for k, v in pairs(cjson.decode(ARGV[1])) do parsed[k] v end local updated cjson.encode(parsed) return redis.call(HSET, KEYS[1], KEYS[2], updated) const updateKeyScriptSha1 await client.SCRIPT_LOAD(script)脚本语义与客户端源码中的对象展开完全对应HGET读出该字段的现有 JSON可能不存在cjson.decode解析然后把ARGV[1]本次更新的部分列 JSON逐键覆盖进去最后HSET写回。Redis 保证 Lua 脚本在单线程中原子执行因此读-改-写不会被其他命令打断。SCRIPT_LOAD把脚本载入 Redis 并返回其 SHA1之后用EVALSHA调用可以省去每次传输脚本体。3. 订阅 shape 日志ShapeStreamconst itemsStream new ShapeStream({ url: http://localhost:3000/v1/shape, params: { table: items, }, }) itemsStream.subscribe(async (messages: Message[]) { /* ... */ })ShapeStream来自electric-sql/clientmonorepo 内对应 packages/typescript-client 包package.json 中声明为electric-sql/client: workspace:*。构造参数只有两个urlElectric 的 shape 端点/v1/shape提供 SSE 流params.table要同步的表名这里只同步items一张表。subscribe回调收到的是一个Message[]数组一个批次的日志消息Message是一个联合类型定义在 packages/typescript-client/src/types.tsexport type MessageT extends Rowunknown Row | ControlMessage // 控制消息up-to-date / must-refetch / snapshot-end / subset-end | EventMessage // 事件消息move-in / move-out子集查询的行进出 | ChangeMessageT // 变更消息带 key/value/headers.operation其中变更消息携带写入目标所需的全部信息export type ChangeMessageT extends Rowunknown Row { key: string // 行的主键哈希中的 field value: T // 变更后的列值update 时仅含被修改的列 old_value?: PartialT // 仅当 replica 为 full 时的更新旧值 headers: Header { operation: insert | update | delete txids?: number[] tags?: MoveTag[] removed_tags?: MoveTag[] active_conditions?: boolean[] } }注意headers.operation是insert|update|delete三元字面量——这正是下一节switch的三个分支依据。4. 逐消息翻译成 Redis 命令批量事务执行itemsStream.subscribe(async (messages: Message[]) { const pipeline client.multi() // 开启一个 Redis 事务 messages.forEach((message) { if (!isChangeMessage(message)) return // 只处理变更消息 switch (message.headers.operation) { case delete: pipeline.hDel(items, message.key) break case insert: pipeline.hSet(items, String(message.key), JSON.stringify(message.value)) break case update: pipeline.evalSha(updateKeyScriptSha1, { keys: [items, String(message.key)], arguments: [JSON.stringify(message.value)], }) break } }) try { await pipeline.exec() // 整批作为单个事务执行 } catch (error) { console.error(Error while updating hash:, error) } })几个值得注意的实现细节isChangeMessage类型守卫批次里混有控制消息如up-to-date守卫的作用就是把它们过滤掉、并在 TypeScript 层面把Message收窄为ChangeMessage。其实现非常简洁见 packages/typescript-client/src/helpers.tsL28-L32export function isChangeMessageT extends Rowunknown Row( message: MessageT ): message is ChangeMessageT { return message ! null key in message }判断依据就是变更消息特有的key字段。三种操作映射到三种 Redis 语义Postgres 的DELETE→ 哈希字段的HDEL缓存项被移除INSERT→ 整行 JSONHSETUPDATE→EVALSHA调用前面那段合并脚本。映射是精确的镜像关系Redis 哈希items在任意时刻都与 Postgres 表items的行集合保持一致。事务批量化client.multi()开启 MULTI/EXEC整个批次的命令在exec()时作为一个不可分割的单元执行。这保证了一批 shape 日志要么全部落进 Redis、要么都不落避免读到批次内半更新的中间态。源码注释中还保留了 Redis 官方文档的建议// FIXME The Redis docs suggest only sending 10k commands at a time // to avoid excess memory usage buffering commands.即生产环境若要自己控制批大小单事务命令数不宜超过约 1 万条。错误处理exec()包在 try/catch 里单批失败只记录错误日志而不终止订阅——shape 流本身会继续推送后续变更Redis 侧的最终一致由后续批次逐步追平。这套方案能省掉什么、需要注意什么省掉的部分相对传统缓存模式手工缓存失效不需要在业务代码里写更新 DB 后DEL缓存键删除/更新都由 shape 日志驱动TTL 管理缓存项的存活由源表决定——行还在表里就一直在哈希里行被删了HDEL自然跟上快照一致性ShapeStream初始会推送快照含 schema 与控制消息因此冷启动时 Redis 会先被灌入存量数据再跟随增量而不是只收到从现在开始的变更。需要注意的限制示例中的连接参数localhost:3000/v1/shape、localhost:6379、Postgres 连接串都是本地开发环境的固定值部署到真实环境需要相应替换且需保证 Electric 已正确配置 Postgres 逻辑复制见 monorepo 内 website/docs/sync 下的同步文档系列每个update都要执行一次 Lua 脚本热点表高频更新时脚本调用会成为热点生产环境可以评估用HSET多字段写整行覆盖替代合并语义但前提是你能保证拿到整行值示例面向单表、单哈希的简单场景。多表、多级缓存、按条件订阅子集等需求需要在此基础上扩展ShapeStream的参数与写入策略。关键文件索引文件作用website/sync/demos/redis.md官方站点的 Redis 演示说明页本文对应的关联文档examples/redis/src/index.ts同步进程全部核心代码Redis 连接、Lua 合并脚本、ShapeStream 订阅与事务写回examples/redis/README.md启动、观察与验证步骤examples/redis/db/migrations/01-create_items_table.sqlitems表建表与初始数据迁移examples/redis/package.json依赖与dev/backend:up/db:migrate脚本定义packages/typescript-client/src/types.tsMessage/ChangeMessage/Operation等消息类型定义packages/typescript-client/src/helpers.tsisChangeMessage类型守卫实现packages/typescript-client/src/shape.tsShape类对 update 部分列的内存合并逻辑Lua 脚本的原型【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考