# 增量同步流程详解(Incremental Sync Flow)
> 本文档梳理 Access → SQL Server 增量同步的完整流程,逐节点说明「发生了什么、读了/写了什么、状态如何流转」,作为重构「Insert 降级 Delete」逻辑的决策依据。
>
> 涉及代码:`src/sync/service.py`(主循环 cycle)、`src/sync/capture.py`(捕获)、`sql/02_sync_apply.sql`(应用)、`src/sync/cleanup.py`(清理)、`src/sync/sql_writer.py`(SQL 端读写)、`src/sync/access_reader.py`(Access 端读取)。
---
## 一、整体架构:一轮 cycle 的三段流水线
每个 cycle(默认间隔 `poll_interval_seconds=10s`)跑一遍三段,顺序固定、不可调换:
```mermaid
flowchart TD
START([cycle 开始
分配 cycle_id]) --> CAP
CAP["【1. Capture 捕获】
逐库读 TableChangeLog → 回读整行 → 入队 SyncQueue
src/sync/capture.py"]
CAP --> APPLY
APPLY["【2. Apply 应用】
调 usp_SyncApply 把 SyncQueue pending 行落到镜像表
sql/02_sync_apply.sql"]
APPLY --> HEALTH
HEALTH["【2.5 队列健康检查】
查 error/dead 卡死行并告警"]
HEALTH --> CLEAN
CLEAN["【3. Cleanup 清理】
删 Access 已应用日志 → SyncQueue 行标 cleaned
src/sync/cleanup.py"]
CLEAN --> PURGE
PURGE["【3.5 Purge 回收】
删 SyncQueue 中超保留期的 cleaned 行"]
PURGE --> DONE([cycle 结束
休眠 poll_interval])
style CAP fill:#e3f2fd,stroke:#1976d2
style APPLY fill:#fff3e0,stroke:#f57c00
style HEALTH fill:#fce4ec,stroke:#c2185b
style CLEAN fill:#e8f5e9,stroke:#388e3c
style PURGE fill:#f3e5f5,stroke:#7b1fa2
```
**关键设计约束**(决定重构可行性的红线):
- **Access 的 `TableChangeLog` 是 append-only,无状态字段**——它只是个待处理队列,无法在上面记录「已重试几次」。
- **日志清除的唯一依据是 SyncQueue 的 `Status='applied'`**——只要一条日志对应的队列行不是 applied,cleanup 就不会删它(详见第三节)。
- **SyncQueue 有唯一索引 `(SourceFile,SourceTable,SourceLogID)` 去重**——同一日志第二次入队会被静默跳过,不会覆盖原行、不会自增计数。
---
## 二、Capture 阶段(数据捕获)—— ❗重构的核心战场
逐库处理,每个 Access 文件独立隔离(单库失败不影响其它)。
```mermaid
flowchart TD
A([开始 capture 某个文件]) --> B["读 Access TableChangeLog
最旧 N 条(按 ID 升序)
N = capture_batch_size=500"]
B --> C{有日志行?}
C -- 否 --> Z([capture 结束])
C -- 是 --> D[逐行处理]
D --> E{"表是否在同步范围内?
is_synced_table"}
E -- 否 --> F["跳过(out_of_scope)
不计入,下轮仍会读到"]
F --> D
E -- 是 --> G{"OperateType?"}
G -- Insert/Update --> H["🔑 回读整行
read_row(table, record_id)
SELECT * FROM 表 WHERE ID=?"]
H --> I{读到行?}
I -- 是 --> J["row_data = JSON 序列化
op 保持 Insert/Update"]
I -- ❌否 --> K["⚠️ 降级 op = Delete
row_data = None
写 DOWNGRADE 警告日志"]
G -- Delete --> L["op = Delete
row_data = None
(不回读,Delete 无需数据)"]
G -- 其它未知 --> M["跳过(unknown_op)
下轮仍会读到"]
J --> N["写 SyncLogArchive(永久审计)
记录 OriginalOperateType + ProcessedOperateType"]
K --> N
L --> N
N --> O["入队 SyncQueue(去重插入)
INSERT...WHERE NOT EXISTS"]
O --> P{插入成功?}
P -- 是 --> Q["enqueued +1"]
P -- 否(去重命中)--> R["dedup_skipped +1
说明上轮 apply 失败残留"]
Q --> S{还有下一行?}
R --> S
F --> S
M --> S
S -- 是 --> D
S -- 否 --> T["打印 capture summary
read/enqueued/downgraded/..."]
T --> Z
style K fill:#ffcdd2,stroke:#c62828,stroke-width:3px
style H fill:#fff9c4,stroke:#f9a825
style I fill:#fff9c4,stroke:#f9a825
```
### 🔴 问题节点:降级 Delete(`capture.py:93-108`)
```python
if op in ("Insert", "Update"):
d = reader.read_row(lr.table_name, lr.record_id)
if d is None:
op = "Delete" # ← 问题根源:读不到就降级
st.downgraded += 1
log.warning("capture DOWNGRADE %s->Delete ...")
```
**为什么读不到?两个场景无法区分:**
1. **真删除**:行被 Insert 后又 Delete(Access 客户端先插后删)→ 这时降级 Delete 是「碰巧正确」。
2. **可见性延迟**(本次事件的根因):行已插入但 ACE 引擎尚未对其他 ODBC 连接可见(批量插入时窗口可达 11 秒)→ 这时降级 Delete 是**有害的误伤**。
**代码当前无法区分这两种情况**,统一降级为 Delete。这就是要重构的核心。
---
## 三、Apply 阶段(数据应用)
调用存储过程 `usp_SyncApply`,**按表分组、集合化处理**所有 pending 行。
```mermaid
flowchart TD
A([call_apply
EXEC usp_SyncApply]) --> B["重置 error 行
RetryCount < max_retries 的 → pending"]
B --> C["按 TargetSchema+TargetTable 分组
遍历每个目标表"]
C --> D{该表有 pending 行?}
D -- 否 --> C
D -- 是 --> E["统计 pending 数 / distinct RecordID 数"]
E --> F["BEGIN TRAN"]
F --> G["🔑 构建列清单
从 sys.columns 读目标表所有列
(排除 ID/computed/identity/timestamp)"]
G --> H["构建动态 SQL"]
H --> I["分支1: Upsert(MERGE)
ranked CTE: 按 RecordID 分区,
SourceLogID DESC 取 rn=1
仅 OperateType∈Insert/Update 且 RowData 非空"]
I --> J["SET IDENTITY_INSERT ON
MERGE 目标表
匹配则 UPDATE, 不匹配则 INSERT
@merged = @@ROWCOUNT"]
J --> K["分支2: Delete
同一 ranked CTE 的 rn=1 行
仅 OperateType=Delete"]
K --> L["DELETE 目标表 WHERE ID IN (...)
@deleted = @@ROWCOUNT"]
L --> M["把该表所有 pending 行
Status → applied, AppliedAt = now
@applied = @@ROWCOUNT"]
M --> N[COMMIT]
N --> O["写 SyncApplyRunLog 审计
pending/merged/deleted/applied/
error/dead + CycleID + 耗时"]
O --> C
N -.失败.-> X["ROLLBACK"]
X --> Y["超 max_retries → dead
否则 → error(下轮重试)"]
Y --> O
style I fill:#e3f2fd,stroke:#1976d2
style K fill:#ffcdd2,stroke:#c62828
style M fill:#fff3e0,stroke:#f57c00
```
### 关键:保序「最后操作胜」(`02_sync_apply.sql:95-135`)
单个 `ranked` CTE 同时供 Upsert 和 Delete 两个分支使用:
```sql
ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn
```
- `rn=1` 是该 RecordID **真正的最后一条日志**。
- Upsert 分支:`rn=1 AND OperateType IN ('Insert','Update')`
- Delete 分支:`rn=1 AND OperateType='Delete'`
所以**降级成 Delete 的行,在这里会真的去 SQL 端执行 DELETE**。本次事件中 5 行 Delete 的 `@deleted=0`(SQL 里本就没这些行,空打),但 `@applied=5`(队列行照样被标 applied)。
---
## 四、Cleanup 阶段(日志清除)—— ❗决定「重试」能否成立的命脉
```mermaid
flowchart TD
A([cleanup 某个文件]) --> B["查 SyncQueue 中
Status=applied 的 SourceLogID 列表
applied_log_ids(file)"]
B --> C{有 applied 行?}
C -- 否 --> Z([cleanup 结束, 返回 0])
C -- 是 --> D["DELETE FROM Access.TableChangeLog
WHERE ID IN (上述列表)
分批 + 锁重试"]
D --> E{删除数 == 预期?}
E -- 否 --> F["WARNING: 部分日志已不在
(被外部或中断的运行删过)"]
E -- 是 --> G["mark_cleaned:
这些队列行 Status → cleaned
CleanedAt = now"]
F --> G
G --> H["INFO: 删除了 N 条 Access 日志"]
H --> Z
style D fill:#e8f5e9,stroke:#388e3c,stroke-width:2px
style B fill:#fff9c4,stroke:#f9a825
```
### 🔑 cleanup 的判定条件是重构的支点
**cleanup 只删 `Status='applied'` 的日志**(`sql_writer.py:264-272` 的 `applied_log_ids`)。这意味着:
| capture 对该日志的处理 | SyncQueue 行状态 | cleanup 是否删 Access 日志 | 后果 |
|------------------------|------------------|----------------------------|------|
| 降级 Delete(现状) | applied | **删除** | ❌ 日志消失,再无重试机会 |
| **跳过不入队(重构后)** | (无对应行) | **不删** | ✅ 日志保留,下轮重读 |
**结论:重构只要做到「读不到 → 不入队」,cleanup 这一段天然会把日志保留下来,无需改动 cleanup.py。** 这是「跨 cycle 重试」能够成立的根基。
---
## 五、数据流转全景:一条日志的完整生命周期
```mermaid
flowchart LR
subgraph Access["Access 端(.accdb)"]
T1[(业务表
如 接收)]
TCL[(TableChangeLog
变更日志 append-only)]
end
subgraph SQL["SQL Server 端(CompanyDB)"]
SQ[(SyncQueue
待处理队列)]
ARCH[(SyncLogArchive
永久审计)]
RL[(SyncApplyRunLog
apply 运行日志)]
MIRROR[(镜像表
如 接收_YEAR2026)]
end
T1 -- "数据宏 After I/U/D
写入一行日志" --> TCL
TCL -- "① capture 读取" --> CAP[Capture]
CAP -- "回读整行" --> T1
CAP -- "② 入队(去重)" --> SQ
CAP -- "② 审计存档" --> ARCH
SQ -- "③ apply 处理" --> APPLY[usp_SyncApply]
APPLY -- "MERGE/DELETE" --> MIRROR
APPLY -- "pending→applied" --> SQ
APPLY -- "记录运行结果" --> RL
SQ -- "④ cleanup 查 applied" --> CLEAN[Cleanup]
CLEAN -- "⑤ 删已应用日志" --> TCL
CLEAN -- "applied→cleaned" --> SQ
style CAP fill:#e3f2fd,stroke:#1976d2
style APPLY fill:#fff3e0,stroke:#f57c00
style CLEAN fill:#e8f5e9,stroke:#388e3c
```
---
## 六、本次事件的重放(在上述流程中的路径)
5 条 Insert 日志(RecordID 16255-16259)在一轮 cycle 中的遭遇:
```mermaid
flowchart TD
A["8/3 09:01 Access 批量插入 5 行
数据宏写 5 条 Insert 日志
SourceLogID 15533-15537"] --> B
B["09:00:50 Capture 读取这 5 条日志
(注: CapturedAt 比 OriginalTime 早 11s
= ACE 可见性窗口)"] --> C
C["read_row 回读 5 行
SELECT * FROM 接收 WHERE ID=16255..16259"] --> D
D["❌ 全部返回 None
(行尚未对其他连接可见)"] --> E
E["🔴 降级为 Delete
入队 SyncQueue, OperateType=Delete
RowData=null"] --> F
F["09:00:51 Apply: Delete 分支
DELETE 接收_YEAR2026 WHERE ID IN(16255..16259)
@deleted=0(SQL 本就没这5行)
@applied=5(队列行标 applied)"] --> G
G["队列行 Status = applied"] --> H
H["Cleanup: 查到这5条 applied
DELETE Access.TableChangeLog ID=15533-15537"] --> I
I["🔴 日志被删,5 条 Insert 证据消失
这5行从此再无机会被同步进 SQL"] --> J
J["8/4 05:08 compare: missing_in_sql=5
ID 16255-16259"]
style D fill:#ffcdd2,stroke:#c62828
style E fill:#ffcdd2,stroke:#c62828,stroke-width:3px
style H fill:#e8f5e9,stroke:#388e3c
style I fill:#ffcdd2,stroke:#c62828
style J fill:#fff9c4,stroke:#f9a825
```
**链式灾难**:降级 Delete(capture)→ 空打 Delete 但标 applied(apply)→ 删 Access 日志(cleanup)。三段配合下,这 5 行被「合法地」从同步链路中抹除。只要在 capture 段断开第一环(不降级、不入队),后面 apply/cleanup 就不会碰它们,日志保留,下轮自然重试。
---
## 七、重构决策点(待定)
基于上述流程,重构的核心是改造 capture 阶段的「降级」分支。需要决策的问题:
1. **读不到时的动作**:跳过不入队(日志保留,下轮重读)vs 入队但标记 defer 状态?
2. **重试上限的判定依据**:用「重试次数」(需要存储计数)vs 用「日志年龄时间窗」(无需存储,靠 `now - OriginalTime > 阈值`)?
3. **终态处理**:超限后该 Access 日志保留(占队列头持续重读)还是删除(丢证据)?
4. **Update 日志**:是否套用同一套 defer 逻辑?
我倾向的方案:**读不到 → 不入队 + 写 defer 审计;用日志年龄(如 120s)作终态判据;超限则入队标 dead(保留 Access 日志不删,待人工)**。零 schema 改动、零状态存储。等你审完这份流程图确认方向后,我再动手。