skip to content

How does a custom Airbyte source emit cursor state so the next sync resumes correctly?

level: middleimportance: must knowfreq 52%

answer

  1. The platform hands back what you last emitted
  2. A checkpoint is a promise about what came before
  3. Emit records first, then the checkpoint
  4. Every N records, not only at stream end
  5. Inclusive resume means boundary duplicates

basics

~20 s

The stream declares a cursor field, then emits STATE messages during read holding the highest cursor value whose records have already been emitted. The platform persists the last STATE it saw and hands it back as the next run's input, so the connector resumes from there.

solid answer

~50 s

In the Python CDK you declare `cursor_field` on the stream and implement `IncrementalMixin`: a `state` getter the framework reads to emit checkpoints and a `state` setter the framework calls with the previous run's value. You update the getter's value only *after* emitting the records it covers — a STATE message means "everything before this is safely handed off", so emitting it early loses rows on a mid-sync crash. Setting `state_checkpoint_interval` makes the framework emit STATE every N records instead of only at the end of the stream, which is what stops a six-hour sync from restarting at zero. In a declarative manifest the same job is done by a `DatetimeBasedCursor` with `cursor_field`, `datetime_format`, `start_datetime` and optionally `step` and `lookback_window`. Resume is inclusive in practice, so records land twice at the boundary — declare a primary key so the destination can dedupe.

code

text · 5 lines
text
RECORD  users  {id: 1, updated_at: 2024-03-01T09:00Z}
RECORD  users  {id: 2, updated_at: 2024-03-01T09:30Z}
STATE   users  {updated_at: 2024-03-01T09:30Z}   <- covers ids 1 and 2
RECORD  users  {id: 3, updated_at: 2024-03-01T10:00Z}
-- crash here: next run resumes at 09:30 and re-reads id 3, losing nothing

go deeper

for a junior

Know that an incremental connector declares a cursor field and emits state, and that the platform gives that state back on the next run so it does not start from scratch.

for a middle

Explain the mechanics: the state getter and setter, the checkpoint interval, and why records must be emitted before the state that covers them.

for a senior

Diagnose the real failures — full re-reads every run, rows lost after a crash, state advancing past unordered records — and show how slicing gives you provable checkpoint boundaries on a long backfill.

for a principal

Own the guarantee you are offering downstream: at-least-once with boundary duplicates, the primary-key requirement that implies, the lateness margin you choose, and what a cursor-field change means for every table already loaded.

