拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Electric 分片(Sharding)实战指南:多实例部署与路由策略

Electric 分片Sharding实战指南多实例部署与路由策略【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric本指南围绕 website/docs/sync/guides/sharding.md 展开讲解如何让 Electric 与分片sharded的 PostgreSQL 数据库协同工作。当数据被横向拆分到多个 Postgres 分片时核心思路是为每个分片部署一个独立的 Electric 实例再通过客户端、代理或边缘层的路由逻辑把请求送到正确的实例。读完本文你将掌握分片架构下的多实例部署配置、三种路由方案、分片映射策略、与认证的结合方式以及健康检查与故障处理的最佳实践。为什么需要一实例一分片Electric 的一个关键约束是每个 Electric 实例只连接一个 PostgreSQL 数据库见 配置文件 中DATABASE_URL的语义。如果你的数据分布在多个 Postgres 分片上正确的做法不是让单个实例尝试连接多个库而是按分片数量部署多个 Electric 实例——每个实例负责一个分片然后根据数据所在位置把请求路由到对应实例。这种每分片一实例的模式带来三个核心收益独立扩展Independent scaling某个分片的负载上升时只需单独扩容该分片对应的 Electric 实例其他分片不受影响故障隔离Fault isolation一个分片的实例宕机不会波及其他分片上的用户灵活性Flexibility可以完全沿用你现有的分片方案按用户 ID、租户、地域等Electric 不强加自己的分片规则。架构示意┌─────────────────┐ │ Your App / │ │ Routing Proxy │ └────────┬────────┘ │ ┌─────────────────┼─────────────────┐ │ │ │ ▼ ▼ ▼ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ Electric │ │ Electric │ │ Electric │ │ (shard 0) │ │ (shard 1) │ │ (shard 2) │ └──────┬──────┘ └──────┬──────┘ └──────┬──────┘ │ │ │ ▼ ▼ ▼ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ Postgres │ │ Postgres │ │ Postgres │ │ (shard 0) │ │ (shard 1) │ │ (shard 2) │ └─────────────┘ └─────────────┘ └─────────────┘应用或一个路由代理负责判断请求的数据落在哪个分片再把 shape 请求转发给对应的 Electric 实例。Electric 本身不感知分片的存在路由是分片方案的上层职责。多实例部署每个分片的独立配置部署时每个分片对应一个 Electric 实例。由于这些实例彼此独立运行甚至可能运行在同一台宿主机上每个实例都需要一组唯一的配置以避免命名与存储冲突。每个实例所需的配置项配置项用途示例DATABASE_URL连接本分片对应的 Postgres 数据库postgresql://...shard-0/dbELECTRIC_INSTANCE_ID实例的唯一标识用于遥测telemetry区分electric-shard-0ELECTRIC_REPLICATION_STREAM_ID复制槽replication slot/ 发布publication名称的唯一后缀shard-0ELECTRIC_STORAGE_DIR持久化存储路径/data/shard-0其中两个配置项与唯一性直接相关值得深入理解其底层机制ELECTRIC_REPLICATION_STREAM_ID在 packages/sync-service/lib/electric/application.ex 中可以看到它是复制发布名与复制槽名的后缀。默认值为default见 packages/sync-service/lib/electric/config.ex最终会生成形如electric_publication_stream_id的发布名和electric_slot_stream_id的复制槽名。如果多个实例连到同一个数据库却不设置不同的 stream ID它们会争用同名的发布和复制槽产生冲突。分片场景下各实例连接的是不同数据库本不会冲突但设置唯一后缀仍然是规范做法——它让复制槽命名可读、可管理也便于日后排查。此外 config/runtime.exs 显示该值在启动时还会经过Electric.Postgres.Identifiers.parse_unquoted_identifier/1的标识符合法性校验。ELECTRIC_INSTANCE_ID默认是一个随机 UUIDElectric.Utils.uuid4()。为每个分片实例设置一个有意义的名字如electric-shard-0可以让 OpenTelemetry 指标、追踪信息在监控平台中按分片区分开。它还会被写入持久化存储见 packages/sync-service/config/runtime.exs 中的persist_installation_id调用。Docker Compose 示例以下docker-compose配置为两个分片各部署一个 Electric 实例。注意每个实例的端口、存储卷、环境变量都是独立的services: # Electric for Shard 0 electric-shard-0: image: electricsql/electric:latest environment: DATABASE_URL: postgresql://postgres:passwordpostgres-shard-0:5432/myapp ELECTRIC_INSTANCE_ID: electric-shard-0 ELECTRIC_REPLICATION_STREAM_ID: shard-0 ELECTRIC_STORAGE_DIR: /var/lib/electric/data ELECTRIC_SECRET: ${ELECTRIC_SECRET} ports: - 3001:3000 volumes: - electric_data_0:/var/lib/electric/data # Electric for Shard 1 electric-shard-1: image: electricsql/electric:latest environment: DATABASE_URL: postgresql://postgres:passwordpostgres-shard-1:5432/myapp ELECTRIC_INSTANCE_ID: electric-shard-1 ELECTRIC_REPLICATION_STREAM_ID: shard-1 ELECTRIC_STORAGE_DIR: /var/lib/electric/data ELECTRIC_SECRET: ${ELECTRIC_SECRET} ports: - 3002:3000 volumes: - electric_data_1:/var/lib/electric/data # Add more shards as needed... volumes: electric_data_0: electric_data_1:要点说明ELECTRIC_SECRET通过环境变量注入所有实例共享同一个密钥即可也可以每个分片单独管理ELECTRIC_STORAGE_DIR必须指向持久化存储Electric 依赖它保存 shape 日志与元数据重启后可以从中断处继续同步详见 配置文档 中的ELECTRIC_STORAGE_DIR说明每个实例通过不同的宿主机端口对外暴露 HTTP API3001、3002……容器内部统一监听 3000由ELECTRIC_PORT控制默认 3000。Kubernetes 示例在 Kubernetes 中可以为每个分片使用独立的 Deployment或使用 StatefulSet 按序管理。下面是electric-shard-0的 Deployment 清单apiVersion: apps/v1 kind: Deployment metadata: name: electric-shard-0 spec: replicas: 1 selector: matchLabels: app: electric shard: 0 template: metadata: labels: app: electric shard: 0 spec: containers: - name: electric image: electricsql/electric:latest env: - name: DATABASE_URL valueFrom: secretKeyRef: name: electric-shard-0-secrets key: database-url - name: ELECTRIC_INSTANCE_ID value: electric-shard-0 - name: ELECTRIC_REPLICATION_STREAM_ID value: shard-0 - name: ELECTRIC_STORAGE_DIR value: /var/lib/electric/data volumeMounts: - name: electric-storage mountPath: /var/lib/electric/data volumes: - name: electric-storage persistentVolumeClaim: claimName: electric-shard-0-pvc实践建议DATABASE_URL从 Secretelectric-shard-0-secrets中读取避免明文暴露数据库凭据每个分片使用独立的 PVC如electric-shard-0-pvc保证存储隔离与实例唯一性其他分片只需复制该清单并替换shard标签、名称、Secret 与 PVC 引用即可。路由策略把请求送到正确的实例分片方案能否落地关键在于路由——即如何确定某个 shape 请求属于哪个分片并把它转发到对应实例。你的应用天然知道每个用户的数据在哪个分片利用这个信息即可完成路由。原文档给出了三种层次的路由方案。方案一客户端路由Client-side routing最简单的方式客户端自己计算分片直接连接正确的 Electric 实例。import { ShapeStream } from electric-sql/client // Your sharding logic - determines which shard has the users data function getShardUrl(userId: string): string { // Option 1: Lookup from a shard directory const shardId shardDirectory.get(userId) // Option 2: Consistent hashing // const shardId hash(userId) % NUM_SHARDS return https://electric-shard-${shardId}.example.com } // Create a shape stream to the correct shard function createUserStream(userId: string) { const shardUrl getShardUrl(userId) return new ShapeStream({ url: ${shardUrl}/v1/shape, params: { table: user_data, where: user_id $1, params[1]: userId, }, }) }其中where: user_id $1与params[1]: userId是 Electric HTTP API 支持的参数化 WHERE 过滤where引用占位符$1params提供实际值详见 HTTP API 文档 中的where/params参数说明。这能确保客户端只拉取属于自己的数据行。客户端路由适合以下场景分片映射可以在客户端获取如发布时下发到客户端愿意向客户端暴露多个 Electric 端点希望服务端基础设施保持最简。方案二代理路由Proxy-based routing需要更多控制权时可以在服务端放一个代理由它决定分片归属对客户端隐藏分片细节。代理接收所有/v1/shape请求从中提取用户标识查得目标分片后把请求原样转发到对应 Electric 实例。// proxy/server.ts import express from express const app express() // Shard URL mapping const SHARD_URLS: Recordnumber, string { 0: http://electric-shard-0:3000, 1: http://electric-shard-1:3000, 2: http://electric-shard-2:3000, // ... add all shards } // Your sharding logic function getShardId(userId: string): number { // Lookup from database, cache, or compute via hashing return userShardMap.get(userId) ?? hashToShard(userId) } app.get(/v1/shape, async (req, res) { // Extract user identifier from request // Could come from: query params, JWT claims, headers, etc. const userId req.query.user_id as string || extractUserIdFromToken(req.headers.authorization) if (!userId) { return res.status(400).json({ error: user_id required for shard routing }) } // Determine target shard const shardId getShardId(userId) const targetUrl SHARD_URLS[shardId] if (!targetUrl) { return res.status(500).json({ error: Unknown shard: ${shardId} }) } // Build upstream URL with Electric protocol parameters only const upstreamUrl new URL(/v1/shape, targetUrl) Object.entries(req.query).forEach(([key, value]) { if (key user_id) return // Dont forward routing param if (Array.isArray(value)) { value.forEach(v upstreamUrl.searchParams.append(key, String(v))) } else if (value ! null) { upstreamUrl.searchParams.set(key, String(value)) } }) // Add Electric API secret as query parameter (not header) if (process.env.ELECTRIC_SECRET) { upstreamUrl.searchParams.set(secret, process.env.ELECTRIC_SECRET) } // Forward request to correct Electric instance const response await fetch(upstreamUrl) // Stream response back to client res.status(response.status) response.headers.forEach((value, key) { // Skip headers that shouldnt be forwarded if (![content-encoding, content-length].includes(key.toLowerCase())) { res.setHeader(key, value) } }) // Note: In Node.js 18, fetch returns a Web ReadableStream. // Use Readable.fromWeb() to pipe to Express response. if (response.body) { const { Readable } await import(node:stream) Readable.fromWeb(response.body as any).pipe(res) } else { res.end() } }) app.listen(3000)代理逻辑中有三个容易被忽略的关键细节路由参数不外传user_id仅用于代理侧路由判断转发前必须从 query 中剔除避免把它作为 shape 参数传给 Electric认证方式特殊Electric 的 API secret 是作为查询参数secret传递的而不是 HTTP Header这与 HTTP API 文档 的约定一致代理需要从环境变量中取出并注入到上游 URL流式转发shape 响应是持续的数据流长轮询必须用Readable.fromWeb()把 fetch 返回的 Web ReadableStream 管道到 Express 的响应对象而不是await response.json()一次性读完。客户端在代理模式下完全不知道分片的存在只需指向统一的 API 地址import { ShapeStream } from electric-sql/client // Client doesnt need to know about shards const stream new ShapeStream({ url: https://api.example.com/v1/shape, params: { table: user_data, user_id: currentUserId, // Proxy uses this for routing }, })[!Warning] 此示例仅演示路由 上面的代理只处理分片路由并不强制执行授权authorization。生产环境中代理还应校验用户身份并设置合适的where子句来限制数据访问范围。完整示例请参考下文与认证结合一节以及 认证指南。方案三边缘路由Edge routing如果前端有 CDN可以在边缘节点如 Cloudflare Worker实现路由进一步降低延迟同时把分片逻辑保留在服务端。边缘 Worker 与代理路由思路一致但运行位置更靠近客户端// Cloudflare Worker or similar edge function export default { async fetch(request: Request, env: Env): PromiseResponse { const url new URL(request.url) // Extract user ID from request const userId url.searchParams.get(user_id) || getUserIdFromJWT(request.headers.get(Authorization)) if (!userId) { return new Response(user_id required, { status: 400 }) } // Determine shard (edge KV lookup or compute) const shardId await getShardForUser(userId) // Remove routing-only params before forwarding url.searchParams.delete(user_id) // Add Electric API secret (injected at edge, never from client) if (env.ELECTRIC_SECRET) { url.searchParams.set(secret, env.ELECTRIC_SECRET) } // Route to correct Electric instance const targetOrigin https://electric-shard-${shardId}.internal const targetUrl new URL(url.pathname url.search, targetOrigin) return fetch(targetUrl, { headers: request.headers, }) }, }边缘路由的两点安全建议ELECTRIC_SECRET应作为 Worker 的环境变量secret binding注入绝不能由客户端传入分片映射可以通过边缘 KV 存储查询如await getShardForUser(userId)内部使用 KV或直接在边缘计算哈希。分片映射策略用户如何映射到分片无论采用哪种路由方式都需要回答同一个问题给定一个用户 ID它属于哪个分片选择哪种映射方式取决于你现有的分片方案。目录式映射Directory-based mapping把用户 → 分片的映射关系存放在一个快速查询服务中迁移灵活但需要一次额外查询// Using Redis async function getShardId(userId: string): Promisenumber { const shardId await redis.get(shard:${userId}) return parseInt(shardId ?? 0, 10) } // Using a database table async function getShardId(userId: string): Promisenumber { const result await db.query( SELECT shard_id FROM user_shards WHERE user_id $1, [userId] ) return result.rows[0]?.shard_id ?? 0 }哈希式映射Hash-based mapping对用户 ID 计算哈希并取模无需查询、延迟最低但新增分片时会导致大量数据迁移除非使用一致性哈希function getShardId(userId: string, numShards: number): number { // Simple hash let hash 0 for (let i 0; i userId.length; i) { hash ((hash 5) - hash) userId.charCodeAt(i) hash hash hash // Convert to 32-bit integer } return Math.abs(hash) % numShards } // Or use a proper consistent hashing library import { ConsistentHash } from consistent-hash const ring new ConsistentHash() ring.add(shard-0) ring.add(shard-1) ring.add(shard-2) function getShardId(userId: string): string { return ring.get(userId) // Returns shard-0, shard-1, etc. }ConsistentHash这类一致性哈希实现的最大价值在于当分片数量变化扩缩容时只有少量键需要重新映射其余键保持原分片把重新分片的成本降到最低。范围式映射Range-based mapping按照 ID 的数值范围划分分片规则直观、易调试但容易出现热点分片数据倾斜function getShardId(userId: string): number { const numericId parseInt(userId.replace(/\D/g, ), 10) if (numericId 1000000) return 0 if (numericId 2000000) return 1 if (numericId 3000000) return 2 // ... return 9 // Default shard }与认证结合路由 授权一体完成分片与 Electric 的认证模式可以天然配合。在代理中同时完成身份认证与分片路由并为用户设置数据访问边界app.get(/v1/shape, async (req, res) { // 1. Authenticate const user await validateToken(req.headers.authorization) if (!user) { return res.status(401).json({ error: Unauthorized }) } // 2. Determine shard from authenticated user const shardId getShardId(user.id) const targetUrl SHARD_URLS[shardId] // 3. Build request with authorization constraints const upstreamUrl new URL(/v1/shape, targetUrl) upstreamUrl.searchParams.set(table, req.query.table as string) // Enforce user can only access their own data using parameterized where upstreamUrl.searchParams.set(where, user_id $1) upstreamUrl.searchParams.set(params[1], user.id) // Add Electric API secret as query parameter if (process.env.ELECTRIC_SECRET) { upstreamUrl.searchParams.set(secret, process.env.ELECTRIC_SECRET) } // 4. Forward to Electric const response await fetch(upstreamUrl) // Stream response... })这个模式值得注意的组合逻辑身份来源于认证而非客户端user.id来自validateToken的结果而不是请求参数杜绝了客户端伪造user_id越权访问他人数据where由服务端强制注入无论客户端请求什么上游 URL 的where始终被重写为user_id $1配合参数化params[1]从服务端层面锁死数据范围路由与授权共享同一个user.id认证通过后分片判断和where约束基于同一身份逻辑上闭环。完整的认证方案如 JWT 验证、公开/私有 shape 区分参见 认证指南。健康检查与监控分片部署中每个实例都可能单独故障因此需要聚合的健康检查以便对整个分片集群有一致的可观测性// Health check aggregator async function checkAllShards(): PromiseShardHealth[] { const checks Object.entries(SHARD_URLS).map(async ([shardId, url]) { const controller new AbortController() const timeoutId setTimeout(() controller.abort(), 5000) try { const response await fetch(${url}/v1/health, { signal: controller.signal, }) clearTimeout(timeoutId) const data await response.json() return { shardId, status: data.status, healthy: response.ok, } } catch (error) { clearTimeout(timeoutId) return { shardId, status: unreachable, healthy: false, } } }) return Promise.all(checks) }每个 Electric 实例都暴露/v1/health端点聚合器对全部SHARD_URLS并行探测5 秒超时防止个别分片拖慢整体巡检每个实例通过 OpenTelemetry 导出指标。为每个实例设置唯一的ELECTRIC_INSTANCE_ID如electric-shard-0监控平台即可按分片区分各项指标快速定位是哪个分片出现了延迟或错误如需 Prometheus 指标可参考 遥测参考文档设置ELECTRIC_PROMETHEUS_PORT后暴露/metrics端点注意该端点必须被定期抓取否则内存中的直方图缓冲会无限增长。设计考量与运维要点数据本地性Data localityElectric 的每个实例只同步单个分片的数据。如果一个查询需要跨分片的数据客户端必须分别向涉及到的每个分片发起独立的 shape 请求在客户端合并结果。对大多数按用户或租户分区的业务来说这不成问题——某个用户的全部数据都落在同一个分片上一次 shape 订阅即可覆盖。故障转移Failover每个 Electric 实例彼此独立。某个分片的实例宕机时其他分片继续正常工作不受影响只有落在该分片上的用户受到波及重启实例后得益于持久化存储同步会从中断处恢复而不是从头开始。因此分片部署天然具备故障域隔离故障的影响半径被限制在单个分片内。重新分片Resharding当需要把用户迁移到新分片时标准流程是更新分片映射目录、配置等新请求开始路由到新分片客户端从旧分片收到must-refetch控制消息这是 shape 失效时的标准行为详见 HTTP API 文档 中控制消息的说明客户端丢弃本地 shape 数据向新分片重新发起完整同步。从 API 表面看切换分片是透明的URL 结构一致但客户端代码需要处理must-refetch并重建本地物化数据。electric-sql/client的ShapeStream会向下游发出must-refetch控制消息应用层应根据该信号触发重新订阅。这也意味着分片方案在设计时应保证客户端具备全量重同步的兜底能力它属于任何 shape 失效时的标准 Electric 行为。小结与后续方向分片的关键路径可以总结为四步按分片部署独立实例唯一配置→ 选择路由层次客户端 / 代理 / 边缘→ 确定映射策略目录 / 哈希 / 范围→ 叠加认证与可观测性。每一步都对应本指南中的可运行示例可以直接作为落地起点。进一步阅读部署指南生产环境下的完整配置含存储磁盘优化升级指南滚动部署与临时复制槽策略其中也对比了 sharding 与同一数据库多实例场景下 stream ID 的用法差异认证指南为分片部署加上身份校验与数据访问控制基准测试参考了解单个分片上的性能预期同步服务配置本文涉及的DATABASE_URL、ELECTRIC_INSTANCE_ID、ELECTRIC_REPLICATION_STREAM_ID、ELECTRIC_STORAGE_DIR等环境变量的完整说明。【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门