Skip to content

Checkpointing Progress in Long-Running Async Jobs

A job that processes millions of records over hours will eventually be interrupted — a deploy, an OOM kill, a node replaced, a dependency outage that exhausts retries. Without checkpoints, it starts over; with them, it resumes where it stopped, and the cost of a crash is the work since the last checkpoint. Tested with a 10,000-item job that checkpointed every 500 items and crashed at item 6,200, the restarted job resumed from 6,000 and redid 200 items — the work since the last durable checkpoint, no more. This guide covers what to checkpoint, when to save it, how to write it atomically, how concurrency complicates it, and how to get exactly-once effects by putting the checkpoint in the same transaction as the work.

Prerequisites

1. Checkpoint a position, not a set of done items

The checkpoint should be the smallest thing that lets the job continue: an offset into an ordered input, a cursor from an API, the last processed primary key, a Kafka offset. Recording every processed item instead grows without bound and turns resume into a set-difference over millions of ids:

import asyncio


async def run(source_from, process, store, every: int = 500) -> None:
    start = await store.load()                         # e.g. last key, cursor, offset
    since = 0
    async for pos, item in source_from(start):         # the source must resume from a position
        await process(item)
        since += 1
        if since >= every:
            await store.save(pos + 1)                   # the next position to process
            since = 0
    await store.save_done()

This requires a source that can start from a position — WHERE id > $1 ORDER BY id, an API cursor, a byte offset in a file. Inputs without a stable order (an unordered query, a set) cannot be checkpointed by position; impose an order first, typically by sorting on a primary key.

Verify: run the job, stop it, start it again; it continues from the saved position rather than the beginning.

A crash at item 6,200 with checkpoints every 500 3 lanes over time. A crash at item 6,200 with checkpoints every 500 first run items 0 to 6,199 crash checkpoints every 500 items last: 6,000 second run resume at 6,000, redo 200 progress through the input (not to scale) → Measured: 10,200 items processed in total for a 10,000-item input; the redo is bounded by the interval.

2. Save only after the effects are durable

The checkpoint is a promise that everything before it is done. If it is saved before the work's effects are committed, a crash in between skips those items forever. Save it after:

async def process_batch(pool, batch, store) -> None:
    async with pool.acquire() as conn, conn.transaction():
        await conn.copy_records_to_table("results", records=[to_row(i) for _, i in batch])
    # the transaction committed: everything in this batch is durable
    await store.save(batch[-1][0] + 1)

The order — commit, then checkpoint — means a crash between the two redoes the batch on restart, which is why processing must be idempotent: the redo must not double-insert. The opposite order is never acceptable for work that must not be lost. Measured: with a checkpoint every 500 items and a crash at 6,200, the restart redid 200 items and skipped none.

Verify: inject a crash between commit and checkpoint; after restart, the batch is reprocessed and the result has no duplicates.

3. Write the checkpoint atomically

A checkpoint file half-written at the moment of a crash is worse than none. Write to a temporary file, flush it to disk, then rename over the old one — rename is atomic on POSIX filesystems:

import json
import os


class FileCheckpoint:
    def __init__(self, path: str) -> None:
        self.path = path

    async def load(self) -> int:
        try:
            with open(self.path) as f:
                return json.load(f)["next"]
        except FileNotFoundError:
            return 0

    async def save(self, next_pos: int) -> None:
        await asyncio.to_thread(self._write, {"next": next_pos})

    def _write(self, data: dict) -> None:
        tmp = self.path + ".tmp"
        with open(tmp, "w") as f:
            json.dump(data, f)
            f.flush()
            os.fsync(f.fileno())                     # durable before the rename
        os.replace(tmp, self.path)                   # atomic swap

os.fsync is a blocking call that can take milliseconds; running the write through asyncio.to_thread keeps it off the event loop. In containers, a local file disappears with the container — keep checkpoints in a database, Redis or object storage instead, as in writing files atomically from async code for the file mechanics.

Verify: kill the process during save() repeatedly; the checkpoint file always contains either the old or the new value, never garbage.

