每天都全量拉取,却每次都少几行
目标
制作一个把不断变化的源分页拉取的同步程序。亲手重现偏移量方式为什么会悄悄漏掉行,换成 (updated_at, id) 游标,用水位线从中断点续传,最后对照源和副本。
为什么重要
在把 60 万条分成每次 1000 条、共 600 次拉取的期间,源也在不断变化。已经翻过去的页中的行被更新、在排序顺序中跑到最后面,后面的行就会向前移动,它们之间的行就在没有被任何人读取的情况下过去了。
这类事故不会报错,条数大致对得上,而且每次漏掉的行都不一样。所以无法重现,几个月之后会以投诉的形式回来。
修复的办法,是用值而不是编号来确定位置。不过,如果排序键不唯一,还有另一个陷阱在等着——只用 updated_at 一个值,就会丢失相同时刻的行,或者无限地停在同一个位置打转。所以游标要增加栏位直到唯一为止。
评分器不会相信你写的句子。评分器会把源服务器启动在评分器所选的端口上,一边修改源,一边实际运行你的同步程序,统计漏掉了什么、重叠了什么。
步骤
- 创建 /root/sync/source.py 并在端口 8019 上启动,把全部数据做成快照保存到 /root/sync/snapshot.json。
- 创建 /root/sync/offset_sync.py,用偏移量遍历,同时在途中修改源,用数字展示遗漏与重复同时出现。
- 创建 /root/sync/cursor_sync.py,用
(updated_at, id)游标遍历,使在同样的情况下遗漏变为 0。 - 给 cursor_sync.py 加上
--state,把水位线保存到文件中,让第二次运行只拉取有变化的部分。 - 即使让页恰好切在拥有相同
updated_at的行的正中间,也没有遗漏和重复,把这一点写入 /root/sync/tie_report.json。 - 接续被中断的同步,把全部数据汇集到 /root/sync/sink.jsonl。
- 用 /root/sync/reconcile.py 对照源和副本,生成 /root/sync/sync_result.json。
- 在 /root/sync/sync_report.md 中分四节进行汇报。
参考
- 源服务器的运行契约:
python3 /root/sync/source.py --port <포트>(占位符为端口)。/health返回{"ok": true, "rows": 60},/all按(updated_at, id)顺序返回全部数据,/rows?offset=&limit=是偏移量方式,/rows?since=&since_id=&limit=是游标方式。/mutate?ids=R-0001,R-0002会把这些行的updated_at推到当前最大值之后,并让value加 1。 - 共有 60 行,有
id·updated_at·value三栏。第 21 行到第 25 行这五行的updated_at是相同的——第 5 步的边界就出在这里。 - 遍历器的运行契约(两者通用):
--base <URL> --limit <n> [--snapshot <파일>] [--mutate-after <페이지>] [--mutate-ids <a,b,c>] [--out <jsonl>](占位符依次为文件名、页码)。输出为{"pages": ..., "fetched": ..., "unique": ..., "missing": [...], "duplicated": [...]},cursor_sync.py 在此之外还有watermark,并额外接收--state <파일>(占位符为文件名)和--max-pages <k>。 pages是拉取到一行以上的请求的数量。missing只有在指定了--snapshot时才会填充,不指定则为空列表。--out会把收到的行以每行一个 JSON 的方式追加写入。--mutate-after <페이지>(占位符为页码)会在收到该页之后立刻更新--mutate-ids中的行。这是我们用来模拟“在遍历期间源发生变化”这一情形的装置。- 比对器的运行契约:
python3 reconcile.py --snapshot <원본 JSON> --sink <사본 JSONL> --out <결과 JSON>(占位符依次为源 JSON、副本 JSONL、结果 JSON)会返回{"source_rows": ..., "sink_rows": ..., "unique": ..., "missing": [...], "extra": [...], "value_mismatch": [...], "match": ...}。sink_rows是副本文件的行数,unique是不同 id 的数量。如果同一个 id 出现多次,就取updated_at最大的那一条。 - 边界报告格式:
{"limit": ..., "pages": ..., "fetched": ..., "missing": [...], "duplicated": [...], "tie_updated_at": ..., "tie_ids": [...]}。 - 常见错误:只用
updated_at一个值作游标;在处理之前保存水位线;把重复当作 bug 并想消除它(至少一次才是正常的);只核对条数而不核对值。 - 服务器要放在后台启动,等到
/health返回 200 之后再继续。评分器不会查看你启动的进程,而是直接重新启动脚本。
启动不断变化的源
创建 /root/sync/source.py 并在端口 8019 上启动,把 /all 保存到 /root/sync/snapshot.json。共有 60 行,并且第 21 行到第 25 行这五行的 updated_at 必须相同。
排序始终按 (updated_at, id) 两栏进行。/rows 在收到 since 时用游标方式应答,否则用偏移量方式应答即可。/mutate 要把指定行的 updated_at 推到比当前最大值更靠后,才会在排序中跑到最末尾。
偏移量遍历吞掉的三行
创建 /root/sync/offset_sync.py,用 --limit 10 遍历,并在第一页之后立刻更新已经翻过去的 3 行(R-0001、R-0002、R-0003)。结果中的 missing 和 duplicated 必须各是 3 条。
已经读过的行被更新、跑到排序最后面时,它后面的行会向前移动。偏移量不知道这一点,所以会跳过被移动的那部分。而跑到最后面的行会在最后一页被再抓到一次。这两种现象同时发生,就是关键。
用值而不是编号来确定位置
创建 /root/sync/cursor_sync.py,用 (updated_at, id) 游标遍历。与第 2 步一样,在第一页之后立刻更新三行,missing 也必须是 0。被更新的行再次被抓到(duplicated)是正常的。
每次请求都传入最后看到的那一行的 (updated_at, id),源从比这个值更大的部分开始返回。比较必须两栏一起做——只比较一个时刻,会丢失相同时刻的行,或者无限地停在同一个位置打转。不要试图消除重复。那个重复正是“期间发生了变化”这一事实的如实反映。
记住下一次运行从哪里开始
给 cursor_sync.py 加上 --state <파일>(占位符为文件名)。如果文件存在,就从那个位置开始,并且每拉取一页都要更新并保存水位线。同一条命令运行两次时,第二次的 fetched 必须是 0。
水位线必须在磁盘上,中断后才能保留下来。保存的时点很重要——如果在处理收到的行之前保存,中断时就会丢失那一页;如果在处理之后保存,就会重新拉取。重新拉取比丢失要好。覆盖文件时,要先写入临时文件再替换,这样即使中途挂掉也不会损坏。
拥有相同时刻的五行
用 --limit 3 把全部数据扫一遍,让页恰好切在拥有相同 updated_at 的五行的正中间。即便如此,也必须没有遗漏和重复。把结果以 limit·pages·fetched·missing·duplicated·tie_updated_at·tie_ids 写入 /root/sync/tie_report.json。
如果只用一个时刻作游标,这里就会出现两种情况之一。只取更大的值,相同时刻的其余行就会整批消失;取相等或更大的值,就会重新拉取已经读过的行,在同一个位置打转。两栏一起比较,就不会出现这两种情况。tie 行可以在快照中按 updated_at 分组统计来找到。
接续被中断的同步
先用 --max-pages 2 --limit 10 --out sink.jsonl --state resume.json 在中途停下,再用同一个状态文件重新运行,拉取其余部分,使 /root/sync/sink.jsonl 中汇集全部 60 行。
检查是否能接续的方法很简单。第二次运行不从头重新拉取,只拉取剩下的部分就行。副本文件是追加写入的,所以两次运行的结果会累积在同一个文件里。开始之前要删掉旧文件,统计数才会对得上。
数一数是否全部拉取到了
创建 /root/sync/reconcile.py,对照 /root/sync/snapshot.json 和 /root/sync/sink.jsonl,并把结果写入 /root/sync/sync_result.json。不仅要看条数,还必须看相同 id 的 value 是否一致。
只是条数对得上,并不算对。少了一条又重复了一条,条数就不会变。如果副本中同一个 id 有多行,updated_at 最大的那一条才是当前值。只有 missing、extra 和值不一致全都为空时,才把 match 设为 true。
同步检查报告
在 /root/sync/sync_report.md 中分为 ## 오프셋 순회가 무엇을 빠뜨렸나 ## 커서로 옮긴 뒤 무엇이 달라졌나 ## 같은 시각의 경계 ## 중단과 재개, 그리고 대사 四节来写(韩文,依次意为“偏移量遍历漏掉了什么”“换成游标之后有什么不同”“相同时刻的边界”“中断与续传,以及对账”)。tie_report.json 和 sync_result.json 中的数字必须写进正文。
读者是那个在问“明明每天全部拉取,为什么还会缺几条”的人。不要把原因笼统地说成“时序问题”,请按顺序写出哪一行为什么漏了过去。也要写明重复不是 bug,这样下一个人才不会想去消除它。