An incremental job reads a source table and loads only what changed since its last run. Its saved state is {"watermark": <ISO timestamp or None>} — the largest updated_at it has loaded (None before the first run).
Write plan_load(state, batch) returning (rows_to_load, new_state):
- •
batch is a list of row dicts, each with an id and an updated_at (ISO strings like "2026-06-01T10:00:00", all in the same format). - •Load only rows with
updated_at strictly after the watermark (all rows when it's None) — rows at the watermark were loaded last time. - •Load one row per `id`: the one with the latest
updated_at. If two rows for an id have the same updated_at, the one later in batch wins. - •
rows_to_load is sorted by id. - •
new_state holds the largest updated_at among the rows you load. If you load nothing, the watermark stays as it was. - •Don't modify
state or the rows; return a new state dict.
Running plan_load again on the same batch with the returned state must load nothing and return the same state.
Python 3.13 in your browser — the standard library plus pandas and numpy; no pip installs.