How to checkpoint the steps of a handler¶
This guide shows you how to make an individual step inside a task handler durable, so that a retry or a crash-recovery re-run doesn't repeat it.
Wrap a step in checkpoint¶
def import_report(engine, task):
rows = engine.checkpoint(
task_id=task.task_id,
step_name="download",
fn=lambda: fetch_from_sftp(task.params["uri"]),
)
engine.checkpoint(
task_id=task.task_id,
step_name="load",
fn=lambda: load_into_db(rows),
)
return {"rows": len(rows)}
The first time checkpoint runs for a given (task_id, step_name), it
calls fn, stores the result, and returns it. Every later call for that
same task and step name returns the stored result immediately, without
calling fn again. That's true whether the call comes from a retry after
fail_run, or from a worker reclaiming an abandoned run after a crash.
Choose step names carefully¶
Checkpoints are keyed on (task_id, step_name), not on where in the code
the call appears. If a handler runs the same logical step twice with the
same step_name (say, inside a loop), the second call just gets the
first call's result.
Give each logically distinct step its own name. For a step inside a
loop, fold the iteration index into the step name, like f"row-{i}".
Keep fn itself fast to retry, slow to run¶
fn executes outside dura's write transaction, so a slow step (an SFTP
fetch, an S3 upload) doesn't hold the database lock. But fn isn't
guarded against running concurrently with itself. If two workers somehow
claimed the same run at once, which dura's leases otherwise prevent,
whichever commits first wins, and the other's result is discarded.
In practice, a single run is only ever claimed by one worker at a time.
This only matters if you're calling checkpoint directly, outside the
normal claim/complete lifecycle.
Read a checkpoint without risking a computation¶
To check whether a step has already run, without triggering it:
downloaded = engine.get_checkpoint(task_id=task.task_id, step_name="download")
if downloaded is None:
... # hasn't run yet
Difference from durable state¶
Checkpoints belong to a task and are deleted when that task is cleaned up
by cleanup(). If you need a value that outlives the task that computed
it (a cursor, a watermark, a dedup record), use durable state instead; see
How to use durable state.