📝 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>
This commit is contained in:
Misaka_Company
2026-06-12 11:08:06 +08:00
parent 17b6b246c7
commit bc461c12ce
2 changed files with 425 additions and 0 deletions

215
docs/full-sync.md Normal file
View File

@@ -0,0 +1,215 @@
# 全量同步机制 (init_full_sync.py)
## 概述
全量同步用于初始化或重建 SQL Server 目标表,将 Access 数据源中的**全部数据**一次性加载到 SQL Server。通常在系统初始化、数据修复或新增表映射时执行。
- **入口脚本**: `init_full_sync.py`
- **计划任务**: `AutoRun-init_full_sync`(用户登录时自动启动)
- **批量大小**: 10,000 行(`BATCH_SIZE`
## 整体架构
```mermaid
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
```
## 同步流程
```mermaid
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
```
## 数据传输细节
### 批量读取与插入
```mermaid
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 中不存在时,系统自动创建:
```mermaid
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`
```mermaid
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` 驱动,新增表映射后首次运行全量同步即可自动建表并填充数据。

210
docs/incremental-sync.md Normal file
View File

@@ -0,0 +1,210 @@
# 增量同步机制 (run_incremental_sync.py)
## 概述
增量同步是一个**长轮询服务**,持续监控 SQL Server 中的 `TableChangeLog` 变更日志表,将 Access 数据源的变更实时同步到 SQL Server 目标表。
- **入口脚本**: `run_incremental_sync.py`
- **计划任务**: `AutoRun-run_incremental_sync`(用户登录时自动启动)
- **轮询间隔**: 30 秒(`POLL_INTERVAL`
## 整体架构
```mermaid
graph LR
subgraph 数据源侧
A[Access .accdb 文件] -->|VBA 宏写入日志| B[TableChangeLog]
end
subgraph 增量同步服务
B -->|轮询 Synced=0| C[run_incremental_sync.py]
C -->|读最新数据| A
C -->|DELETE + INSERT| D[SQL Server 目标表]
C -->|标记 Synced=1| B
end
subgraph 监控
C -->|心跳| E[Uptime Kuma]
end
```
## 变更日志驱动
Access 端的 VBA 宏在数据变更时,向 `TableChangeLog` 表写入一条记录:
| 字段 | 说明 |
|------|------|
| `LogID` | 自增主键 |
| `TableAddress` | Access 文件路径(多种格式) |
| `TableName` | 发生变更的 Access 表名 |
| `RecordID` | 变更记录的主键值 |
| `Synced` | 同步标记0=未同步1=已同步 |
### 路径匹配策略
VBA 端写入的 `TableAddress` 有多种格式,系统构造 4 种候选值进行匹配:
```
;DATABASE=\\192.168.110.114\生产进度表\成品入库.accdb ← 网络路径前缀
LOCAL=\\server\share\file.accdb ← 本地等号前缀
LOCAL:\\server\share\file.accdb ← 本地冒号前缀
\\192.168.110.114\生产进度表\成品入库.accdb ← 裸路径
```
## 核心同步流程
```mermaid
flowchart TD
START([服务启动]) --> POLL[轮询 TableChangeLog<br/>WHERE Synced=0]
POLL -->|无记录| SLEEP[休眠 POLL_INTERVAL 秒]
SLEEP --> POLL
POLL -->|发现未同步记录| GROUP[按 TableName 分组<br/>合并同一表的 record_ids]
GROUP --> CONNECT[连接 Access 数据库]
CONNECT --> FOR_EACH{{遍历每个表}}
FOR_EACH --> CHECK_TABLE{目标表是否存在?}
CHECK_TABLE -->|不存在| AUTO_CREATE[自动建表<br/>ensure_schema + create_table_from_access]
CHECK_TABLE -->|存在| READ_ACCESS
AUTO_CREATE --> READ_ACCESS
READ_ACCESS[A. 从 Access 读取最新数据<br/>SELECT WHERE PK IN ...]
READ_ACCESS --> EXTRACT_PK[提取 inserted_pk_set<br/>Access 实际返回的主键集合]
EXTRACT_PK --> DELETE[B. 删除目标表旧记录<br/>DELETE WHERE PK IN ...]
DELETE --> HAS_NEW{{有新数据?}}
HAS_NEW -->|有| IDENTITY_CHECK{含 IDENTITY 列?}
IDENTITY_CHECK -->|是| ID_ON[SET IDENTITY_INSERT ON]
IDENTITY_CHECK -->|否| INSERT
ID_ON --> INSERT[C. 批量插入新记录<br/>executemany]
INSERT --> ID_OFF[SET IDENTITY_INSERT OFF]
ID_OFF --> VERIFY
HAS_NEW -->|无| VERIFY
VERIFY[D. 提交前逐主键校验] --> VERIFY_OK{校验通过?}
VERIFY_OK -->|通过| MARK_SYNCED[E. 标记 Synced=1]
MARK_SYNCED --> COMMIT[F. 提交事务]
COMMIT --> POST_VERIFY
POST_VERIFY{提交后复核<br/>ENABLE_POST_COMMIT_VERIFY} -->|关闭| SUCCESS
POST_VERIFY -->|开启| POST_OK{复核通过?}
POST_OK -->|通过| SUCCESS[记录同步成功日志]
POST_OK -->|失败| REVERT[回滚 Synced=0<br/>下一轮重试]
VERIFY_OK -->|失败| ROLLBACK[回滚事务<br/>保持 Synced=0]
ROLLBACK --> NEXT_TABLE
SUCCESS --> NEXT_TABLE
REVERT --> NEXT_TABLE
NEXT_TABLE{{下一张表?}} -->|是| FOR_EACH
NEXT_TABLE -->|否| CLOSE_ACC[关闭 Access 连接]
CLOSE_ACC --> SUMMARY[输出文件级汇总]
SUMMARY --> POLL
style VERIFY fill:#f9f,stroke:#333
style POST_VERIFY fill:#f9f,stroke:#333
style AUTO_CREATE fill:#bbf,stroke:#333
```
## 逐主键校验机制
传统的"数量对比"方法存在缺陷:少插和漏删的错误可能互相抵消。本系统采用**逐主键确认**方式:
```mermaid
flowchart LR
subgraph 输入
A[record_ids<br/>日志中的主键列表]
B[inserted_pk_set<br/>Access 实际读到的主键]
end
subgraph 删除校验
C[差集: record_ids - inserted_pk_set<br/>= 应删除的主键]
D[查询目标表<br/>这些主键是否还存在?]
C --> D
D -->|还存在任一条| E[❌ 校验失败]
D -->|全部不存在| F[✅ 删除校验通过]
end
subgraph 插入校验
G[查询目标表<br/>inserted_pk_set 是否都在?]
G -->|有任一条查不到| E
G -->|全部存在| H[✅ 插入校验通过]
end
A --> C
B --> C
B --> G
```
### 校验函数说明
| 函数 | 作用 |
|------|------|
| `fetch_existing_pks()` | 分批查询目标表返回实际存在的主键集合IN 子句每批不超过 900 个参数) |
| `verify_sync_result()` | 执行删除校验 + 插入校验,失败时抛出 `SyncVerificationError` |
| `_norm_key()` | 主键归一化为字符串,规避 Access/SQL Server 类型差异 |
| `_chunked()` | 将列表分批,避免超出 SQL Server 参数上限 |
## 故障恢复
```mermaid
stateDiagram-v2
[*] --> Pending: VBA 写入日志 Synced=0
Pending --> Syncing: 同步服务读取
Syncing --> Verified: 校验通过 + 提交
Verified --> Committed: Synced=1 标记生效
Committed --> [*]: 同步完成
Syncing --> Rollback: 校验失败 / 异常
Rollback --> Pending: 事务回滚 Synced=0<br/>下一轮自动重试
Committed --> Pending: 提交后复核失败<br/>Synced 撤回为 0
```
### 三层保障
1. **提交前校验** — 删/插完成后、标记 Synced=1 前,逐主键确认结果正确
2. **事务回滚** — 校验失败或异常时回滚事务,保持 `Synced=0`,下一轮自动重试
3. **提交后复核**`commit()` 后再查一次数据库确认数据已持久化(可通过 `ENABLE_POST_COMMIT_VERIFY` 开关控制)
## 轮询策略
```
while True:
has_work = process_sync_task()
if has_work:
sleep(0.1) # 有积压,快速重试
else:
sleep(30) # 无工作,标准间隔
```
有未处理数据时以 0.1 秒间隔快速处理积压;无数据时按 `POLL_INTERVAL` 休眠。
## 自动建表
当目标表在 SQL Server 中不存在时,系统根据 Access 表的列定义自动建表:
```
Access cursor.description → 类型映射 → CREATE TABLE 语句
```
| Access 类型 | SQL Server 类型 |
|-------------|----------------|
| `int` | `INT`(主键时追加 `IDENTITY(1,1) PRIMARY KEY` |
| `float` | `FLOAT` |
| `bool` | `BIT` |
| `datetime.datetime` | `DATETIME` |
| `decimal.Decimal` | `DECIMAL(p, s)` |
| `str` (size ≤ 4000) | `NVARCHAR(size)` |
| `str` (size > 4000) | `NVARCHAR(MAX)` |