Insert/Update 日志回读不到源行时,旧逻辑无条件降级为 Delete,在 ACE 引擎可见性延迟(批量插入约11s窗口)下误伤,导致该 Insert 的行不仅没进 SQL 反而被空打 Delete 并清理 Access 日志,数据永久丢失(8/4 事件根因: 氩弧焊.接收 16255-16259 缺失)。 改为按 Access 日志年龄(capture_defer_seconds=60)决定动作: - 年龄 < 阈值:DEFER,跳过不入队,Access 日志保留,下一轮 cycle 重读; - 年龄 >= 阈值:AGED-OUT,入队直接标 dead 待人工(usp_SyncApply 只选 pending,不会执行任何破坏性 SQL),归档表记 ProcessedOperateType=AgedOut。 零 schema/SQL 改动:cleanup 只删 Status=applied 的日志,DEFER 行不入队 则无 applied 行、AGED-OUT 标 dead 非 applied,两者 Access 日志均保留, 与既有 cleanup/apply/service 逻辑天然自洽。 详见 docs/plan-capture-defer-by-age.md 与 docs/incremental-sync-flow.md。
282 lines
12 KiB
Markdown
282 lines
12 KiB
Markdown
# 增量同步流程详解(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 开始<br/>分配 cycle_id]) --> CAP
|
||
|
||
CAP["【1. Capture 捕获】<br/>逐库读 TableChangeLog → 回读整行 → 入队 SyncQueue<br/><i>src/sync/capture.py</i>"]
|
||
CAP --> APPLY
|
||
|
||
APPLY["【2. Apply 应用】<br/>调 usp_SyncApply 把 SyncQueue pending 行落到镜像表<br/><i>sql/02_sync_apply.sql</i>"]
|
||
APPLY --> HEALTH
|
||
|
||
HEALTH["【2.5 队列健康检查】<br/>查 error/dead 卡死行并告警"]
|
||
HEALTH --> CLEAN
|
||
|
||
CLEAN["【3. Cleanup 清理】<br/>删 Access 已应用日志 → SyncQueue 行标 cleaned<br/><i>src/sync/cleanup.py</i>"]
|
||
CLEAN --> PURGE
|
||
|
||
PURGE["【3.5 Purge 回收】<br/>删 SyncQueue 中超保留期的 cleaned 行"]
|
||
PURGE --> DONE([cycle 结束<br/>休眠 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<br/>最旧 N 条(按 ID 升序)<br/>N = capture_batch_size=500"]
|
||
B --> C{有日志行?}
|
||
C -- 否 --> Z([capture 结束])
|
||
C -- 是 --> D[逐行处理]
|
||
|
||
D --> E{"表是否在同步范围内?<br/>is_synced_table"}
|
||
E -- 否 --> F["跳过(out_of_scope)<br/>不计入,下轮仍会读到"]
|
||
F --> D
|
||
E -- 是 --> G{"OperateType?"}
|
||
|
||
G -- Insert/Update --> H["🔑 回读整行<br/>read_row(table, record_id)<br/>SELECT * FROM 表 WHERE ID=?"]
|
||
H --> I{读到行?}
|
||
I -- 是 --> J["row_data = JSON 序列化<br/>op 保持 Insert/Update"]
|
||
I -- ❌否 --> K["⚠️ 降级 op = Delete<br/>row_data = None<br/>写 DOWNGRADE 警告日志"]
|
||
|
||
G -- Delete --> L["op = Delete<br/>row_data = None<br/>(不回读,Delete 无需数据)"]
|
||
|
||
G -- 其它未知 --> M["跳过(unknown_op)<br/>下轮仍会读到"]
|
||
|
||
J --> N["写 SyncLogArchive(永久审计)<br/>记录 OriginalOperateType + ProcessedOperateType"]
|
||
K --> N
|
||
L --> N
|
||
|
||
N --> O["入队 SyncQueue(去重插入)<br/>INSERT...WHERE NOT EXISTS"]
|
||
O --> P{插入成功?}
|
||
P -- 是 --> Q["enqueued +1"]
|
||
P -- 否(去重命中)--> R["dedup_skipped +1<br/>说明上轮 apply 失败残留"]
|
||
|
||
Q --> S{还有下一行?}
|
||
R --> S
|
||
F --> S
|
||
M --> S
|
||
S -- 是 --> D
|
||
S -- 否 --> T["打印 capture summary<br/>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<br/>EXEC usp_SyncApply]) --> B["重置 error 行<br/>RetryCount < max_retries 的 → pending"]
|
||
B --> C["按 TargetSchema+TargetTable 分组<br/>遍历每个目标表"]
|
||
|
||
C --> D{该表有 pending 行?}
|
||
D -- 否 --> C
|
||
D -- 是 --> E["统计 pending 数 / distinct RecordID 数"]
|
||
|
||
E --> F["BEGIN TRAN"]
|
||
|
||
F --> G["🔑 构建列清单<br/>从 sys.columns 读目标表所有列<br/>(排除 ID/computed/identity/timestamp)"]
|
||
|
||
G --> H["构建动态 SQL"]
|
||
|
||
H --> I["分支1: Upsert(MERGE)<br/>ranked CTE: 按 RecordID 分区,<br/>SourceLogID DESC 取 rn=1<br/>仅 OperateType∈Insert/Update 且 RowData 非空"]
|
||
I --> J["SET IDENTITY_INSERT ON<br/>MERGE 目标表<br/>匹配则 UPDATE, 不匹配则 INSERT<br/>@merged = @@ROWCOUNT"]
|
||
|
||
J --> K["分支2: Delete<br/>同一 ranked CTE 的 rn=1 行<br/>仅 OperateType=Delete"]
|
||
K --> L["DELETE 目标表 WHERE ID IN (...)<br/>@deleted = @@ROWCOUNT"]
|
||
|
||
L --> M["把该表所有 pending 行<br/>Status → applied, AppliedAt = now<br/>@applied = @@ROWCOUNT"]
|
||
|
||
M --> N[COMMIT]
|
||
|
||
N --> O["写 SyncApplyRunLog 审计<br/>pending/merged/deleted/applied/<br/>error/dead + CycleID + 耗时"]
|
||
O --> C
|
||
|
||
N -.失败.-> X["ROLLBACK"]
|
||
X --> Y["超 max_retries → dead<br/>否则 → 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 中<br/>Status=applied 的 SourceLogID 列表<br/>applied_log_ids(file)"]
|
||
B --> C{有 applied 行?}
|
||
C -- 否 --> Z([cleanup 结束, 返回 0])
|
||
C -- 是 --> D["DELETE FROM Access.TableChangeLog<br/>WHERE ID IN (上述列表)<br/>分批 + 锁重试"]
|
||
D --> E{删除数 == 预期?}
|
||
E -- 否 --> F["WARNING: 部分日志已不在<br/>(被外部或中断的运行删过)"]
|
||
E -- 是 --> G["mark_cleaned:<br/>这些队列行 Status → cleaned<br/>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[(业务表<br/>如 接收)]
|
||
TCL[(TableChangeLog<br/>变更日志 append-only)]
|
||
end
|
||
|
||
subgraph SQL["SQL Server 端(CompanyDB)"]
|
||
SQ[(SyncQueue<br/>待处理队列)]
|
||
ARCH[(SyncLogArchive<br/>永久审计)]
|
||
RL[(SyncApplyRunLog<br/>apply 运行日志)]
|
||
MIRROR[(镜像表<br/>如 接收_YEAR2026)]
|
||
end
|
||
|
||
T1 -- "数据宏 After I/U/D<br/>写入一行日志" --> 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 行<br/>数据宏写 5 条 Insert 日志<br/>SourceLogID 15533-15537"] --> B
|
||
|
||
B["09:00:50 Capture 读取这 5 条日志<br/>(注: CapturedAt 比 OriginalTime 早 11s<br/> = ACE 可见性窗口)"] --> C
|
||
|
||
C["read_row 回读 5 行<br/>SELECT * FROM 接收 WHERE ID=16255..16259"] --> D
|
||
|
||
D["❌ 全部返回 None<br/>(行尚未对其他连接可见)"] --> E
|
||
|
||
E["🔴 降级为 Delete<br/>入队 SyncQueue, OperateType=Delete<br/>RowData=null"] --> F
|
||
|
||
F["09:00:51 Apply: Delete 分支<br/>DELETE 接收_YEAR2026 WHERE ID IN(16255..16259)<br/>@deleted=0(SQL 本就没这5行)<br/>@applied=5(队列行标 applied)"] --> G
|
||
|
||
G["队列行 Status = applied"] --> H
|
||
|
||
H["Cleanup: 查到这5条 applied<br/>DELETE Access.TableChangeLog ID=15533-15537"] --> I
|
||
|
||
I["🔴 日志被删,5 条 Insert 证据消失<br/>这5行从此再无机会被同步进 SQL"] --> J
|
||
|
||
J["8/4 05:08 compare: missing_in_sql=5<br/>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 改动、零状态存储。等你审完这份流程图确认方向后,我再动手。
|