## What state is in the protocol A `read` may receive `--state`, and it emits `STATE` messages as it goes. The platform stores the most recent STATE message a sync produced and passes it back on the next run. Modern connectors emit per-stream state: a state message carrying a `stream_descriptor` naming the stream and a `stream_state` blob whose contents are entirely the connector's business — `{"updated_at": "2024-03-01T10:00:00Z"}`, an id high-water mark, whatever the source lets you resume from. The framework never interprets it. That freedom is why state bugs are the most common class of custom-connector defect: nothing validates that your blob means what you think. ## Declaring the cursor in Python A stream declares `cursor_field`, the record field the state tracks. Older connectors implemented `get_updated_state(current_stream_state, latest_record)`, returning the merged state after each record. Newer ones implement `IncrementalMixin`, exposing `state` as a property with a getter and setter: the setter receives the previous run's blob at start-up, and the getter is read by the framework whenever it wants to checkpoint. Either way the arithmetic is the same — take the maximum of the stored cursor and the value on records you have already emitted. ## The ordering rule that actually matters A STATE message is a promise about the records that precede it in the output stream. The platform will, on the next run, start from that state and never re-read anything the state claims is done. So the only safe sequence is: emit records, then emit the state covering them. Two ways teams break it: **Advancing state before yielding.** If you compute the maximum cursor over a page and update state, then yield the page's records, a crash between the two loses that page permanently — the next sync starts after records that were never emitted. **Advancing state on unordered data.** If the API returns records in arbitrary order and you take the running maximum, a checkpoint mid-page can record a cursor higher than records still to come. Resume from it and those laggards are skipped. The fix is either to request a sorted read from the source (order by the cursor field ascending) or to only advance state at a boundary you can prove is complete — the end of a time slice, for instance. ## Checkpoint frequency If state is emitted only when a stream finishes, a sync that fails at hour five re-does five hours of work and re-hammers the API's rate limit. `state_checkpoint_interval` tells the framework to emit STATE every N records, giving you bounded rework. The trade is that more frequent checkpoints mean more state writes and, with an append destination, a longer tail of partially-loaded data on failure. Reasonable connectors checkpoint on the order of thousands of records, or once per completed slice. Slicing interacts here. `stream_slices` (or a declarative `step` on the cursor) splits a large read into windows — a month at a time on a two-year backfill. Each slice is a natural, provably-complete checkpoint boundary, which is why slicing plus per-slice state is the pattern that makes long backfills survivable. ## The declarative equivalent A manifest stream gets a `DatetimeBasedCursor`: `cursor_field` names the record field, `datetime_format` says how to parse it, `start_datetime` gives the floor for a first run, `step` sets the window size, `cursor_granularity` sets the smallest unit the source distinguishes, and `lookback_window` deliberately re-reads a margin behind the stored cursor. The cursor's value is injected into requests through start-time and end-time request options, so each slice becomes a filtered API call. The framework handles checkpointing per slice. ## Inclusive bounds, duplicates and the primary key Most sources filter with `>=`, and a cursor with second granularity cannot distinguish rows written in the same second. Both mean the boundary record set is re-read on the next run. That is the right default — losing rows is worse than repeating them — but it means an incremental append destination accumulates duplicates unless the stream declares a primary key and the user picks a deduplicating sync mode. A connector that reports no primary key on a stream it expects to be deduplicated has shipped a bug. `lookback_window` is the deliberate version of the same idea, for sources where rows can be written with a timestamp slightly in the past — you accept extra re-reads to stop losing late arrivals. ## Symptoms to recognise - **Every sync is a full sync.** State was never emitted (no checkpoint interval and the stream never completes), or the setter ignores the incoming state, or the cursor field is absent from records so the maximum stays null. - **Rows go missing after a failure.** State advanced ahead of emitted records, or advanced on unordered data. - **Duplicates every run at the boundary.** Expected with inclusive filters; the fix is a primary key and dedupe, not narrowing the filter. - **State goes backwards or resets.** A schema or cursor-field change altered the blob's meaning; treat cursor changes as a reset, not a silent upgrade.

  • Why is emitting a STATE message before the records it covers a data-loss bug?
    The platform treats state as "everything up to here is handed off" and starts the next sync after it. If the process dies between the checkpoint and the records, those records are never emitted and never re-read. Always yield records first and checkpoint after them.
  • A source returns records in arbitrary order. What breaks if you checkpoint on the running maximum cursor?
    A mid-stream checkpoint can record a value higher than records still to be emitted, so a failure loses those laggards on resume. Either request an ascending sort on the cursor field, or only advance state at a boundary you can prove is complete, such as the end of a time slice.
  • Why do incremental syncs produce duplicate rows at the cursor boundary, and what fixes it?
    Sources are normally filtered inclusively and timestamps have limited granularity, so the boundary rows are read again. That is intentional — repeating beats losing. Declare a primary key on the stream so a deduplicating destination sync mode collapses the repeats rather than trying to make the filter exclusive.
  • What does a lookback window buy you and what does it cost?
    It re-reads a margin behind the stored cursor, catching rows that were written with a timestamp slightly in the past or committed after a previous read passed their timestamp. The cost is extra API calls and more duplicate rows for the destination to dedupe, so size it to the source's observed lateness.

saying these in an interview costs you the question

  • Thinking the platform interprets the state blob's contents
  • Updating the cursor before yielding the records
  • Only emitting state when the whole stream finishes
  • Trying to eliminate boundary duplicates by filtering exclusively
  • Assuming state persists if the sync crashes after reading rows

context