数据同步设计文档
1.1 背景
Section titled “1.1 背景”有巢云 SaaS 平台采用”一主多子”的运营模式:运营主应用作为 SaaS 运营中心,管理所有租户的注册、功能配置和数据维护;各租户的子应用独立部署(独立数据库 + 独立服务实例),实现数据完全隔离。
租户在运营主应用中选择多个功能模块,组合成一个定制化的子应用系统。主应用维护的账号和组织机构数据需要准实时单向同步到对应租户的子应用。
1.2 术语定义
Section titled “1.2 术语定义”| 术语 | 含义 |
|---|---|
| 主应用 | SaaS 运营中心,管理租户/模板/账号/组织 |
| 子应用系统 | 租户专属系统,B/C/D…,独立数据库 + 独立进程 |
| 功能模板 | 预定义的插件组合方案,租户选择应用来初始化子应用 |
| 同步通道 | 基于 Redis Stream 的消息队列,承载主→子数据同步 |
| 事件发布器 | 主应用侧插件,监听数据变更并发布事件到 Stream |
| 同步消费者 | 子应用侧插件,消费 Stream 事件并写入本地数据库 |
1.3 现有系统基础
Section titled “1.3 现有系统基础”- 框架:数智基建
- 依赖中间件:PostgreSQL / Redis
2.1 功能目标
Section titled “2.1 功能目标”- 功能组合:租户在SAAS运营系统中选择若干插件,组合成一个完整的子应用系统
- 账号同步:主应用注册/修改/删除的租户账号,准实时同步到该租户的子应用
- 组织机构同步:主应用维护的租户组织机构(含员工/身份关系),准实时同步到子应用
- 子应用独立运行:各子应用独立部署、独立数据库、独立进程,互不影响
3. 架构设计
Section titled “3. 架构设计”3.1 整体架构
Section titled “3.1 整体架构”flowchart TB
subgraph SaaS运营主应用
TM[租户管理]
FM[功能模板管理]
OM[账号/组织机构管理]
EPS[事件发布器<br/>plugin-sync-sender]
end
SaaS运营主应用 -->|发布同步事件| R["Redis Stream 同步通道<br/>tenant:{tid}:sync<br/>plugin-sync-receiver"]
R --> TA["租户A<br/>子应用<br/>独立DB_A"]
R --> TB["租户B<br/>子应用<br/>独立DB_B"]
R --> TC["租户C<br/>子应用<br/>独立DB_C"]
style SaaS运营主应用 fill:#e6f7ff,stroke:#1890ff
style R fill:#fff7e6,stroke:#fa8c16
3.2 分层说明
Section titled “3.2 分层说明”| 层级 | 职责 | 实现 |
|---|---|---|
| 主应用层 | 租户管理、功能模板、账号/组织维护、事件发布 | 运营系统+ plugin-sync-sender |
| 同步通道层 | 消息有序传输、消费组隔离、可重播 | Redis Stream |
| 子应用层 | 按模板启用插件、消费事件、ID映射、幂等写入 | 实例 B/C/D + plugin-sync-receiver |
4. 数据同步设计
Section titled “4. 数据同步设计”4.1 同步范围
Section titled “4.1 同步范围”| 表名 | 主键 | 关联字段处理 |
|---|---|---|
| users | id (bigInt) | 无外键依赖 |
| organizations | id (snowflakeId) | parentId → 查本地映射 |
| employees | id (snowflakeId) | userId → 查本地映射 |
| identities | id (snowflakeId) | employeeId / departmentId / userId → 查映射 |
| rolesUsers | 复合 | userId → 查映射;roleName 不映射 |
4.2 同步通道设计
Section titled “4.2 同步通道设计”技术选型:Redis Stream
Section titled “技术选型:Redis Stream”Stream 设计
Section titled “Stream 设计”- Stream Key:
tenant:{tenantId}:sync,每个租户独立 Stream - Consumer Group:
subapp-{instanceId},每个子应用实例独立消费组 - 消息上限:
MAXLEN ~ 10000,防止 OOM
4.3 事件设计
Section titled “4.3 事件设计”事件 Payload 结构
Section titled “事件 Payload 结构”{ "tenantId": 1001, "entityType": "users", "entityId": 10086, "action": "create", "data": { "id": 10086, "username": "zhangsan", "nickname": "张三", "email": "zhangsan@example.com", "phone": "13800138000" }, "eventId": "evt_1724409600000_abc123", "timestamp": 1724409600000, "sourceUser": 1}| action | 触发 Hook | 子应用处理 |
|---|---|---|
| create | afterCreate | upsert(有则更新无则插入) |
| update | afterUpdate | upsert |
| delete | afterDestroy | 按主应用ID查本地记录并删除 |
4.4 同步流程
Section titled “4.4 同步流程”flowchart TB
A[主应用数据变更] --> B[钩子事件触发]
subgraph 钩子事件触发
H1[users.afterCreate/afterUpdate/afterDestroy]
H2[organizations.afterCreate/afterUpdate/afterDestroy]
H3[employees.afterCreate/afterUpdate/afterDestroy]
H4[identities.afterCreate/afterUpdate/afterDestroy]
H5[rolesUsers.afterCreate/afterUpdate/afterDestroy]
end
B --> C[事件发布器构建payload]
C --> D["XADD 写入 Redis Stream"]
D --> E["子应用消费器 XREADGROUP 读取消息"]
E --> F1["① 幂等检查:eventId 是否已处理<br/>(syncIdempotent 表)"]
F1 --> F2["② ID映射:主应用ID → 子应用本地ID<br/>(syncIdMap 表)"]
F2 --> F3["③ 关联字段替换:parentId/userId/employeeId等"]
F3 --> F4["④ 执行写入:upsert 或 delete"]
F4 --> G["XACK确认消费 + 记录幂等"]
style A fill:#e6f7ff,stroke:#1890ff
style D fill:#fff7e6,stroke:#fa8c16
style G fill:#f0fff4,stroke:#52c41a
4.5 ID 映射策略
Section titled “4.5 ID 映射策略”主应用和子应用的数据库独立,自增 ID 不一致。每张同步表在子应用侧增加 sourceId 列,存储主应用中的原始 ID,通过 sync_id_map 表维护映射关系。
| 表 | 需映射的外键 | 映射方式 |
|---|---|---|
| organizations | parentId | 查 syncIdMap 获取本地父节点 ID |
| employees | userId | 查 syncIdMap 获取本地 users.id |
| identities | employeeId, organiseId, userId | 分别查 syncIdMap |
| rolesUsers | userId | 查 syncIdMap 获取本地 users.id |
特殊处理:树形结构(organizations)
Section titled “特殊处理:树形结构(organizations)”organizations 表有 parentId 自引用。若子节点先于父节点到达,parentId 映射会失败。
处理策略:暂存到待处理队列,定时重试(指数退避 1s→2s→4s),父节点到达后自动消费。
4.6 幂等机制
Section titled “4.6 幂等机制”每条事件携带全局唯一 eventId,子应用消费前检查 syncIdempotent 表,已处理则直接跳过并 ACK。
防止场景:
- Redis Stream 消息重投递
- 消费者崩溃后消息重新可见
- 网络抖动导致的重复发送
4.7 异常处理与重试
Section titled “4.7 异常处理与重试”| 异常场景 | 处理方式 |
|---|---|
| 子应用离线 | 消息在 Stream 中积压,子应用上线后从 lastDeliveredId 重播 |
| 消费失败 | 不 ACK,消息留存 Stream,定时任务重试(3 次) |
| 重试超阈值 | 进入死信队列,触发告警 |
| Redis 宕机 | AOF 持久化恢复;主应用 afterHook 事务内先落库,Redis 恢复后扫表补发 |
| 树形顺序乱 | 待处理队列 + 指数退避重试 |
5.1 子应用数据库表结构
Section titled “5.1 子应用数据库表结构”-- ID 映射表CREATE TABLE syncIdMap ( id BIGINT PRIMARY KEY AUTO_INCREMENT, entityType VARCHAR(50) NOT NULL, sourceId BIGINT NOT NULL, localId BIGINT NOT NULL, UNIQUE KEY uk_entity_source (entityType, sourceId), INDEX idx_lookup (entityType, sourceId));
-- 幂等表CREATE TABLE syncIdempotent ( id BIGINT PRIMARY KEY AUTO_INCREMENT, eventId VARCHAR(64) NOT NULL UNIQUE, entityType VARCHAR(50), sourceId BIGINT, action VARCHAR(20), status VARCHAR(20) DEFAULT 'done', error_msg TEXT, INDEX idx_entity (entityType, sourceId));5.2 同步流程图
Section titled “5.2 同步流程图”flowchart TB
A["管理员在主应用注册用户 zhangsan"]
A --> B["users.afterCreate 钩子触发"]
B --> C["XADD tenant:1001:sync * event"]
C -->|消息内容<br/>tenantId:1001<br/>entityType:users<br/>entityId:10086<br/>action:create<br/>eventId:evt_xxx| D["Redis Stream"]
D --> E["子应用消费器 XREADGROUP 读取消息"]
E --> F{"幂等检查<br/>查询 sync_idempotent<br/>evt_xxx 是否已处理?"}
F --已处理--> Z["XACK 直接跳过"]
F --未处理--> G["ID映射处理<br/>无外键依赖,直接upsert"]
G --> H["写入本地users表<br/>sourceId=10086,本地id=1"]
H --> I["写入 sync_id_map<br/>(users, 主ID10086 ➜ 本地ID1)"]
I --> J["写入 sync_idempotent<br/>(evt_xxx, done)"]
J --> K["XACK 确认消费"]
style A fill:#e6f7ff,stroke:#1890ff
style D fill:#fff7e6,stroke:#fa8c16
style K fill:#f0fff4,stroke:#52c41a