TT Lab
Get started
Learn Learning paths Courses

Idempotency — Two Clicks, One Charge

Resume imports from the last committed checkpoint

Continue in TT Lab

Goal

Combine the source fingerprint and the checkpoint to prevent wrongly resuming with a different file.

Why it matters

A job importing thousands of rows died partway. The operator swapped the file and ran it again with the same job id, and the result had the first part from the old file and the rest from the new file. If you store only the number of the processed row, you cannot confirm the identity of the input. You must keep the fingerprint of the source bytes together with the last committed position.

Steps

  1. In /root/work/idem-batch-checkpoint-lab/service.py, parse_rows(raw) reads JSON array bytes. Each item has an id (a non-empty str) and a value (an int, not bool), and duplicate ids are forbidden. A violation is a ValueError. It returns a list of rows that have only {id,value}.

Prepare it once at the beginning. It does not overwrite an existing file.

mkdir -p /root/work/idem-batch-checkpoint-lab
test -e /root/work/idem-batch-checkpoint-lab/service.py || cp /opt/fixtures/ten_labs/idem-batch-checkpoint-lab/service.py /root/work/idem-batch-checkpoint-lab/service.py
cd /root/work/idem-batch-checkpoint-lab
  1. In /root/work/idem-batch-checkpoint-lab/service.py, source_digest(raw) is the hex string obtained by applying SHA-256 to the bytes. It does not normalize the JSON.

  2. In /root/work/idem-batch-checkpoint-lab/service.py, init_db(path) idempotently creates imports(id TEXT PRIMARY KEY,digest TEXT NOT NULL,next_index INTEGER NOT NULL) and items(batch TEXT NOT NULL,id TEXT NOT NULL,value INTEGER NOT NULL,PRIMARY KEY(batch,id)).

  3. In /root/work/idem-batch-checkpoint-lab/service.py, begin(path,batch,digest) stores next_index=0 and returns 0 for a new job, returns next_index for an existing job with the same fingerprint, and raises ValueError for a different fingerprint.

  4. In /root/work/idem-batch-checkpoint-lab/service.py, apply_chunk(path,batch,start,rows,fault=lambda index:None) runs only when the current next_index==start, and otherwise raises ValueError. It puts rows into items in order and calls fault(overall index) after each insertion. If everything succeeds, it stores next_index=start+len(rows) and returns it.

  5. In /root/work/idem-batch-checkpoint-lab/service.py, checkpoint(path,batch) returns next_index, and None for a job that does not exist.

  6. In /root/work/idem-batch-checkpoint-lab/service.py, values(path,batch) returns the (id,value) tuples of that batch in ascending id order.

  7. In /root/work/idem-batch-checkpoint-lab/service.py, import_all(path,batch,raw,size=2,fault=lambda index:None) validates that size is a positive int (not bool) and uses parse_rows, source_digest, and begin. It processes the remaining rows size at a time with apply_chunk and returns the final checkpoint.

Notes

Check the contract of the input rows

In /root/work/idem-batch-checkpoint-lab/service.py, parse_rows(raw) reads JSON array bytes. Each item has an id (a non-empty str) and a value (an int, not bool), and duplicate ids are forbidden. A violation is a ValueError. It returns a list of rows that have only {id,value}.

Prepare it once at the beginning. It does not overwrite an existing file.

mkdir -p /root/work/idem-batch-checkpoint-lab
test -e /root/work/idem-batch-checkpoint-lab/service.py || cp /opt/fixtures/ten_labs/idem-batch-checkpoint-lab/service.py /root/work/idem-batch-checkpoint-lab/service.py
cd /root/work/idem-batch-checkpoint-lab

If a duplicate id is silently overwritten by the last row, the result of the import cannot be predicted.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/01-contract.sh.

Fix the fingerprint of the source bytes

In /root/work/idem-batch-checkpoint-lab/service.py, source_digest(raw) is the hex string obtained by applying SHA-256 to the bytes. It does not normalize the JSON.

The contract allows resuming only for the same source.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/02-contract.sh.

Store the input and the checkpoint

In /root/work/idem-batch-checkpoint-lab/service.py, init_db(path) idempotently creates imports(id TEXT PRIMARY KEY,digest TEXT NOT NULL,next_index INTEGER NOT NULL) and items(batch TEXT NOT NULL,id TEXT NOT NULL,value INTEGER NOT NULL,PRIMARY KEY(batch,id)).

Store the same row id from different import jobs separately.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/03-contract.sh.

Keep it from resuming with a different source

In /root/work/idem-batch-checkpoint-lab/service.py, begin(path,batch,digest) stores next_index=0 and returns 0 for a new job, returns next_index for an existing job with the same fingerprint, and raises ValueError for a different fingerprint.

Even with the same job id and row number, the input file may be different.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/04-contract.sh.

Commit the batch and the position atomically

In /root/work/idem-batch-checkpoint-lab/service.py, apply_chunk(path,batch,start,rows,fault=lambda index:None) runs only when the current next_index==start, and otherwise raises ValueError. It puts rows into items in order and calls fault(overall index) after each insertion. If everything succeeds, it stores next_index=start+len(rows) and returns it.

If you commit separately for every row, the checkpoint and the row state get out of step.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/05-contract.sh.

Look up the current position

In /root/work/idem-batch-checkpoint-lab/service.py, checkpoint(path,batch) returns next_index, and None for a job that does not exist.

Read the last committed position, not the position of the last attempt to process.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/06-contract.sh.

Separate the results per job

In /root/work/idem-batch-checkpoint-lab/service.py, values(path,batch) returns the (id,value) tuples of that batch in ascending id order.

Make batch a condition so that rows with the same id from another job do not get mixed into the result.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/07-contract.sh.

Continue safely after a failure in the middle

In /root/work/idem-batch-checkpoint-lab/service.py, import_all(path,batch,raw,size=2,fault=lambda index:None) validates that size is a positive int (not bool) and uses parse_rows, source_digest, and begin. It processes the remaining rows size at a time with apply_chunk and returns the final checkpoint.

Check the scenario of resuming after the second batch fails while preserving the success of the first batch.

After saving, check with bash /opt/lab/checks/idem-batch-checkpoint-lab/08-contract.sh.