Files
ProductionDataBaseSync_Data…/docs/plan-capture-defer-by-age.md
Misaka_Company 7cef5153e6 fix(capture): 读不到源行时按日志年龄延迟重试,根治降级Delete致数据丢失
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。
2026-08-05 11:00:58 +08:00

11 KiB
Raw Permalink Blame History

重构方案capture 读不到行时按「日志年龄」延迟重试

目标根治「Insert 日志回读不到行 → 降级 Delete → 数据丢失」的缺陷8/4 事件根因)。 核心改动:读不到 → 不入队、不降级,按 Access 日志的年龄决定「下轮重试」还是「判死待人工」。 改动面:capture.py(主)、config.py新增2个配置项capture.py 的 CaptureStats新增计数器零 schema 改动、零 SQL 改动、不动 apply/cleanup。


一、当前问题行为 vs 重构后行为

场景 当前行为(缺陷) 重构后行为
Insert/Update 读不到行(可见性延迟,本次事件) 降级 Delete → 入队 → apply 空打 Delete + 标 applied → cleanup 删 Access 日志 → 数据永久丢失 年龄 < 阈值:跳过不入队,日志保留,下轮重读 → 读到后正常 Insert
Insert/Update 读不到行(真删除:先插后删) 降级 Delete碰巧正确 年龄 < 阈值:跳过;后续 Delete 日志会兜底正确清理;超龄判死待人工
Insert/Update 一直读不到(真异常/数据损坏) 降级 Delete错误 年龄 ≥ 阈值:入队标 dead保留 Access 日志,告警待人工

