Online feature publication needs a watermark, entity-level freshness checks and a repair plan that preserves event ordering.
Feature materialization lag: detect stuck updates and repair safely
Track source and publication progress separately
A materialization job can report success while reading no new source partitions. Record the latest source event time seen, latest source partition completed and latest value published online. Compare those watermarks with expected arrival patterns. A global watermark can hide a stalled region, so monitor a bounded set of meaningful slices. Request-time freshness still protects individual decisions when a dashboard is late to notice lag.
Preserve ordering under late arrival
An old merchant event arriving after a newer one belongs in historical storage but should not automatically replace the online value. Compare event time and a deterministic tie-breaker before publishing. If a correction changes a past event, recompute affected historical rows and decide whether current online state must change. Replay rules show why correcting history and changing current operations are separate transitions.
Repair a stuck slice
Pause or isolate the faulty publisher, identify the last verified partition and replay from that point into a staging view. Reconcile entity counts, duplicate keys, event-time maxima and feature values before swapping the view. A blind “refresh all” can overwrite good current values with stale records or double downstream work. Keep old and new view versions long enough to compare requests. If the source itself is late, refreshing more often does not create information that has not arrived.
Test a realistic outage
Simulate a region whose source partition stops for 82 minutes while other regions remain current. Confirm freshness fallback appears only where needed, alert ownership points to the source or publisher, and repaired values become usable after event-time checks. The receipt project measures both time to detect the lag and time to restore safe scoring, not just job success.
Implementation
def choose_online_record(current, arriving):
if current is None:
return arriving
current_key = (current["event_at"], current["revision"])
arriving_key = (arriving["event_at"], arriving["revision"])
return arriving if arriving_key > current_key else current
current = {"event_at": 82, "revision": 2, "value": 470}
late = {"event_at": 47, "revision": 9, "value": 471}
correction = {"event_at": 82, "revision": 3, "value": 472}
assert choose_online_record(current, late) == current
assert choose_online_record(current, correction) == correction
Performance and operating cost
The per-entity comparison is O(1) time and space. A repair over n source records costs at least O(n) scan time plus online writes; staged reconciliation may need O(u) state for u entities. The example assumes revision order is meaningful within an event timestamp, which the producer contract must enforce.
Common Mistakes
- Trusting a successful job status without checking source and online watermarks.
- Letting a late old event overwrite a newer online value.
- Refreshing every entity when only one partition is stuck.
- Confusing historical correction with permission to change a past production decision.
Read next
- Online feature freshness: use event and availability clocks
- Project: keep receipt scoring safe during feature publication lag
- Batch replay: supersede outputs without duplicating downstream actions
- Feature schema evolution: keep producers and rollback models compatible
- Model alerts: page on customer symptoms with a named owner
Continue the workflow: Project: recover a receipt stream scorer after offline uploads.
