Mehdi Akiki
Published on

High-Water Marks, Overlap Windows, and the Records Between Them

Authors
  • Mehdi Akiki avatar
    Name
    Mehdi Akiki
    Twitter

Article · Derived state

A high-water mark says, “I have completely processed source changes up to this position.” The word completely carries most of the design.

When the source provides a durable log offset or opaque change token, the position can have a strong contract. When the source provides only updated_at, I cannot honestly claim that everything before one time is complete. Records can share timestamps, commit late or become visible out of order.

An overlap window makes a timestamp-based reader safer by rereading recent history. It does not magically turn a wall clock into a log. The complete design needs total ordering, idempotent application, an atomic checkpoint and a separate repair path.

The unsafe algorithm

read rows WHERE updated_at > saved_time
apply rows
save maximum updated_at

This skips equal timestamps and any record that appears later with an older time. Why updated_at is dangerous simulates those failures.

If time is the only source primitive, I use two positions:

  • H: the committed high-water tuple (time, stable_id);
  • L: the lower bound obtained by subtracting an overlap duration from H.time.

The window algorithm

At the start of a run:

H = committed high-water tuple
L = H.time - overlap
U = a fixed upper bound captured for this run

Then read a stable order:

SELECT *
FROM source_records
WHERE updated_at >= :lower_time
  AND updated_at <= :upper_time
ORDER BY updated_at ASC, id ASC;

If the API paginates, every page uses the same logical U where the provider supports it. The worker upserts every record idempotently. It tracks the greatest tuple observed within the run, but it does not commit the new checkpoint until all pages and durable writes succeed.

In pseudocode:

H = load_checkpoint()
L = subtract(H.time, overlap)
U = source_now_or_snapshot_boundary()
candidate = H

for page in read_pages(L, U, order = [updated_at, id]):
    transaction:
        for record in page:
            apply_idempotently(record.identity, record.version, record.payload)
            candidate = max(candidate, (record.updated_at, record.id))

transaction:
    record_run_complete(L, U, candidate)
    commit_checkpoint(max(candidate, (U, minimum_id)))

That final expression depends on the source contract. Some sources guarantee the scan covers every row through U, so the checkpoint can advance to the boundary even if no row exists there. Others only allow advancing to the maximum observed record. I write this choice into the adapter instead of hiding it in generic sync code.

A record between the marks

Assume:

committed H: 10:00:00 / id=900
overlap:     5 minutes
run window:  [09:55:00, 10:10:00]

During the previous run, record 731 was not visible. A long transaction later commits it with updated_at=09:58:00.

The strict cursor misses it because 09:58 < 10:00. The next overlap read includes it because 09:58 >= 09:55. The upsert safely sees many other records again and applies only a newer version or the same version as a no-op.

Now imagine the record appears with updated_at=09:40. A five-minute overlap still misses it. This is why overlap controls bounded lateness; it is not proof against arbitrary lateness.

Choose overlap from evidence

I measure visibility delay:

lateness = first_observed_at - source_updated_at

Then I inspect the distribution per provider, tenant and operation type. A useful overlap may exceed a high percentile plus clock uncertainty and API-cache delay.

I do not choose only from the average. One long transaction, asynchronous export or nightly backfill creates the tail that matters.

The trade-off is direct:

larger overlap → more duplicate reads and provider quota
smaller overlap → greater probability of an unseen late record

If a provider has 30-day backfills, repeatedly rereading 30 days may be worse than a daily reconciliation scan or provider change feed.

The tuple closes timestamp ties

An overlap alone does not guarantee complete pagination inside one timestamp. I require a deterministic secondary key:

(updated_at, immutable_id)

The exact next-page predicate is:

WHERE updated_at > :page_time
   OR (updated_at = :page_time AND id > :page_id)
ORDER BY updated_at, id

The API must preserve this order. If it exposes only updated_at and can return more same-timestamp records than one page without a stable page token, the endpoint cannot provide complete incremental reads. No client-side trick repairs records the server never makes reachable.

Idempotency is part of overlap

An overlap intentionally creates duplicates. A consumer that sends an email or charges a card for every observation is unsafe.

I separate observation from effect:

provider record
→ deduplicate by provider identity and source version
→ update canonical state conditionally
→ create an outbox effect only for a real state transition

The deduplication key includes the provider/tenant namespace. Version comparison prevents an old overlapped row from replacing a newer state. If the provider has no trustworthy version, I retain a content hash and still rely on reconciliation for uncertainty.

Checkpoint only after durable application

The crash cases are:

apply records → crash before checkpoint: records repeat, safe with idempotency
checkpoint → crash before apply:       records can be lost permanently

I choose repetition. The checkpoint transaction must follow durable application, or share one atomic transaction when storage permits. Durable cursor design covers this boundary and lease ownership in detail.

Reconciliation covers what the window cannot

No finite overlap handles an unbounded delayed backfill, hard deletion without a tombstone, source repair that moves time backward, or a provider bug that omitted a page.

I add a slower loop:

  • periodic full or partitioned comparison;
  • counts plus identity/checksum comparison, not counts alone;
  • explicit tombstone or deletion feed;
  • quarantined repair queue with evidence;
  • ability to reset from a provider-supported snapshot;
  • freshness and oldest-unreconciled-age metrics.

The fast incremental loop gives freshness. Reconciliation gives eventual repair.

What I record per run

{
  "provider": "calendar-x",
  "tenant": "tenant-41",
  "lowerBound": ["2026-09-25T09:55:00Z", ""],
  "upperBound": "2026-09-25T10:10:00Z",
  "previousHighWater": ["2026-09-25T10:00:00Z", "900"],
  "committedHighWater": ["2026-09-25T10:10:00Z", ""],
  "recordsRead": 1842,
  "duplicates": 611,
  "newOrChanged": 23,
  "maximumObservedLatenessSeconds": 173,
  "pages": 19,
  "status": "complete"
}

This record makes an apparently stuck watermark, growing lateness or expensive overlap visible. It also explains which interval a failed run attempted.

My practical rule

A high-water mark is a claim about completeness, not merely the maximum value seen.

With a timestamp source, I use a stable tuple, a fixed run boundary, a measured overlap, idempotent version-aware upserts and checkpoint-after-apply. Then I run reconciliation for failures no finite overlap can cover.

The best improvement remains moving to a source change token, revision or log position with a documented ordering contract. Until then, I keep the uncertainty visible in the algorithm and the metrics.