跳转到内容

数据同步设计文档

有巢云 SaaS 平台采用”一主多子”的运营模式:运营主应用作为 SaaS 运营中心,管理所有租户的注册、功能配置和数据维护;各租户的子应用独立部署(独立数据库 + 独立服务实例),实现数据完全隔离。

租户在运营主应用中选择多个功能模块,组合成一个定制化的子应用系统。主应用维护的账号和组织机构数据需要准实时单向同步到对应租户的子应用。

术语含义
主应用SaaS 运营中心,管理租户/模板/账号/组织
子应用系统租户专属系统,B/C/D…,独立数据库 + 独立进程
功能模板预定义的插件组合方案,租户选择应用来初始化子应用
同步通道基于 Redis Stream 的消息队列,承载主→子数据同步
事件发布器主应用侧插件,监听数据变更并发布事件到 Stream
同步消费者子应用侧插件,消费 Stream 事件并写入本地数据库
  • 框架:数智基建
  • 依赖中间件:PostgreSQL / Redis

  1. 功能组合:租户在SAAS运营系统中选择若干插件,组合成一个完整的子应用系统
  2. 账号同步:主应用注册/修改/删除的租户账号,准实时同步到该租户的子应用
  3. 组织机构同步:主应用维护的租户组织机构(含员工/身份关系),准实时同步到子应用
  4. 子应用独立运行:各子应用独立部署、独立数据库、独立进程,互不影响

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
层级职责实现
主应用层租户管理、功能模板、账号/组织维护、事件发布运营系统+ plugin-sync-sender
同步通道层消息有序传输、消费组隔离、可重播Redis Stream
子应用层按模板启用插件、消费事件、ID映射、幂等写入实例 B/C/D + plugin-sync-receiver


表名主键关联字段处理
usersid (bigInt)无外键依赖
organizationsid (snowflakeId)parentId → 查本地映射
employeesid (snowflakeId)userId → 查本地映射
identitiesid (snowflakeId)employeeId / departmentId / userId → 查映射
rolesUsers复合userId → 查映射;roleName 不映射
  • Stream Key: tenant:{tenantId}:sync,每个租户独立 Stream
  • Consumer Group: subapp-{instanceId},每个子应用实例独立消费组
  • 消息上限: MAXLEN ~ 10000,防止 OOM
{
"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子应用处理
createafterCreateupsert(有则更新无则插入)
updateafterUpdateupsert
deleteafterDestroy按主应用ID查本地记录并删除
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

主应用和子应用的数据库独立,自增 ID 不一致。每张同步表在子应用侧增加 sourceId 列,存储主应用中的原始 ID,通过 sync_id_map 表维护映射关系。

需映射的外键映射方式
organizationsparentId查 syncIdMap 获取本地父节点 ID
employeesuserId查 syncIdMap 获取本地 users.id
identitiesemployeeId, organiseId, userId分别查 syncIdMap
rolesUsersuserId查 syncIdMap 获取本地 users.id

特殊处理:树形结构(organizations)

Section titled “特殊处理:树形结构(organizations)”

organizations 表有 parentId 自引用。若子节点先于父节点到达,parentId 映射会失败。

处理策略:暂存到待处理队列,定时重试(指数退避 1s→2s→4s),父节点到达后自动消费。

每条事件携带全局唯一 eventId,子应用消费前检查 syncIdempotent 表,已处理则直接跳过并 ACK。

防止场景:

  • Redis Stream 消息重投递
  • 消费者崩溃后消息重新可见
  • 网络抖动导致的重复发送
异常场景处理方式
子应用离线消息在 Stream 中积压,子应用上线后从 lastDeliveredId 重播
消费失败不 ACK,消息留存 Stream,定时任务重试(3 次)
重试超阈值进入死信队列,触发告警
Redis 宕机AOF 持久化恢复;主应用 afterHook 事务内先落库,Redis 恢复后扫表补发
树形顺序乱待处理队列 + 指数退避重试

-- 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)
);
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