We Pull Everything Daily and Still Lose a Few Rows
Goal
You build a synchronization that receives a constantly changing source in pages. You reproduce yourself why the offset method silently drops rows, move to a (updated_at, id) cursor, resume from the interruption point with a watermark, and finally compare the source and the copy.
Why it matters
While receiving 600,000 records in 600 passes of 1,000, the source keeps changing too. If a row on a page we already passed is updated and goes to the very end of the sort order, the rows behind it are pulled forward, and the rows in between pass by without anyone reading them.
This incident raises no error, the count is roughly right, and a different row goes missing each time. So it cannot be reproduced and comes back months later as a report.
The fix is to mark the position by value, not by number. But if the sort key is not unique, another trap awaits — if you use updated_at alone, you lose rows with the same timestamp or go round the same place forever. So the cursor adds slots until it becomes unique.
The grader does not believe your sentences. It starts the source server directly on a port the grader chooses, changes the source as it goes, actually runs your synchronization, and counts what was missed and what overlapped.
Steps
- Create /root/sync/source.py, run it on port 8019, and take a snapshot of everything in /root/sync/snapshot.json.
- Create /root/sync/offset_sync.py so that it traverses by offset but changes the source midway, showing in numbers that omissions and duplicates arise together.
- Create /root/sync/cursor_sync.py so that it traverses with a
(updated_at, id)cursor and the omissions become 0 in the same situation. - Attach
--stateto cursor_sync.py so that it leaves a watermark in a file and the second run receives only what changed. - Even when the page is made to cut through the middle of rows that have the same
updated_at, write in /root/sync/tie_report.json that there are neither omissions nor duplicates. - Resume an interrupted synchronization and gather everything in /root/sync/sink.jsonl.
- With /root/sync/reconcile.py, compare the source and the copy and create /root/sync/sync_result.json.
- Report in four sections in /root/sync/sync_report.md.
Notes
- Source server run contract:
python3 /root/sync/source.py --port <포트>(the placeholder is the port)./healthreturns{"ok": true, "rows": 60},/allreturns everything in(updated_at, id)order,/rows?offset=&limit=is the offset method, and/rows?since=&since_id=&limit=is the cursor method./mutate?ids=R-0001,R-0002pushes theupdated_atof those rows behind the current largest value and raisesvalueby 1. - There are 60 rows with three slots
id,updated_at, andvalue. Five rows, number 21 through 25, have the sameupdated_at— the boundary in step 5 arises here. - Traverser run contract (common to both):
--base <URL> --limit <n> [--snapshot <파일>] [--mutate-after <페이지>] [--mutate-ids <a,b,c>] [--out <jsonl>](the placeholders are file, page, and jsonl). The output is{"pages": ..., "fetched": ..., "unique": ..., "missing": [...], "duplicated": [...]}, and cursor_sync.py adds awatermarkto this and additionally accepts--state <파일>and--max-pages <k>. pagesis the number of requests that received one or more rows.missingis filled only if you give--snapshot, and is an empty list if you do not.--outappends the received rows, one per line as JSON.--mutate-after <페이지>updates the rows of--mutate-idsin the source right after receiving that page. It is a device with which we imitate the situation of the source changing while running.- Comparator run contract:
python3 reconcile.py --snapshot <원본 JSON> --sink <사본 JSONL> --out <결과 JSON>(the placeholders are the source JSON, copy JSONL, and result JSON) outputs{"source_rows": ..., "sink_rows": ..., "unique": ..., "missing": [...], "extra": [...], "value_mismatch": [...], "match": ...}.sink_rowsis the number of lines in the copy file anduniqueis the number of distinct ids. If the same id appears several times, use the one with the largestupdated_at. - Boundary report format:
{"limit": ..., "pages": ..., "fetched": ..., "missing": [...], "duplicated": [...], "tie_updated_at": ..., "tie_ids": [...]}. - Common mistakes: using only
updated_atas the cursor, saving the watermark before processing, treating duplicates as bugs and trying to remove them (at-least-once is normal), and matching only counts and not checking values. - Run the server in the background, wait until
/healthis 200, and then move on. The grader does not look at the process you left running but restarts the scripts directly.
Start a source that keeps changing
Create /root/sync/source.py, run it on port 8019, and save /all to /root/sync/snapshot.json. There are 60 rows, and the updated_at of five rows, number 21 through 25, must be the same.
Always sort by the two slots (updated_at, id). /rows can answer with the cursor method if since comes and the offset method otherwise. /mutate must push the specified rows' updated_at behind the current largest value for them to go to the very end of the sort.
The three rows the offset traversal swallowed
Create /root/sync/offset_sync.py so that it traverses with --limit 10 but, right after the first page, updates 3 rows that have already been passed (R-0001, R-0002, R-0003). The missing and duplicated in the result must each be 3.
If a row that has already been read is updated and goes to the very end of the sort, the rows after it are pulled forward. The offset does not know that, so it skips as many as were pulled forward. And the rows that went to the end are caught once more on the last page. The core is that the two phenomena happen at the same time.
Mark the position by value, not by number
Create /root/sync/cursor_sync.py so that it traverses with a (updated_at, id) cursor. Even if you update three rows right after the first page exactly as in step 2, missing must be 0. Updated rows being caught again (duplicated) is normal.
With each request, pass the (updated_at, id) of the last row seen, and the source returns those greater than that value. The comparison must be on both slots together — if you compare only the timestamp, you lose rows with the same timestamp or go round the same place forever. Do not try to remove the duplicates. Those duplicates are exactly the fact "it changed in the meantime."
Remember where the next run starts
Attach --state <파일> to cursor_sync.py (the placeholder is the file). If the file exists, start from that position, and every time you receive a page you must update and save the watermark. If you run the same command twice, fetched must be 0 in the second run.
The watermark must be on disk to survive an interruption. The moment of saving matters — if you save before processing the received rows, you lose that page on interruption, and if you save after processing, you receive it again. Receiving again is better than losing. When overwriting the file, write to a temporary file and swap, so it does not break even if it dies midway.
Five rows with the same timestamp
Sweep through everything once with --limit 3 so that a page cuts through the middle of the five rows that have the same updated_at. There must still be neither omissions nor duplicates. Write the result to /root/sync/tie_report.json as limit, pages, fetched, missing, duplicated, tie_updated_at, and tie_ids.
If you take the cursor with the timestamp alone, one of two things happens here. If you receive only larger values, the rest of the same timestamp vanishes entirely, and if you receive equal-or-larger values, you receive already-read rows again and go round the same place. If you compare the two slots together, neither happens. You can find the tie rows by grouping and counting by updated_at in the snapshot.
Resume an interrupted synchronization
After stopping midway with --max-pages 2 --limit 10 --out sink.jsonl --state resume.json, run it again with the same state file to receive the rest, so that all 60 rows gather in /root/sync/sink.jsonl.
The way to check that resuming works is simple. It is enough if the second run does not receive from the beginning again but receives only what is left. The copy file is written by appending, so the results of the two runs pile up in one file. You must delete the old file before starting for the count to come out right.
Count whether everything was received
Create /root/sync/reconcile.py to compare /root/sync/snapshot.json and /root/sync/sink.jsonl and write the result to /root/sync/sync_result.json. You must look not only at the count but also whether the value of the same id matches.
Only the count matching is not matching. If one is missing and one is duplicated, the count stays the same. If the copy has several lines for the same id, the one with the largest updated_at is the current value. Set match to true only when missing, extra, and value mismatches are all empty.
Synchronization inspection report
In /root/sync/sync_report.md, write four sections, ## 오프셋 순회가 무엇을 빠뜨렸나, ## 커서로 옮긴 뒤 무엇이 달라졌나, ## 같은 시각의 경계, and ## 중단과 재개, 그리고 대사 (in order: what the offset traversal missed, what changed after moving to the cursor, the boundary of the same timestamp, interruption and resumption and reconciliation). The numbers from tie_report.json and sync_result.json must be in the body.
The reader is someone asking "we receive everything every day, so why are a few records missing?" Do not lump the cause into a "timing issue"; write in order which rows passed by and why. Also write that duplicates are not a bug, so the next person does not try to remove them.