Retailer Catalog Snapshot Ingestion
Published 5 October 2026
Retailers upload a full snapshot of their product catalog every day. The system
ingests each snapshot into a product catalog (rows keyed by retailer_id and
product_id) and serves reads from it. Each snapshot gets a generation number; a
newer snapshot supersedes an older one, and products missing from the newest snapshot
are deleted.
These notes pick up partway through a mock interview, after the basic upload → queue → workers → Postgres pipeline was in place.
Architecture¶
Where the design ends up after the discussion below:
Click the diagram to pause or play.
- A retailer uploads a full snapshot. The upload API stores the file in S3 next to the previous one, and records the new generation as the retailer's latest.
- Diff workers merge-walk the old and new files (both sorted by
product_id) and keep only inserts, updates and deletes, usually 1–5% of the catalog. - The changes are chunked into messages of at most 5,000 rows, on a small-retailer or a large-retailer queue.
- Apply workers process at most 20 chunks per retailer at a time. Before every batch they check that the generation is still the latest, then write with the generation as a fencing token, so a stale or zombie worker can't overwrite newer data.
- Postgres is sharded by
retailer_id. The read service finds the shard for a bareproduct_idthrough a lookup and reads from that shard's replicas.
Superseded snapshots: check before every batch¶
The first idea was to have a worker drop a message if a newer snapshot for that retailer already exists. That check is right, but it has to run before every batch, not just when the worker picks up the message.
Say snapshot S2 arrives while S1 is already 40% done. A check at pickup alone never catches it, because S1's worker checked long before S2 existed.
Zombie workers happen on any queue¶
SQS FIFO does guarantee order within a message group, so ordering isn't the real problem. The problem is zombie workers, and switching to Kafka doesn't remove them: a consumer that stalls past its session timeout gets kicked out, its partition goes to another consumer, and then the old one wakes up and keeps writing.
Any queue can end up with two live writers. Fencing is what makes that harmless.
Fencing with the generation number¶
Fencing comes from distributed locking. A worker holding a lease freezes (GC pause, network stall, slow disk), the lease expires, a second worker takes over and writes newer data, and then the first one wakes up still believing it holds the lock. It can't detect this itself, since from its point of view no time passed, so the storage has to reject it. Every lock grant carries a monotonically increasing fencing token, every write carries the token, and the storage refuses any write whose token is lower than the highest it has already seen.
In this design, the generation number is the fencing token. The product row remembers the highest generation that wrote it, and every write is conditional:
UPDATE products
SET retailer_attributes = :attrs, snapshot_generation = :gen
WHERE retailer_id = :r AND product_id = :p
AND snapshot_generation < :gen;
Replaying the bad case:
- S2 (generation 8) writes product X. The row is now stamped 8.
- A zombie S1 worker (generation 7) wakes up and tries to write X.
- The condition
8 < 7is false, so zero rows are updated. The stale write is silently dropped.
The order the writes land in no longer matters, because the row itself decides who wins. The "drop the message if a newer version exists" check is an optimization that saves wasted work; fencing is the correctness guarantee for the cases that check misses.
The same idea protects the delete phase: DELETE ... WHERE snapshot_generation < :gen
can never remove rows a newer snapshot has touched.
Scale¶
Question: 100K retailers upload full snapshots daily, with about 22K read QPS.
- Estimate the daily write load and say where a single Postgres leader breaks first.
- What do you shard or partition on, and how do you avoid rewriting unchanged products?
- How do you stop one giant retailer's 5M-row upload from starving thousands of small retailers?
My answer¶
- Load: 100K retailers × ~1,000 products = 100M rows a day, about 1,157 writes/sec; 4× for peak days gives ~4,000 writes/sec.
- Sharding: on
retailer_id, so a retailer's rows land on one shard. Not onproduct_id, which would make batch uploads hard. - Unchanged products: some kind of diffing to find changed rows — deferred.
- Fairness: split the 5M rows into uniform chunks of at most 5,000 rows per queue message so workers pick up evenly sized work. A thousand chunks can still starve others, so use two queues, one for big-catalog retailers and one for small ones, with more partitions and bigger workers on the big queue.
Feedback¶
The estimate is off by 20×. The average catalog is about 20,000 products, not 1,000:
- 100K retailers × 20K products ≈ 2 billion product writes a day, about 23K writes/sec on average.
- Uploads aren't spread evenly. Most retailers export overnight, so peaks of 50–100K writes/sec are realistic.
- At ~2 KB per row that's about 4 TB of row writes a day, before indexes and write-ahead log, which multiply it further.
A single Postgres leader can't take that, and that should be stated plainly. It breaks on write-ahead log and replication throughput first, then on vacuum: every rewritten row leaves a dead tuple, and the table bloats faster than vacuum can clean it.
Diffing isn't "later", it's the main answer. Most products don't change day to day; maybe 1–5% do. Writing only changed rows turns 2B writes a day into 20–100M. Nothing else buys a 20–100× reduction. Two ways to get it:
- Content hash per product: skip the write when the hash matches. This still needs a read per product, but reads are much cheaper than writes plus WAL plus vacuum.
- Diff the files themselves: keep the previous snapshot in S3, sort both files by
product_id, and merge-walk them in the workers. Only inserts, updates and deletes reach the database. This also answers "how do you compare two snapshots", and it moves the heavy work off the database onto cheap, horizontally scalable compute.
Sharding on retailer_id is a good call — browse and ingestion both stay on one
shard. But the other read pattern was skipped: "get product by product_id" arrives
without a retailer_id, so which shard do you query? Either make product_id encode the
retailer or shard, or keep a separate product_id → retailer_id lookup. Skew is
manageable: even a 5M-product retailer is only about 10 GB, which fits comfortably on
one shard.
Fairness: uniform chunks and size-split queues are both solid. To make it stronger, cap in-flight chunks per retailer, say 20 at a time. The giant retailer gets steady throughput but can never occupy the whole worker pool, and small retailers always find free workers. This is per-tenant fair scheduling — naming it explicitly scores well.
Follow-up: observability¶
Assume the system is built: upload, diff, chunked workers with fencing, sharded Postgres, read replicas, and a read service.
Question: How would you know it's working correctly?
- What metrics and alerts would you set up for ingestion and for serving?
- How would you detect silent correctness bugs, where the pipeline reports success but the catalog doesn't match the retailer's snapshot? Consider the fencing and delete bugs above.
- A retailer emails: "I uploaded my file six hours ago and my products still show the old price." What would you look at, in order, to find out why?
Not answered yet.