intermediate 3 min answer

Review this design. Each of 18 services embeds a Kafka Streams GlobalKTable of a 9-million-row customer topic so that lookups are local, each instance holding it in RocksDB on local disk. Instances take 11 minutes to become ready. What would you remove, what would you change, and what would you leave alone even though it looks odd?

globalktablematerialised viewrestore timedeployscoupling
Show the full answer Hide the answer

What is actually required

Two questions decide most of this review. How many of the 18 services need customer data on a latency-critical path, and how much of the 9 million rows does each one touch? In most estimates of this shape, two or three services do high-rate lookups and the rest do a few hundred per second while handling requests that already have a 200 ms budget.

The relevant number is not memory. At roughly 300 bytes per row the table is about 2.7 GB; at 60 instances across the fleet that is 160 GB of duplicated state, which is affordable. The real cost is that a deploy performs 60 restores of a 2.7 GB table, and an unplanned rebalance does the same during an incident. Eleven minutes of readiness turns a rolling deploy into an hour and makes the fleet slowest exactly when it is being recovered.

What I would remove, and why it is safe

Remove the GlobalKTable from every service whose lookup rate is under a few hundred per second and whose latency budget can absorb a 1 to 2 ms round trip to a shared cache. Those services gain nothing from locality and pay the full restore cost. The test is arithmetic, not taste: if lookups per second times the latency saved is less than the restore time amortised over the deploy interval, the local copy is a loss.

The one change that matters

Attack readiness time, not storage. In order: keep warm standby replicas so a rebalance promotes rather than restores; bootstrap new instances from a periodic state snapshot instead of replaying the topic from offset zero; and check that the topic is actually compacting — a compacted topic still contains every uncompacted record in its active segment and in any segment below the cleaner's dirty-ratio threshold, so a "9-million-row" topic can easily be twice the size of its distinct key set, and the restore reads all of it.

What I would leave alone, even though it looks odd

The two services that genuinely need local lookups. Replacing their table with a remote call adds a network hop to every event, couples their availability to the cache, and replaces a predictable restore cost with an unpredictable tail. A fully replicated local table is the right tool when lookup rate is high and the data set is bounded — which is what a customer table is. Leave the odd-looking duplication where it is paying for itself.

How I would argue this in review

Not as "this is over-engineered". As a number: restore time times instance count times deploys per week is the cost, and it is charged during deploys and incidents rather than at steady state. That reframes the discussion from a preference about architecture to a measurable operational tax, and it tells the two teams that keep their tables why they are different.

When the whole pattern is wrong

If the reference data were 200 million rows, or unbounded, a replicated local table stops being an option at any restore time and the answer is a shared store with a cache. The pattern's validity is a function of the data set's size and its churn, and nobody re-checks it after the data set grows.