diff --git a/docs/full-sync.md b/docs/full-sync.md new file mode 100644 index 0000000..6e05d06 --- /dev/null +++ b/docs/full-sync.md @@ -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
成品入库.accdb"] + A2["Access 文件 2
25年压力表合同数据.accdb"] + A3["Access 文件 N
..."] + 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
中的每个 Access 文件}} + + FILE_LOOP --> FILE_CHECK{文件是否存在?} + FILE_CHECK -->|不存在| SKIP[跳过该文件] + FILE_CHECK -->|存在| ACC_CONN[连接 Access 数据库] + + ACC_CONN --> TABLE_LOOP{{遍历该文件下的
每个表映射}} + + TABLE_LOOP --> ACC_READ[读取 Access 表结构
SELECT TOP 1 * FROM table] + + ACC_READ --> TARGET_CHECK{目标表是否存在?} + + TARGET_CHECK -->|不存在| AUTO_CREATE[自动建表
根据 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[批量数据传输
fetchmany BATCH_SIZE] + DATA_TRANSFER --> HAS_MORE{{还有数据?}} + + HAS_MORE -->|有| BATCH[executemany 插入一个批次
记录进度日志] + BATCH --> HAS_MORE + + HAS_MORE -->|无| ID_OFF_CHECK{IDENTITY_INSERT
是否已开启?} + ID_OFF_CHECK -->|是| ID_OFF[SET IDENTITY_INSERT OFF] + ID_OFF_CHECK -->|否| COMMIT + ID_OFF --> COMMIT[提交事务] + + COMMIT --> STATS[记录表级统计
行数 / 速率 / 用时] + + 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[记录失败日志
回滚事务
清理 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 行
记录一次进度 + 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
类型映射] + B --> C["CREATE TABLE [schema].[table] (
[col1] INT IDENTITY(1,1) PRIMARY KEY,
[col2] NVARCHAR(100),
...
)"] + C --> D[DDL 单独提交
确保后续参数绑定正常] +``` + +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: 会话级设置
不随事务回滚 +``` + +关键点: +- `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` 驱动,新增表映射后首次运行全量同步即可自动建表并填充数据。 diff --git a/docs/incremental-sync.md b/docs/incremental-sync.md new file mode 100644 index 0000000..a0e1f0b --- /dev/null +++ b/docs/incremental-sync.md @@ -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
WHERE Synced=0] + + POLL -->|无记录| SLEEP[休眠 POLL_INTERVAL 秒] + SLEEP --> POLL + + POLL -->|发现未同步记录| GROUP[按 TableName 分组
合并同一表的 record_ids] + + GROUP --> CONNECT[连接 Access 数据库] + + CONNECT --> FOR_EACH{{遍历每个表}} + + FOR_EACH --> CHECK_TABLE{目标表是否存在?} + + CHECK_TABLE -->|不存在| AUTO_CREATE[自动建表
ensure_schema + create_table_from_access] + CHECK_TABLE -->|存在| READ_ACCESS + AUTO_CREATE --> READ_ACCESS + + READ_ACCESS[A. 从 Access 读取最新数据
SELECT WHERE PK IN ...] + + READ_ACCESS --> EXTRACT_PK[提取 inserted_pk_set
Access 实际返回的主键集合] + + EXTRACT_PK --> DELETE[B. 删除目标表旧记录
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. 批量插入新记录
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{提交后复核
ENABLE_POST_COMMIT_VERIFY} -->|关闭| SUCCESS + POST_VERIFY -->|开启| POST_OK{复核通过?} + POST_OK -->|通过| SUCCESS[记录同步成功日志] + POST_OK -->|失败| REVERT[回滚 Synced=0
下一轮重试] + + VERIFY_OK -->|失败| ROLLBACK[回滚事务
保持 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
日志中的主键列表] + B[inserted_pk_set
Access 实际读到的主键] + end + + subgraph 删除校验 + C[差集: record_ids - inserted_pk_set
= 应删除的主键] + D[查询目标表
这些主键是否还存在?] + C --> D + D -->|还存在任一条| E[❌ 校验失败] + D -->|全部不存在| F[✅ 删除校验通过] + end + + subgraph 插入校验 + G[查询目标表
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
下一轮自动重试 + + Committed --> Pending: 提交后复核失败
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)` |