二、配置项新增(src/sync/config.py 的 RuntimeConfig

class RuntimeConfig(BaseModel):
    # ... 既有字段 ...
    # Insert/Update 日志回读不到行时,按日志年龄延迟重试:年龄小于此秒数则
    # 跳过不入队(保留 Access 日志,下一轮 cycle 重新捕获),度过 ACE 引擎
    # 的可见性窗口;超过此年龄仍读不到则判定为真删除/异常,入队标 dead 待人工。
    # 设为 0 可关闭延迟重试(退化为旧的"立即判死"语义,但不再降级 Delete
    capture_defer_seconds: int = 120
  • 默认值 120 秒:覆盖本次事件观察到的 ~11 秒可见性窗口,并留足 10 倍余量;对 poll_interval=10s 意味着最多重试约 12 轮。
  • 单一配置项,config.yaml 无需改动即可生效(用默认值)。defer_dead_op 不暴露为配置(实现细节,固定为 dead 状态,见下)。

三、capture.py 改动(核心)

3.1 CaptureStats 新增两个计数器

@dataclass
class CaptureStats:
    # ... 既有字段 ...
    read: int = 0
    enqueued: int = 0
    dedup_skipped: int = 0
    downgraded: int = 0      # 保留字段,重构后恒为 0兼容旧日志解析
    deferred: int = 0        # 【新】读不到行但年龄 < 阈值,跳过待下轮重试
    aged_out: int = 0        # 【新】读不到行且年龄 ≥ 阈值,判死待人工
    out_of_scope: int = 0
    unknown_op: int = 0
    # ... merge() / ops_str() 同步更新 ...

3.2 降级分支重构(capture_file 中 Insert/Update 回读逻辑)

替换 现有的第 93-108 行(整段 if op in ("Insert", "Update"): 块):

if op in ("Insert", "Update"):
    d = reader.read_row(lr.table_name, lr.record_id)
    if d is None:
        # 读不到行:不再降级 Delete。按 Access 日志年龄决定延迟重试还是判死。
        # 必须计算日志年龄(原始设计无此逻辑)。
        age_seconds = _log_age_seconds(lr.time)
        if age_seconds < cfg.runtime.capture_defer_seconds:
            # 暂态读不到ACE 可见性窗口):跳过,不入队、不清理。
            # Access 日志因无对应 applied 队列行cleanup 不会删除,下轮重试。
            st.deferred += 1
            log.warning(
                "capture DEFER %s file=%s table=%s record_id=%s log_id=%s "
                "log_time=%s age=%ds (source row unreadable; will retry next "
                "cycle while log age < %ds)",
                lr.operate_type, fm.file, lr.table_name, lr.record_id,
                lr.id, lr.time, int(age_seconds),
                cfg.runtime.capture_defer_seconds,
            )
            continue   # ← 关键跳过本行archive/queue 都不写
        else:
            # 超龄仍读不到:真删除或真异常。入队标 dead不降级、不执行任何
            # 破坏性操作。保留 Access 日志(无 applied 行 → cleanup 不删),
            # 队列健康检查会告警,等待人工介入。
            st.aged_out += 1
            log.warning(
                "capture AGED-OUT %s file=%s table=%s record_id=%s log_id=%s "
                "log_time=%s age=%ds >= %ds -- enqueuing as dead for manual "
                "review (source row still unreadable after defer window)",
                lr.operate_type, fm.file, lr.table_name, lr.record_id,
                lr.id, lr.time, int(age_seconds),
                cfg.runtime.capture_defer_seconds,
            )
            op = "dead"  # 仅用于入队时的状态标记,见下
    else:
        row_data = json.dumps(d, ensure_ascii=False)

3.3 超龄行的入队方式(aged_out 分支)

超龄行需要进 SyncQueue 但不能被 apply 执行任何 SQL 操作(没数据可插,也不能 Delete。两种实现可选我倾向 A

方案 A推荐入队后直接标 deadOperateType 保留原 Insert/Update 真相

  • 入队时 OperateType 仍写 Insert/Update(保留原始意图,便于审计),RowData=null
  • 入队后立即 UPDATE ... SET Status='dead', ErrorMsg='source row unreadable after {age}s defer'
  • apply 的游标只选 Status='pending'dead 行不会被处理 → 不会误删。
  • queue 健康检查已有 dead 告警(service.py:166-184),自动浮现。

方案 B新增 OperateType='Noop' — 改动面更大apply 存储过程需识别),不推荐。

需要在 sql_writer.py 新增一个方法 insert_dead_row(row, error_msg),逻辑 = 先 insert_queue_row(去重插入)再 UPDATE ... SET Status='dead', RetryCount=<对应>, ErrorMsg=?

3.4 日志年龄计算辅助函数 _log_age_seconds

import datetime as _dt

def _log_age_seconds(log_time: object) -> float:
    """Access 日志行 Time 字段距今的秒数。log_time 是 pyodbc 返回的 datetime。
    异常时返回一个大数(视为已超龄),确保宁可判死也不无限重试。"""
    try:
        if isinstance(log_time, _dt.datetime):
            return (_dt.datetime.now() - log_time).total_seconds()
        # Access via ODBC 通常返回 datetime兜底处理 naive/其它类型
        return float("inf")
    except Exception:
        return float("inf")

⚠️ 时区/时钟注意lr.time 是 Access 端写入的本地时间,datetime.now() 也是本机本地时间,两者同在 114 主机同一时区,可直接相减。本次事件中 OriginalTime/CapturedAt 的"倒挂"现象差11秒不影响此逻辑——因为按年龄判断即便 lr.time 偏早age 只会被算得更大,倾向判死而非误伤,方向安全。


四、归档表SyncLogArchive的处理

分支 是否写 archive 理由
DEFER跳过 不写 日志保留在 Access下轮 capture 会重新读到并正常归档;此时写 archive 反而会在去重表里留下"读不到"的半成品记录
AGED-OUT判死 超龄是终态需留永久审计OriginalOperateType=Insert/Update, ProcessedOperateType='AgedOut', RowData=null, OriginalTime=lr.time

ProcessedOperateType 新增值 'AgedOut'(仅 archive 表用varchar(10) 放得下 7 字符)。这是纯审计标记,不影响任何执行逻辑。


五、不改动的地方(明确边界)

模块 是否改动 原因
cleanup.py 不改 只删 Status='applied' 的日志DEFER 行不入队无 applied 行 → 日志保留AGED-OUT 标 dead 非 applied → 日志也保留。天然自洽。
sql/02_sync_apply.sql 不改 apply 游标只选 pendingdead 行天然跳过DEFER 行根本不入队。
sql/01_sync_queue.sql 不改 不新增列,不改索引。
service.py 不改 队列健康检查已有 error/dead 告警(queue_error_samplesAGED-OUT 的 dead 行会自动被它捕获并 WARNING。capture summary 日志格式已包含新计数器(由 CaptureStats.ops_str/merge 驱动)。
access_reader.py 不改 read_row 行为不变。
config.yaml 不改 用默认值 120s 即可。

六、本次事件 5 行的重放(重构后)

cycle N (09:00:50): capture 读到 5 条 Insert 日志
  → read_row 返回 None
  → age = now(09:00:50) - log_time(09:01:01) → 注: 因 Access 时间戳特性 age 可能算成负或小
  → 即便按最保守计算age 远 < 120s
  → DEFER跳过不入队写 WARNINGAccess 日志保留

cycle N+1 (09:01:00): capture 再次读到这 5 条日志
  → read_row 此刻可见性窗口已过11s > 窗口)→ 读到行 ✅
  → 正常入队 OperateType=Insert带完整 RowData
  → apply MERGE → SQL 正确写入 5 行 ✅
  → cleanup 删 Access 日志(这次是 applied合理

结果8/4 compare 不再出现 missing_in_sql=5。


七、验证计划

  1. 单元测试tests/test_capture.py
    • mock read_row 返回 None + lr.time 为近时 → 断言 deferred=1, enqueued=0, aged_out=0,不调用 insert_queue_row / insert_archive_row
    • mock read_row 返回 None + lr.time 为 200s 前 → 断言 aged_out=1,调用 insert_dead_rowStatus=dead。
    • mock read_row 返回 dict + 任意时间 → 断言正常入队(回归测试)。
    • capture_defer_seconds=0 → 任何读不到都立即判死(边界)。
  2. 现有测试回归pytest 全绿(确保去重/正常 Insert/真 Delete 路径不受影响)。
  3. 集成验证(部署后观察 1-2 天):
    • 关注 capture summary 日志的 deferred= 计数,确认批量插入场景下有 defer 发生且下轮 enqueued。
    • 关注 queue health WARNING确认 dead 行(如有)被正确告警。
    • compare --granularity ids,确认无 missing_in_sql。

八、改动文件清单

文件 改动类型 说明
src/sync/config.py 新增字段 RuntimeConfig 加 capture_defer_seconds: int = 120
src/sync/capture.py 核心重构 降级分支 → defer/aged_out 分支CaptureStats 加 2 计数器;新增 _log_age_seconds
src/sync/sql_writer.py 新增方法 insert_dead_row(row, error_msg):去重插入后立即标 dead
tests/test_capture.py 新增用例 覆盖 DEFER / AGED-OUT / 正常 / 边界 4 种情况

总计 4 个文件,零 SQL/零 schema 改动。


九、待你确认的决策点

  1. capture_defer_seconds 默认值 120s 是否合适覆盖11s窗口×10倍余量
  2. AGED-OUT 行的处理:入队标 dead方案A推荐vs 其它?
  3. archive 表 ProcessedOperateType 新增值 'AgedOut' 是否可接受?
  4. downgraded 计数器保留为恒0向后兼容旧日志解析还是直接删除

确认后我即按此方案执行。