Training on Served Features Makes Every Deploy a Training-Data Write
Uber's fix relocates the cost into stream-join state and gives a bad model version a direct path into retraining, so version tags and eviction metrics must ship with it.
The expensive part is the join, not the log
Uber's reported volume numbers explain why few teams have tried this. Logging every feature for every scored candidate would run roughly 1.7 PB/day, so Uber stacked reductions. A Flink join of predictions against impression events keeps only the ~5% of candidates users are actually shown. Feature allow lists and integer aliases for feature names cut the payload a further 4–5x. Composed multiplicatively, logged output lands near 17–21 TB/day (our derivation from the stated reductions).
Buffered join state grew by more than 12 TB per hour under default Flink/RocksDB checkpointing, and that is the line item that decides the design. A prediction×impression join holds every prediction until its impression arrives or the window closes, and about 95% never match. If the rate holds, 288 TB/day. Uber swapped the default for custom state handling and aggressive eviction. For sizing elsewhere, state scales as prediction rate × logged payload × join window, and the eviction TTL decides whether the design is affordable.
Eviction carries a statistical cost worth pricing separately. An impression arriving after its prediction was evicted becomes a lost label. If lateness correlates with slow devices, poor networks or particular regions, the loss is non-random and the labels tilt toward fast paths. Uber did not publish the drop rate.
Why this makes rollback harder
SRE Weekly featured Balu Kambala's argument, built on a CircleCI example, that rolling back code does not erase state that has already spread downstream. In an ML system, prediction logs that feed the next training set are one of those surfaces. A registry rollback restores the weights and stops there.
Uber's design puts that surface at the centre. Once logged serving values are the training source of truth, every deploy writes training data. A bad model version decides which candidates become impressions. A buggy upstream feature at serving time is logged faithfully as “what the model saw.” That interaction is our own inference. Those rows reach next week's training set unless each carries a model_version tag and a quarantine query can drop the deploy window.
Both sources converge on a second rule. SRE Weekly's real-time pricing example runs Kafka → Redis Pub/Sub → .NET Channels → SSE at 1,900+ messages per second. Redis Pub/Sub delivers each message at most once and does not persist it, so any feature derived after that fan-out diverges from the log whenever a subscriber disconnects. The divergence resurfaces later as unexplained train/serve skew. Training features need a durable, replayable record of what served, which is what Uber built.
What the design takes away
- Schema freedom. With an allow list, a feature not logged today has no history tomorrow. New-feature experiments slow from an afternoon's backfill to weeks of waiting unless an offline backfill path stays open for candidate features.
- Off-policy data. Impressions-only logging discards the ~95% of unshown candidates, which are the raw material for off-policy evaluation and exposure-bias correction (reweighting for items users were never shown). A small random slice of unshown candidates keeps both options cheaply.
Training on what the model saw removes skew. It also makes every deploy, good or bad, the author of the next training set.
The sensible order is measure, tag, then pilot. A one-week skew audit says whether the mismatch is anywhere near Uber's pre-fix level before anyone commits to running streaming join state.
What to do
Run a one-week skew audit this sprint: shadow-log served values for your top 20–50 ranking features on ~1% of requests, recompute them with your offline ETL, and report the mismatch rate per feature.
Add a model_version tag to every row your serving path logs, and build a quarantine query that drops any deploy window from the next training set, before a logging pilot starts.
Size Flink join state as prediction rate × logged payload × join window before any pilot this quarter, set explicit eviction TTLs, and track the late-impression drop rate by device and region from day one.