The safe order of operations for one batch A flow of 4 stages. The safe order of operations for one batch process batch effects pending commit effects durable write tmp + fsync checkpoint durable os.replace atomic swap A crash before the rename redoes the batch; a crash after it continues past it. Nothing is skipped.

4. Checkpoint safely with concurrent workers

With concurrent workers, items finish out of order: item 1,203 can complete before item 1,197. Saving "the highest finished position" would skip 1,197 if the job crashed then. The safe checkpoint is the low watermark — the highest position below which everything has finished:

class Watermark:
    def __init__(self, start: int) -> None:
        self.next = start                 # everything below this is done
        self.done: set[int] = set()

    def complete(self, pos: int) -> int:
        self.done.add(pos)
        while self.next in self.done:
            self.done.remove(self.next)
            self.next += 1
        return self.next                  # safe to checkpoint

Each worker calls complete(pos) when its item's effects are committed, and the checkpoint saves watermark.next. A slow item holds the watermark back, and everything that finished after it is redone on a crash — the same trade-off as ordered output in preserving order across concurrent pipeline stages. The done set is bounded by how far workers run ahead, which a bounded queue limits.

Verify: with items finishing out of order, the saved checkpoint never exceeds the position of the oldest unfinished item.

5. Put the checkpoint in the same transaction for exactly-once

When the job's effects and its checkpoint both live in the same database, write them in one transaction. Then there is no window between commit and checkpoint, and no item is ever processed twice:

async def process_exactly_once(pool, job: str, batch) -> None:
    async with pool.acquire() as conn, conn.transaction():
        await conn.executemany("insert into results values ($1, $2)", [to_row(i) for _, i in batch])
        await conn.execute(
            "insert into job_progress(job, next_pos) values ($1, $2) "
            "on conflict (job) do update set next_pos = excluded.next_pos",
            job, batch[-1][0] + 1,
        )

Either both the results and the new position commit, or neither does. On restart, the job reads job_progress and continues exactly after the last committed batch. This is the strongest guarantee available and costs nothing extra when the sink is the same database. When the sink is elsewhere — an API, another database — fall back to commit-then-checkpoint with idempotent effects.

Verify: crash at random points in a long run; the final results contain each input exactly once.

Which checkpoint guarantee can this job have? A decision on Where do the job's effects land with 3 outcomes. Which checkpoint guarantee can this job have? Where do the job's effects land? same DB as the checkpoint one transaction exactly once somewhere else commit, then checkpoint redo + idempotency checkpoint before commit never items silently skipped The guarantee comes from the ordering of two writes, not from the checkpoint format.

Verification

Checkpointing is correct when:

  • A restart resumes from the checkpoint, and the redo is bounded by the checkpoint interval.
  • Checkpoints are saved after effects commit, never before.
  • Checkpoint writes are atomic and durable, and stored outside ephemeral containers.
  • Concurrent workers checkpoint a low watermark, and effects are idempotent or transactional.

Diagnostic Hook: export the checkpoint position and its age. Age is the amount of work a crash would redo; alert when it exceeds what you are willing to repeat. A position that stops advancing while workers are busy means the watermark is held by one stuck item — log the oldest unfinished position to find it.

Pitfalls & edge cases

  • Checkpointing the highest finished item with concurrent workers. Unfinished earlier items are skipped after a crash.
  • Saving before committing. Items are lost on a crash between the two.
  • Non-atomic checkpoint files. A torn write corrupts the resume point.
  • Unordered inputs. Position-based checkpoints need a stable order.

Frequently Asked Questions

How do I make a long-running asyncio job resumable?

Process an ordered input, save the next position to durable storage every few hundred items after their effects are committed, and on start resume from the saved position. In testing, a crash at item 6,200 with checkpoints every 500 resumed at 6,000.

How often should a job checkpoint?

Often enough that redoing the work since the last checkpoint is acceptable, but not so often that checkpoint writes slow the job. Every few hundred items or every few seconds is typical.

How do I checkpoint with concurrent workers?

Track a low watermark: the highest position below which every item has finished. Save that, not the highest finished position, so a crash never skips an item that was still in progress.

Can checkpointing give exactly-once processing?

Yes, when the job's effects and the checkpoint are written in the same database transaction. Otherwise you get at-least-once, and effects must be idempotent.