advanced 3 min answer

A team replays 60 days of order events through the current job to correct a tax calculation. Side effects are already handled - the replay writes to a shadow table and emits no notifications. The job enriches each order by calling the pricing service for the item's list price. What does the shadow table contain when the replay finishes?

replaydeterminismenrichmenttemporal-joinreprocessing
Show the full answer Hide the answer

What happens, step by step

The replay reads history at perhaps 50 times live rate. For each order it asks the pricing service for a price, and the service answers with today's price. Sixty days of orders are repriced at current prices, and the corrected tax is computed on a base that no customer was ever charged.

The output is therefore wrong in a way that is not visible as wrongness. Every row is populated, every type checks, the tax column finally looks right, and the totals differ from the live table for two reasons at once - the fix and the reprice - so nobody can attribute the difference to either.

Where it amplifies

  • Read amplification into a dependency sized for live traffic. Sixty days of lookups compressed into a few hours is the speed-up factor applied to the pricing service's request rate. Its cache hit rate also collapses, because 60 days of history touches a far wider key set than the live working set, so the calls land on its database.
  • A reconciliation nobody can close. The comparison between shadow and live is now two-dimensional, and the usual outcome is that the team promotes the shadow table anyway because the column they were fixing looks correct.
  • No self-healing. Duplicate output leaves evidence. Non-determinism leaves none. There is no signal anywhere that distinguishes "priced then" from "repriced now", and no later run recovers the original value once the pricing service has moved on.

What stops it

The enrichment has to come from the log, not from a live service. Three concrete forms, in increasing cost:

  1. Put the price in the event. The producer writes the price it used at order time. Cost: larger events and a schema that carries derived data, which some teams resist on purity grounds and which is almost always the right trade.
  2. Publish prices as a versioned changelog topic and do an as-of join on event time. Cost: the job holds price state keyed by item and must retain versions covering the whole replay window.
  3. Read an immutable point-in-time store keyed by entity and valid-time, so the job asks "what was the price at this event's timestamp".

The general rule is worth stating plainly: a job is replayable only if its output is a pure function of the log. Any read of mutable outside state breaks it - a service call, now(), a feature store without as-of semantics, a random seed, a config value read at startup. The Kappa promise that you can reprocess by replaying is a claim about the purity of the job, not about the durability of the log, and teams adopt the log and skip the purity.

Because there is no runtime signal, the enforcement has to be structural: no network reads in the transform stage, checked in review or by an architecture test on the enrichment path.

When this is the wrong answer

Sometimes current enrichment is exactly what you want. Rebuilding a search index or re-scoring historical orders against today's catalogue should call the live service, and an as-of join would produce a stale index. Decide which of the two replays you are running before you start, write it down, and accept that one code path cannot serve both - the correction run and the rebuild run need different enrichment.