Mehdi Akiki
Published on

When to Advance a Sync Checkpoint After Partial Success

Authors
  • Mehdi Akiki avatar
    Name
    Mehdi Akiki
    Twitter

Article · Interrupted execution

A synchronization job fetches page 417, writes 83 of its 100 records, then the process dies. What should its checkpoint contain?

If it advances to page 418, seventeen records may never arrive. If it stays at page 417, the first 83 records will be seen again.

I choose repetition over silent loss, then design the write path to survive that repetition.

The principle is precise:

Advance a checkpoint only after every effect covered by that checkpoint is durable or safely repeatable.

The difficult words are “every effect” and “covered.” They depend on what the checkpoint means.

Define the claim before storing the value

A checkpoint can be an offset, cursor, timestamp, page number, source revision, or compound key. The value alone is not the contract.

For example, page = 417 could mean:

  • page 417 is the next page to fetch;
  • page 417 was fetched but not applied;
  • every record through page 417 was durably applied;
  • page 417 was partially applied and has a separate outcome ledger.

These meanings produce different recovery behaviour. I name the field from its claim, such as next_page_token or applied_through_offset, instead of a vague cursor.

The basic crash table

Consider one unit of source work U and its destination effects E.

fetch U → compute E → persist E → advance checkpoint

There are four important crash windows:

Crash pointRecovery seesRequired behaviour
before fetchold checkpointfetch the same unit
after fetch, before writeold checkpointfetch or replay the unit
during destination writesold checkpointrepeat the unit safely
after all writes, before checkpointold checkpointrepeat completed effects safely
after checkpointnew checkpointstart at the next unit

The fourth row surprises people. Even after the destination work succeeds, the checkpoint update can fail. Unless both writes share one transaction, recovery cannot know that the effects completed.

Protocol 1: checkpoint before applying

This order is attractive because duplicate work disappears:

fetch page 417
save next_page = 418
write page 417

It creates a loss window. If the process dies after saving 418 and before all writes finish, recovery starts after records that never became durable.

I do not use this protocol when the source cannot reconstruct the missing records by another mechanism.

Protocol 2: apply, then checkpoint

The usual at-least-once protocol is:

fetch page 417
write every record from page 417
save next_page = 418

A crash after the writes and before the checkpoint repeats page 417. This is safe when applying the same source fact again converges to the same product state.

A destination write may use a stable identity and source revision:

insert into normalized_accounts (
  tenant_id,
  provider,
  external_id,
  source_revision,
  lifecycle
)
values ($1, $2, $3, $4, $5)
on conflict (tenant_id, provider, external_id)
do update set
  source_revision = excluded.source_revision,
  lifecycle = excluded.lifecycle
where normalized_accounts.source_revision <= excluded.source_revision;

The exact comparison depends on the provider's revision semantics. The important part is that retrying an older observation cannot undo a newer one.

Sending an email or charging a card is not naturally convergent. Such effects need an idempotency record, outbox, or separate workflow state. “The database upsert is idempotent” does not make every downstream effect idempotent.

Protocol 3: commit effects and checkpoint together

If destination state and checkpoint live in the same transactional database, I can commit them atomically:

begin;

-- Apply all rows for this source unit.
-- Record an outbox event instead of sending externally here.

update sync_partitions
set next_cursor = $next_cursor,
    checkpointed_at = now()
where tenant_id = $tenant_id
  and partition_key = $partition_key
  and next_cursor = $expected_cursor;

commit;

Now recovery sees either the old state and old checkpoint, or the new state and new checkpoint.

The compare-and-set condition prevents two workers from independently advancing the same partition. A lease alone is not enough because a paused worker can resume after its lease expires. I use a generation or expected checkpoint value to fence stale writers.

This transaction cannot atomically include a third-party API, message broker, and email service. The outbox pattern converts the external effect into durable local intent, which another repeatable publisher completes.

Partial success needs a smaller commit unit

A page is a transport unit chosen by the provider. It is not automatically the right recovery unit.

If records are independent, I can track per-record outcomes:

page 417
  record A → applied
  record B → applied
  record C → quarantined: invalid enum
  record D → retryable: destination timeout

I advance past the page only when every record has reached an allowed durable state. “Quarantined with evidence” can be an allowed state if the product accepts completion with errors. “Logged and forgotten” is not.

Another option is to write the complete fetched page into an observation store first. Then the external page token can advance after the observations are durable, while a separate application checkpoint tracks normalization. This creates two honest stages:

source API → durable observations → product application
              checkpoint A           checkpoint B

The extra storage is valuable when source pages are expensive, tokens expire, or mappings need replay without another provider call.

A checkpoint must include all state needed for replay

A timestamp alone often cannot resume safely. Records can share timestamps, arrive late, or change while pagination is in progress.

A compound checkpoint may include:

{
  "watermark": "2026-09-22T08:41:00Z",
  "last_id_at_watermark": "acct_9182",
  "source_snapshot": "rev_271",
  "mapper_version": 12
}

This records ordering, tie-breaking, source consistency when available, and the interpretation used for stored results.

Opaque continuation tokens need special care. Some are short-lived and cannot be reused after a restart. Some encode a snapshot; others paginate a changing collection. I test the provider's actual contract instead of assuming that a token is a durable offset.

Kafka's small rule generalizes well

The current Kafka consumer API says a committed offset should identify the next record to consume. This avoids ambiguity about whether the named record is completed or pending.

The same convention works for API integrations:

checkpoint = first source unit whose effects are not yet proven durable

Apache Flink's fault-tolerance explanation shows the stronger distributed version: a snapshot combines source positions with operator state. It also makes clear that end-to-end exactly-once needs replayable sources and transactional or idempotent sinks.

I borrow the invariant, not necessarily the infrastructure.

Checkpoint frequency is a recovery trade-off

Checkpointing every record minimizes replay but adds write amplification. Checkpointing every thousand records is cheaper but repeats more work after failure.

I choose the interval from:

  • cost of repeating one unit;
  • maximum acceptable recovery time;
  • checkpoint storage load;
  • probability of interruption;
  • duration and validity of source snapshots or tokens;
  • cost and safety of repeated destination effects.

For a fast idempotent upsert, replaying a page may be cheap. For a slow provider with a narrow rate limit, durable observations and frequent source checkpoints may be worth the storage.

Test the protocol by killing it

Happy-path tests do not prove checkpoint safety. I add failpoints after each boundary:

  1. after fetch;
  2. after the first destination write;
  3. after the final destination write;
  4. before checkpoint commit;
  5. after checkpoint commit but before acknowledgement;
  6. after lease expiry while the old worker is still alive.

After restart, I assert two properties:

  • no eligible source fact is missing;
  • repeated effects remain within the declared semantics.

I also inspect the operator view. A checkpoint should say what is blocked, not only how far a successful run reached.

The rule I keep

A sync checkpoint is not “the last thing I saw.” It is the boundary of a durable claim.

When effects and checkpoint share a transaction, I commit them together. When they do not, I advance after effects and make replay safe. When a transport page is too large, I introduce per-record outcomes or a durable observation stage.

The system may repeat work after a crash. That is visible and repairable. Silently skipping work is neither.