A classifieds and auction marketplace at eBay's listing volume runs CQRS - writes to Postgres and a search-shaped read model built from an append-only change log that now holds 1.2 billion events across 24 partitions. A projection bug has been live for nine days. Work out roughly how long a full rebuild takes and say what that number decides.
Show the full answer Hide the answer
The assumptions, stated
Rebuild time is not the log's read throughput. Reading 1.2 billion compact events off sequential storage at a few hundred megabytes a second is tens of minutes, and that number is a distraction. The binding constraint is the write capacity of the read store, because every event turns into one or more indexed upserts.
Two planning figures carry the estimate, both order-of-magnitude:
- A single projection worker doing one indexed upsert per event against a relational or document store sustains on the order of 2,000 to 10,000 events per second. Call it 5,000.
- A batched sink - a search index taking bulk requests, or a columnar store - sustains roughly an order of magnitude more, so 20,000 to 50,000 per worker.
The arithmetic
Single worker, unbatched: 1.2e9 ÷ 5,000 = 240,000 seconds, about 67 hours.
The log has 24 partitions, so 24 workers is the parallelism ceiling without resharding. Perfect scaling gives 1.2e9 ÷ 120,000 = 10,000 seconds, under three hours. The read store will not absorb 120,000 writes a second while also serving live traffic. If you are willing to spend 30,000 writes a second of its capacity on the rebuild, the answer is 1.2e9 ÷ 30,000 = 40,000 seconds, about 11 hours.
So the honest answer is a range: 11 to 67 hours, and where you land inside it is a choice about how much of the read store's capacity you are willing to take away from live queries.
Which assumption dominates the error
Not the log, and not the worker count. It is the fan-out from one event to index writes. If a ListingUpdated event touches the listing document, a seller-aggregate document and two facet counters, the per-event cost is four writes and every number above divides by four. Measure that fan-out on a thousand events before quoting a rebuild time; it is the single input most likely to be wrong by a factor of three.
What the number decides
If the rebuild takes longer than the time your business will accept a wrong or missing read model, CQRS is not reversible in place, and you have to own a second projection. That is the whole decision.
- Rebuild under an hour: rebuild in place. Point the projector at offset zero over a maintenance window and accept degraded search.
- Rebuild over a few hours: build the corrected projection into a new index beside the live one and flip a pointer when it catches up. The permanent cost is double read-store storage and a deploy path that can address two index versions. The benefit is that rebuild stops being an outage and becomes a background job, which is what also makes a read-model schema change routine.
- Either way, 1.2 billion events is a signal to add periodic snapshots so the rebuild starts from a checkpoint rather than from 2019.
The deeper finding here is not the rebuild. Nine days of silent drift is the failure. A reconciliation job sampling a thousand random keys a minute and alerting on the first mismatch rather than on a rate would have caught it within the hour for a fraction of a core.
When a full rebuild is the wrong answer
If the bug affected a known subset - events in one window, or one event type - and the projection handler is idempotent per key, replay only the affected keys. A bug touching 400,000 listings out of 90 million is minutes of work. The precondition is being able to get from a key to its events, which means a key-partitioned log or an index from key to offsets. Teams that discover they lack this during an incident rebuild everything because it is the only operation they have, which is the argument for building the targeted path first.