TT Lab
开始
学习 学习路径 课程

集成与部署

每天都全量拉取,却每次都少几行

在 TT Lab 中继续学习

目标

制作一个把不断变化的源分页拉取的同步程序。亲手重现偏移量方式为什么会悄悄漏掉行,换成 (updated_at, id) 游标,用水位线从中断点续传,最后对照源和副本。

为什么重要

在把 60 万条分成每次 1000 条、共 600 次拉取的期间,源也在不断变化。已经翻过去的页中的行被更新、在排序顺序中跑到最后面,后面的行就会向前移动,它们之间的行就在没有被任何人读取的情况下过去了。 这类事故不会报错,条数大致对得上,而且每次漏掉的行都不一样。所以无法重现,几个月之后会以投诉的形式回来。 修复的办法,是用值而不是编号来确定位置。不过,如果排序键不唯一,还有另一个陷阱在等着——只用 updated_at 一个值,就会丢失相同时刻的行,或者无限地停在同一个位置打转。所以游标要增加栏位直到唯一为止。 评分器不会相信你写的句子。评分器会把源服务器启动在评分器所选的端口上,一边修改源,一边实际运行你的同步程序,统计漏掉了什么、重叠了什么。

步骤

  1. 创建 /root/sync/source.py 并在端口 8019 上启动,把全部数据做成快照保存到 /root/sync/snapshot.json。
  2. 创建 /root/sync/offset_sync.py,用偏移量遍历,同时在途中修改源,用数字展示遗漏与重复同时出现。
  3. 创建 /root/sync/cursor_sync.py,用 (updated_at, id) 游标遍历,使在同样的情况下遗漏变为 0。
  4. 给 cursor_sync.py 加上 --state,把水位线保存到文件中,让第二次运行只拉取有变化的部分。
  5. 即使让页恰好切在拥有相同 updated_at 的行的正中间,也没有遗漏和重复,把这一点写入 /root/sync/tie_report.json。
  6. 接续被中断的同步,把全部数据汇集到 /root/sync/sink.jsonl。
  7. 用 /root/sync/reconcile.py 对照源和副本,生成 /root/sync/sync_result.json。
  8. 在 /root/sync/sync_report.md 中分四节进行汇报。

参考

启动不断变化的源

创建 /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,这样下一个人才不会想去消除它。