毎日まるごと取り込んでいるのに、毎回いくつか欠ける
目標
絶えず変わる元データを、ページに分けて受け取る同期を作ります。オフセット方式がなぜ黙って行を落とすのかを、自分で再現し、(updated_at, id)のカーソルに移し、ウォーターマークで中断した地点から再開し、最後に、元データとコピーを照合します。
なぜ重要なのか
60万件を1000件ずつ600回に分けて受け取る間にも、元データは変わり続けます。すでに通り過ぎたページの行が更新されて、並び順の一番後ろに行くと、後ろにあった行たちが前に引き寄せられ、その間の行は、誰にも読まれないまま通り過ぎます。
この事故は、エラーが出ず、件数はほぼ合い、毎回違う行が抜けます。そのため、再現できず、数か月後に報告として戻ってきます。
直す方法は、位置を番号ではなく値で指定することです。ただし、並べ替えキーが一意でなければ、別の落とし穴が待っています。updated_at1つで指定すると、同じ時刻の行を失うか、同じ場所を無限に回ります。そのため、カーソルは、一意になるまで欄を増やします。
採点ツールは、作成した文章を信じません。元データのサーバーを、採点ツールが選んだポートで直接起動し、採点ツールが元データを変えながら、作成した同期を実際に動かして、何が抜けて、何が重なったかを数えます。
ステップ
- /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を付けて、ウォーターマークをファイルに残し、2回目の実行が、変わったものだけを受け取るようにしてください。 - ページが、同じ
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に、4つの節で報告してください。
参考
- 元データのサーバーの実行契約:
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の3つの欄です。21行目から25行目までの5行は、updated_atが同じです。ステップ5の境界は、ここで生じます。 - 巡回器の実行契約(2つに共通):
--base <URL> --limit <n> [--snapshot <파일>] [--mutate-after <페이지>] [--mutate-ids <a,b,c>] [--out <jsonl>](プレースホルダーは順にファイル、ページ、JSONLです)。出力は{"pages": ..., "fetched": ..., "unique": ..., "missing": [...], "duplicated": [...]}で、cursor_sync.pyは、ここにwatermarkがさらに付き、--state <파일>と--max-pages <k>をさらに受け取ります(プレースホルダーはファイルです)。 pagesは、1行以上を受け取ったリクエストの数です。missingは、--snapshotを渡したときだけ埋め、渡さなければ空のリストです。--outは、受け取った行を、1行に1つずつ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_at1つだけで指定すること、ウォーターマークを処理の前に保存すること、重複をバグと見てなくそうとすること(少なくとも1回が正常です)、件数だけを合わせて、値を合わせてみないこと。 - サーバーはバックグラウンドで起動し、
/healthが200になるまで待ってから、次に進みます。採点ツールは、起動しておいたプロセスを見ず、スクリプトを直接起動し直します。
絶えず変わる元データを起動する
/root/sync/source.pyを作成してポート8019で起動し、/allを/root/sync/snapshot.jsonに保存してください。行は60個で、21行目から25行目までの5行のupdated_atが、同じである必要があります。
並べ替えは、いつも(updated_at, id)の2つの欄で行います。/rowsは、sinceが来ればカーソル方式、そうでなければオフセット方式で答えればよいです。/mutateは、指定した行たちのupdated_atを、現在の最大値より後ろに押しやってこそ、並べ替えで一番最後に行きます。
オフセットの巡回が飲み込んだ3行
/root/sync/offset_sync.pyを作成して、--limit 10で巡回しつつ、最初のページの直後に、すでに通り過ぎた行3つ(R-0001、R-0002、R-0003)を更新してください。結果のmissingとduplicatedが、それぞれ3件である必要があります。
すでに読んだ行が更新されて、並べ替えの一番後ろに行くと、その後ろの行たちが、前に引き寄せられます。オフセットは、その事実を知らないので、引き寄せられた分を飛ばします。そして、一番後ろに行った行たちは、最後のページでもう一度捕まります。2つの現象が同時に起きることが、核心です。
位置を番号ではなく値で
/root/sync/cursor_sync.pyを作成して、(updated_at, id)のカーソルで巡回するようにしてください。ステップ2とまったく同じように、最初のページの直後に3つの行を更新しても、missingが0である必要があります。更新された行が再び捕まること(duplicated)は、正常です。
リクエストごとに、最後に見た行の(updated_at, id)を渡し、元データは、その値より大きいものから返します。比較は、2つの欄を一緒に行う必要があります。時刻1つだけを比較すると、同じ時刻の行を失うか、同じ場所を無限に回ります。重複をなくそうとしないでください。その重複は、「その間に変わった」という事実そのままです。
次の実行がどこからかを覚えておく
cursor_sync.pyに--state <파일>を付けてください(プレースホルダーはファイルです)。ファイルがあれば、その位置から始め、ページを受け取るたびに、ウォーターマークを更新して保存する必要があります。同じコマンドを2回実行すると、2回目はfetchedが0である必要があります。
ウォーターマークは、ディスクにあってこそ、中断されても生き残ります。保存の時点が重要です。受け取った行を処理する前に保存すると、中断時にそのページを失い、処理したあとに保存すると、受け取り直します。失うよりは、受け取り直すほうがましです。ファイルを上書きするときは、一時ファイルに書いて入れ替えなければ、途中で死んでも壊れません。
同じ時刻を持つ5行
--limit 3で全体を1回なぞって、同じupdated_atを持つ5行の真ん中を、ページが分けるようにしてください。それでも、欠落も重複もない必要があります。結果を、/root/sync/tie_report.jsonに、limit・pages・fetched・missing・duplicated・tie_updated_at・tie_idsで書いてください。
カーソルを時刻1つだけで指定すると、ここで2つのうちどちらかが起きます。大きい値だけを受け取れば、同じ時刻の残りがまるごと消え、同じか大きい値を受け取れば、すでに読んだ行をまた受け取って、同じ場所を回ります。2つの欄を一緒に比較すれば、どちらも起きません。tieの行たちは、スナップショットで、updated_atでまとめて数えれば見つかります。
中断された同期を引き継ぐ
--max-pages 2 --limit 10 --out sink.jsonl --state resume.jsonで、途中で止めたあと、同じ状態ファイルで再び実行して、残りを受け取り、/root/sync/sink.jsonlに、全体の60行が集まるようにしてください。
引き継ぎができているかを見る方法は、簡単です。2回目の実行が、最初から受け取り直さず、残りだけを受け取ればよいのです。コピーのファイルは、追記して書くので、2回の実行の結果が、1つのファイルに積み重なります。始める前に、古いファイルを消さなければ、数が合いません。
全部受け取ったかを数えてみる
/root/sync/reconcile.pyを作成して、/root/sync/snapshot.jsonと/root/sync/sink.jsonlを照合し、結果を/root/sync/sync_result.jsonに書いてください。件数だけでなく、同じidのvalueまで合っているかを、見る必要があります。
件数だけが合っているのは、合っているのではありません。1件が抜けて1件が重複しても、件数は変わりません。コピーに同じidが複数行あれば、updated_atが最大のものが、現在の値です。missingとextraと値の不一致が、すべて空のときだけ、matchをtrueにしてください。
同期の点検報告書
/root/sync/sync_report.mdに、## 오프셋 순회가 무엇을 빠뜨렸나、## 커서로 옮긴 뒤 무엇이 달라졌나、## 같은 시각의 경계、## 중단과 재개, 그리고 대사(韓国語の見出しで、順にオフセットの巡回が何を落としたか、カーソルに移したあと何が変わったか、同じ時刻の境界、中断と再開そして突合、を意味します)の4つの節で書いてください。tie_report.jsonとsync_result.jsonの数字が、本文に入っている必要があります。
読む人は、「毎日全部受け取っているのに、なぜ数件ずつ抜けるのか」と尋ねる人です。原因を「タイミングの問題」とひとまとめにせず、どの行がなぜ通り過ぎたかを、順に書いてください。重複がバグではないという点も、一緒に書いておかなければ、次の人が、それをなくそうとしてしまいます。