Files
BLD_sync/docs/full-sync.md
Misaka_Company bc461c12ce 📝 docs: add sync mechanism documentation with Mermaid diagrams
Add documentation for incremental sync (change-log driven polling,
per-PK verification, fault recovery) and full sync (batch transfer,
auto table creation, IDENTITY_INSERT handling).

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-06-12 11:08:06 +08:00

6.6 KiB
Raw Permalink Blame History

全量同步机制 (init_full_sync.py)

概述

全量同步用于初始化或重建 SQL Server 目标表,将 Access 数据源中的全部数据一次性加载到 SQL Server。通常在系统初始化、数据修复或新增表映射时执行。

  • 入口脚本: init_full_sync.py
  • 计划任务: AutoRun-init_full_sync(用户登录时自动启动)
  • 批量大小: 10,000 行(BATCH_SIZE

整体架构

graph TD
    subgraph 数据源
        A1["Access 文件 1<br/>成品入库.accdb"]
        A2["Access 文件 2<br/>25年压力表合同数据.accdb"]
        A3["Access 文件 N<br/>..."]
    end

    subgraph 全量同步
        S[init_full_sync.py]
    end

    subgraph SQL Server - CompanyDB
        T1["[schema1].[table1]"]
        T2["[schema2].[table2]"]
        T3["[schemaN].[tableN]"]
    end

    subgraph 通知
        N[ntfy 推送通知]
    end

    A1 & A2 & A3 -->|SELECT *| S
    S -->|TRUNCATE + INSERT| T1 & T2 & T3
    S -->|汇总通知| N

同步流程

flowchart TD
    START([开始全量同步]) --> SQL_CONN[连接 SQL Server]
    SQL_CONN --> FILE_LOOP{{遍历 SYNC_MAPPING<br/>中的每个 Access 文件}}

    FILE_LOOP --> FILE_CHECK{文件是否存在?}
    FILE_CHECK -->|不存在| SKIP[跳过该文件]
    FILE_CHECK -->|存在| ACC_CONN[连接 Access 数据库]

    ACC_CONN --> TABLE_LOOP{{遍历该文件下的<br/>每个表映射}}

    TABLE_LOOP --> ACC_READ[读取 Access 表结构<br/>SELECT TOP 1 * FROM table]

    ACC_READ --> TARGET_CHECK{目标表是否存在?}

    TARGET_CHECK -->|不存在| AUTO_CREATE[自动建表<br/>根据 Access schema 创建]
    TARGET_CHECK -->|存在| TRUNCATE[TRUNCATE 目标表]
    AUTO_CREATE --> DDL_COMMIT[DDL 单独提交]
    DDL_COMMIT --> IDENT_CHECK
    TRUNCATE --> IDENT_CHECK

    IDENT_CHECK{含 IDENTITY 列?}
    IDENT_CHECK -->|是| ID_ON[SET IDENTITY_INSERT ON]
    IDENT_CHECK -->|否| DATA_TRANSFER
    ID_ON --> DATA_TRANSFER

    DATA_TRANSFER[批量数据传输<br/>fetchmany BATCH_SIZE]
    DATA_TRANSFER --> HAS_MORE{{还有数据?}}

    HAS_MORE -->|有| BATCH[executemany 插入一个批次<br/>记录进度日志]
    BATCH --> HAS_MORE

    HAS_MORE -->|无| ID_OFF_CHECK{IDENTITY_INSERT<br/>是否已开启?}
    ID_OFF_CHECK -->|是| ID_OFF[SET IDENTITY_INSERT OFF]
    ID_OFF_CHECK -->|否| COMMIT
    ID_OFF --> COMMIT[提交事务]

    COMMIT --> STATS[记录表级统计<br/>行数 / 速率 / 用时]

    STATS --> NEXT_TABLE{{下一张表?}}
    NEXT_TABLE -->|是| TABLE_LOOP
    NEXT_TABLE -->|否| CLOSE_ACC[关闭 Access 连接]
    CLOSE_ACC --> FILE_SUMMARY[输出文件级汇总]

    SKIP --> FILE_LOOP
    FILE_SUMMARY --> FILE_LOOP

    FILE_LOOP -->|全部文件处理完| FINAL[输出全局汇总]
    FINAL --> NTFY[发送 ntfy 推送通知]
    NTFY --> END([结束])

    TABLE_LOOP -->|异常| ERR_HANDLE[记录失败日志<br/>回滚事务<br/>清理 IDENTITY_INSERT]
    ERR_HANDLE --> NEXT_TABLE

    style AUTO_CREATE fill:#bbf,stroke:#333
    style DATA_TRANSFER fill:#bfb,stroke:#333

数据传输细节

批量读取与插入

sequenceDiagram
    participant ACC as Access 数据库
    participant PY as Python 脚本
    participant SQL as SQL Server

    PY->>ACC: SELECT * FROM [table]
    loop 每 BATCH_SIZE 行
        ACC-->>PY: fetchmany(10000)
        PY->>SQL: executemany(INSERT, rows)
        Note over PY: 每 5 秒或每 10000 行<br/>记录一次进度
    end
    PY->>SQL: COMMIT

进度日志

传输过程中按时间和行数双条件输出进度:

表 [成品入库] → [成品入库记录]
  检测到 25 个列
  已清空目标表
  开始数据传输...
  进度: 10,000 行 | 速率: 45,000 行/秒
  进度: 20,000 行 | 速率: 43,500 行/秒
  表 [成品入库记录] 完成: 23,456 行 | 速率: 44,200 行/秒 | 用时: 0.5秒

自动建表

当目标表在 SQL Server 中不存在时,系统自动创建:

flowchart LR
    A[Access cursor.description] --> B[_access_col_to_sql<br/>类型映射]
    B --> C["CREATE TABLE [schema].[table] (<br/>  [col1] INT IDENTITY(1,1) PRIMARY KEY,<br/>  [col2] NVARCHAR(100),<br/>  ...<br/>)"]
    C --> D[DDL 单独提交<br/>确保后续参数绑定正常]

DDL 建表后必须单独 commit(),否则 pyodbc 的参数绑定无法获取正确的列元数据。

类型映射

Access 类型 SQL Server 类型 备注
int (主键) INT IDENTITY(1,1) PRIMARY KEY 自增主键
int (非主键) INT
float FLOAT
bool BIT
datetime.datetime DATETIME
decimal.Decimal DECIMAL(p, s) 保留精度
str (size ≤ 4000) NVARCHAR(size)
str (size > 4000) NVARCHAR(MAX) Memo 字段

IDENTITY_INSERT 处理

SQL Server 中含标识列(自增列)的表在插入显式 ID 值时,必须开启 IDENTITY_INSERT

stateDiagram-v2
    [*] --> Off
    Off --> On: SET IDENTITY_INSERT ON
    On --> Inserting: INSERT with explicit ID values
    Inserting --> On: executemany 完成
    On --> Off: SET IDENTITY_INSERT OFF
    Off --> [*]: COMMIT

    note right of On: 会话级设置<br/>不随事务回滚

关键点:

  • IDENTITY_INSERT会话级设置,每个连接同一时刻只能对一张表开启
  • 异常时必须在 except 块中显式关闭,否则后续同表同步会报错
  • 不随事务 ROLLBACK 回滚,必须手动关闭

通知机制

全量同步完成后通过 ntfy 发送汇总通知:

🎯 全量同步完成

✅ 成功 15 张表:
  • 成品入库记录: 23,456 行, 0.5秒
  • 25年压力表合同数据: 12,345 行, 0.3秒
  ...

❌ 失败 1 张表:
  • 某表名: 错误信息摘要...

总计: 35,801 行 | 用时: 2.3分钟
  • 全部成功:普通优先级
  • 存在失败:高优先级 + warning 标签

与增量同步的对比

维度 全量同步 增量同步
触发方式 用户登录 / 手动执行 持续轮询服务
数据范围 全部数据 仅变更记录
目标表处理 TRUNCATE + 全量 INSERT DELETE 旧 + INSERT 新(按主键)
校验机制 无逐条校验 逐主键校验 + 提交后复核
适用场景 初始化、数据重建 日常实时同步
运行时长 一次性执行完毕 常驻后台运行
数据完整性 依赖 TRUNCATE 原子性 事务 + 校验双重保障

配置参考

全量同步的行为由 config.py 中的 SYNC_MAPPING 驱动,新增表映射后首次运行全量同步即可自动建表并填充数据。