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)。你不需要单独处理缓存失效,也不需要为缓存记录设置过期时间(TTL),Electric 会替你处理。
(引自 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.md:
1. 在 monorepo 根目录安装并构建所有 workspace 包:
cd electric # 进入 monorepo 根目录 pnpm install pnpm run -r build2. 进入示例目录,启动后端服务(Electric + Postgres,基于 Docker Compose):
cd examples/redis pnpm backend:up这一步会停掉并删除其他示例容器挂载的卷,保证示例始终从一个干净的数据库和磁盘启动。从 examples/redis/package.json 可以看到backend:up实际是两段动作的串联:
"backend:up": "PROJECT_NAME=redis-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:password@localhost:54321/electric"insert 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. 连接 Redis
import { 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 日志:ShapeStream
const itemsStream = new ShapeStream({ url: `http://localhost:3000/v1/shape`, params: { table: `items`, }, }) itemsStream.subscribe(async (messages: Message[]) => { /* ... */ })ShapeStream来自@electric-sql/client(monorepo 内对应 packages/typescript-client 包,package.json 中声明为"@electric-sql/client": "workspace:*")。构造参数只有两个:
url:Electric 的 shape 端点,/v1/shape提供 SSE 流;params.table:要同步的表名,这里只同步items一张表。
subscribe回调收到的是一个Message[]数组(一个批次的日志消息),Message是一个联合类型,定义在 packages/typescript-client/src/types.ts:
export type Message<T extends Row<unknown> = Row> = | ControlMessage // 控制消息:up-to-date / must-refetch / snapshot-end / subset-end | EventMessage // 事件消息:move-in / move-out(子集查询的行进出) | ChangeMessage<T> // 变更消息:带 key/value/headers.operation其中变更消息携带写入目标所需的全部信息:
export type ChangeMessage<T extends Row<unknown> = Row> = { key: string // 行的主键(哈希中的 field) value: T // 变更后的列值(update 时仅含被修改的列) old_value?: Partial<T> // 仅当 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.ts(L28-L32):export function isChangeMessage<T extends Row<unknown> = Row>( message: Message<T> ): message is ChangeMessage<T> { return message != null && `key` in message }判断依据就是变更消息特有的
key字段。三种操作映射到三种 Redis 语义:Postgres 的
DELETE→ 哈希字段的HDEL(缓存项被移除);INSERT→ 整行 JSONHSET;UPDATE→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.sql | items表建表与初始数据迁移 |
| examples/redis/package.json | 依赖与dev/backend:up/db:migrate脚本定义 |
| packages/typescript-client/src/types.ts | Message/ChangeMessage/Operation等消息类型定义 |
| packages/typescript-client/src/helpers.ts | isChangeMessage类型守卫实现 |
| packages/typescript-client/src/shape.ts | Shape类对 update 部分列的内存合并逻辑(Lua 脚本的原型) |
【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考