Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Idempotent loads: commit target rows and extraction progress together

Last updated: 6 Oct 20265 min read
tutorial
IntermediateBy AITrove Editorial

An idempotent load produces the same target state after a retry and ties its progress marker to committed output.

Expect at-least-once work

A worker can write a receipt batch and crash before acknowledging it. The scheduler then replays the same batch. If the target only appends, the result doubles. Use a stable receipt key or source event ID and an upsert or deduplication rule. An upsert needs conflict semantics: a later version may replace an earlier one, but an old replay must not overwrite a newer value.

Place the checkpoint with the write

In one database transaction, validate a batch, apply its rows, and record the last source position. Either both commit or neither does. If target and checkpoint live in different systems, a single local transaction cannot cover them; design a replay-and-reconcile protocol instead of claiming exactly-once delivery. Compound extraction markers] define which page to read next.

Make mutations observable

Record rows inserted, rows updated, stale events skipped and rejected records separately. A no-op retry should still confirm the batch identity was seen. Test a crash before commit, a crash immediately after commit, and the same batch delivered twice. The final target rows and checkpoint must match the uninterrupted run.

Separate identity from batch position

A batch ID prevents a repeated file from applying twice, but it cannot by itself prevent an older correction file from overwriting a newer receipt. Retain row-level version or source position for mutable records. CDC ordering] makes this especially important.

Implementation

python
with database:
    if not database.execute(
        "SELECT 1 FROM processed_batches WHERE batch_id = ?",
        (batch_id,)).fetchone():
        database.executemany(
            "INSERT INTO receipts(receipt_id, amount_cents) VALUES(?, ?) "
            "ON CONFLICT(receipt_id) DO UPDATE SET amount_cents=excluded.amount_cents",
            receipt_rows)
        database.execute(
            "INSERT INTO processed_batches(batch_id, last_position) VALUES(?, ?)",
            (batch_id, last_source_position))

Performance and operating cost

With indexed keys, batch upserts cost about O(B log N) for B rows in a table of N rows; transaction logs and unique indexes add write amplification. Reconciliation scans can dominate very large backfills.

Common Mistakes

  • Do not acknowledge a batch before its target transaction commits.
  • Do not claim a batch ID protects against out-of-order row versions.
  • Do not claim exactly-once effects across two independent stores without a recovery protocol.

Read next

Continue the workflow: Failure classification and retry budgets.

ai-data
data-engineering
Storage details