A sales ranking system answers one question for every product category: "what are the best-selling items right now?". You have seen the output on shopping sites as "Best Sellers in Headphones" or a "#3 in Running Shoes" badge on a product page. Under the hood it is a counting and sorting problem over a never-ending stream of orders.
Interviewers like this problem because it has a crisp core (top-K per category) wrapped in real engineering choices: batch or streaming, exact or approximate counting, how fresh the ranking must be and what that freshness costs. A good answer moves comfortably between data structures (heaps, sketches) and data pipelines (MapReduce-style jobs, stream processors).
What interviewers typically probe:
- How you count sales over a rolling window, such as "the last 24 hours", rather than all time.
- How you compute top-K per category efficiently when products belong to several categories.
- Batch versus streaming, and when approximate counting (a count-min sketch) is justified.
- How rankings are served at high read volume with low latency.
Background that helps: Message queues, Kafka and event streaming and Caching.
Problem and scope
Build a service that:
- Consumes order events from the checkout system.
- Maintains sales counts per product over a rolling window.
- Produces, for every category, the top 100 products by units sold.
- Serves those rankings to category pages, product pages ("#3 in Headphones") and internal tools.
Out of scope: personalised recommendations, search ranking, pricing, and fraud detection itself (though we will filter out cancelled and fraudulent orders).
Clarifying questions
| Question | Why it matters | Assumption here |
|---|---|---|
| Rank by units sold, revenue or a blended score? | Changes what you aggregate | Units sold, with returns and cancellations subtracted |
| What window? All time, last 24 hours, last 7 days? | Determines state size and expiry logic | Last 24 hours, updated hourly; a 7-day variant later |
| How fresh must rankings be? | Batch hourly vs streaming in seconds | Within 1 hour is acceptable; a few minutes is nice |
| How many products and categories? | Memory and output size | 50 million products, 20,000 categories |
| Can a product be in several categories? | Fan-out per sale | Yes: a leaf category and its ancestors, about 4 on average |
| Must counts be exact? | Allows approximation | Rank order must be right at the top; small count errors are fine |
| Read volume? | Serving design | About 1 billion ranking reads per day |
Interview tip
Ask "how fresh?" early and tie your design to the answer. "Within an hour" lets you start with a simple batch job. "Within a minute" pushes you to streaming. Saying this out loud shows you understand that freshness is a cost dial, not a free requirement.
Functional and non-functional requirements
Functional:
- Ingest order line items (
order_id,product_id,quantity,timestamp, status changes). - Compute units sold per product over the last 24 hours, per category.
- Return the top K (up to 100) products for any category.
- Return a single product's rank in each of its categories.
- Exclude cancelled, refunded and flagged orders.
Non-functional:
- Freshness: rankings no more than 1 hour stale (target: minutes).
- Read latency: under 50 ms at p99 (99th percentile).
- High availability for reads: a stale ranking is far better than an error on a category page.
- Correctness under retries: the same order event delivered twice must not count twice.
- Cost-efficient: rankings are a merchandising feature, not a payment system.
Back-of-the-envelope estimates
Write side (order events).
- 30 million orders per day × 1.5 items per order = 45 million line-item events per day.
- 45 × 10⁶ ÷ 86,400 ≈ 521 events per second on average.
- A big sale day can be 10× higher: about 5,200 per second at peak.
- Each product belongs to about 4 categories (leaf plus ancestors), so category-level updates are about 4 × 521 ≈ 2,083 per second on average.
Event storage. At about 200 bytes per event, 45M × 200 B = 9 GB per day, about 3.3 TB per year. Easy for object storage or a data warehouse.
Counting state. Keep 24 hourly counters per active product. If 10 million products sell at least once in a day: 10M × 24 × 4 bytes ≈ 0.96 GB. Even all 50 million products would need about 4.8 GB. Exact counting fits in memory on a handful of machines.
Output size. 20,000 categories × 100 entries × 16 bytes (8-byte product ID + 8-byte count) = 32 MB. The entire ranking dataset fits in memory on every serving node.
Read side. 1 billion ranking reads per day ÷ 86,400 ≈ 11,600 per second on average and about 35,000 per second at a 3× peak. Since the data is only 32 MB and changes at most every few minutes, this is a caching problem, not a database problem.
What the numbers tell you
Writes are modest (hundreds to thousands per second). Reads are higher but hit tiny, slowly changing data. State fits in memory. So the hard parts are not raw scale: they are windowing, deduplication, category fan-out, and the freshness-versus-cost choice.
API design
GET /v1/rankings/{category_id}?limit=50
-> {
"category_id": "headphones",
"window": "24h",
"computed_at": "2026-10-10T09:00:00Z",
"items": [
{"rank": 1, "product_id": "P3", "units": 1840},
{"rank": 2, "product_id": "P1", "units": 1702}
]
}
GET /v1/products/{product_id}/ranks
-> {"ranks": [{"category_id": "headphones", "rank": 3},
{"category_id": "electronics", "rank": 412}]}
Notes:
computed_attells clients and caches how fresh the data is, and it is useful when debugging "why did my product drop?".- The ingest side has no public API. Checkout publishes
OrderPlaced,OrderCancelledandOrderRefundedevents to an event log, and the ranking pipeline consumes them. - Many sites show exact counts only internally and show only the rank publicly; counts reveal business data.
Data model and storage choice
| Data | Store | Why |
|---|---|---|
| Raw order events | Event log (Kafka-style) for days, plus object storage or a warehouse for years | Replayable, cheap, append-only |
| Product-to-category mapping | Relational database, cached in memory by the pipeline | Small, changes rarely |
| Rolling counters | In-memory state of the stream processor, checkpointed; or hourly aggregate tables | Fast updates, recoverable |
| Hourly aggregates | Columnar table (hour, product_id, units) | Batch jobs and backfills |
| Published rankings | Key-value store rank:<category> and an in-process cache | Tiny, read-heavy |
The hourly aggregate table is the backbone of the batch approach:
CREATE TABLE sales_hourly (
hour TIMESTAMP NOT NULL, -- truncated to the hour
product_id BIGINT NOT NULL,
units INTEGER NOT NULL,
PRIMARY KEY (hour, product_id)
);
Published rankings are written as one versioned value per category, for example rank:headphones:v20261010T09 with a pointer rank:headphones -> v20261010T09. Readers always see a complete ranking, never a half-written one.
High-level design
+----------+ OrderPlaced / Cancelled / Refunded
| Checkout |------------------+
+----------+ |
+------v------+
| Event log | (partitioned by product_id)
+--+-------+--+
| |
stream path | | batch path
| |
+------------v--+ +-v-------------------+
| Stream counter| | Archive to object |
| (dedupe, hour | | storage, hourly |
| buckets) | | MapReduce/Spark job |
+-------+-------+ +----------+-----------+
| |
+---------+-----------+
|
+----------v-----------+
| Top-K per category |
| (min-heaps, merge) |
+----------+-----------+
| publish versioned lists
+----------v-----------+
| Ranking store (KV) |
+----------+-----------+
|
+----------v-----------+ +-----+
| Ranking API + local +----->| CDN |
| in-memory cache | +-----+
+----------------------+
You can build either path alone. A common real-world shape is both: a streaming path for freshness, and a batch path that recomputes the authoritative answer from raw events and corrects drift. That pairing is often called a lambda architecture. A kappa architecture keeps only the streaming path and handles corrections by replaying the log. See Kafka and event streaming.
Request flows
Flow 1: an order is placed
- Checkout commits the order and publishes
OrderPlacedwithorder_id, line items and timestamp (via an outbox, so the event is not lost if publishing fails right after the commit). - The event lands in a log partition chosen by
product_id, so all events for one product go to the same counter worker. - The stream worker checks a deduplication set keyed by
(order_id, line_item_id). Already seen? Skip it. - It adds
quantityto the product's counter for that event's hour. - Periodically (say every minute) the worker emits updated window totals for products that changed.
- The top-K stage maps each product to its categories and updates per-category top lists.
- Every few minutes, new top-K lists are published as a new version.
Flow 2: an order is cancelled
- Checkout publishes
OrderCancelledwith the sameorder_idand items. - The worker subtracts the quantity from the bucket of the original order's hour (the event carries the original timestamp), if that hour is still inside the window.
- Rankings reflect the cancellation at the next publish.
Flow 3: a shopper opens "Best Sellers in Headphones"
- The page calls
GET /v1/rankings/headphones?limit=50. - A CDN or the API's local in-memory cache answers most requests (TTL of a minute or two).
- On a miss, the API reads
rank:headphonesfrom the key-value store and caches it. - The page fetches product details (title, price, image) from the catalogue service by ID.
Deep dive 1: counting over a rolling window
A rolling (sliding) window means "the last 24 hours as of now". At 10:00 it covers 10:00 yesterday to 10:00 today; at 11:00 the oldest hour drops out and the newest is added. You cannot simply keep one counter per product, because you must know what to subtract when old sales leave the window.
Hourly buckets
Keep one counter per product per hour, and a running total of the window:
- On a sale: add to the current hour's bucket and to the total.
- When a new hour starts: for each product with a count in the hour that just left the window, subtract that bucket from the total and delete it.
This is called a bucketed sliding window. Its accuracy is one bucket: the window is really "the last 24 complete hours plus the current partial hour", which is fine for rankings. Smaller buckets (5 minutes) give smoother windows at the cost of 12 times as many counters.
Worked example: a 3-hour window
Category "headphones", window of 3 hourly buckets.
| Hour | P1 | P2 | P3 |
|---|---|---|---|
| 10:00 | 4 | 0 | 0 |
| 11:00 | 7 | 8 | 0 |
| 12:00 | 2 | 0 | 5 |
| 13:00 | 0 | 0 | 6 |
At 12:xx, the window holds hours 10, 11 and 12:
- P1 = 4 + 7 + 2 = 13
- P2 = 0 + 8 + 0 = 8
- P3 = 0 + 0 + 5 = 5
Top 2: P1 (13), P2 (8).
At 13:xx, hour 10 slides out (subtract it) and hour 13 comes in:
- P1 = 13 − 4 + 0 = 9
- P2 = 8 − 0 + 0 = 8
- P3 = 5 − 0 + 6 = 11
Top 2: P3 (11), P1 (9). P3 rose to first place, and P1 dropped not because it sold less recently but because its strong 10:00 hour aged out. That is exactly what a rolling window should do.
A runnable sketch
The Python below implements the bucketed window and per-category top-K with a min-heap. It reproduces the worked example. A production system shards this by product and keeps per-category structures rather than scanning all totals, but the logic is the same.
import heapq
from collections import defaultdict
class RollingCategoryRanker:
"""Exact sales counts over the last `window_hours` hours, per category."""
def __init__(self, window_hours=24):
self.window = window_hours
# buckets[hour][(category, product)] -> units sold in that hour
self.buckets = defaultdict(lambda: defaultdict(int))
self.totals = defaultdict(int) # (category, product) -> window sum
def record(self, hour, category, product, units):
self.buckets[hour][(category, product)] += units
self.totals[(category, product)] += units
def advance_to(self, hour):
"""Drop every bucket that has slid out of the window."""
for old in [h for h in self.buckets if h <= hour - self.window]:
for key, units in self.buckets.pop(old).items():
self.totals[key] -= units
if self.totals[key] == 0:
del self.totals[key]
def top_k(self, category, k):
heap = [] # min-heap of (units, product)
for (cat, product), units in self.totals.items():
if cat != category:
continue
if len(heap) < k:
heapq.heappush(heap, (units, product))
elif units > heap[0][0]: # beats the weakest of the top k
heapq.heapreplace(heap, (units, product))
return sorted(heap, reverse=True)
r = RollingCategoryRanker(window_hours=3)
r.record(10, "headphones", "P1", 4)
r.record(11, "headphones", "P1", 7)
r.record(11, "headphones", "P2", 8)
r.record(12, "headphones", "P1", 2)
r.record(12, "headphones", "P3", 5)
r.advance_to(12)
print(r.top_k("headphones", 2)) # [(13, 'P1'), (8, 'P2')]
r.advance_to(13) # hour 10 slides out
r.record(13, "headphones", "P3", 6)
print(r.top_k("headphones", 2)) # [(11, 'P3'), (9, 'P1')]
Common mistake
Do not bucket by the time the event was processed. Bucket by the time the order was placed (event time). Otherwise a pipeline delay of 30 minutes shifts sales into the wrong hour, and replaying the log after an outage produces different rankings. Allow a grace period for late events, then close the bucket.
Deep dive 2: top-K per category with heaps
Given counts for n products in a category, you want the k largest. Sorting all n costs O(n log n). A min-heap of size k does better: a min-heap is a tree-shaped structure where the smallest element is always at the root and can be read in O(1) and replaced in O(log k).
Algorithm:
- Push the first k products.
- For each remaining product, compare its count with the heap's minimum (the root). If it is larger, replace the root; otherwise skip it.
- At the end the heap holds the top k; sort those k for display.
Worked example: top 3
Counts: P1 = 50, P2 = 20, P3 = 80, P4 = 35, P5 = 65, P6 = 10.
| Step | Product | Heap after step (min first) | Action |
|---|---|---|---|
| 1 | P1 50 | 50 | push |
| 2 | P2 20 | 20, 50 | push |
| 3 | P3 80 | 20, 50, 80 | push (heap full) |
| 4 | P4 35 | 35, 50, 80 | 35 > 20, replace |
| 5 | P5 65 | 50, 65, 80 | 65 > 35, replace |
| 6 | P6 10 | 50, 65, 80 | 10 < 50, skip |
Result: P3 (80), P5 (65), P1 (50).
| Operation | Time |
|---|
Finding the top k of n counts
Distributing it
Counts are sharded by product, so one category's products are spread across many workers. Use a two-stage top-K:
- Each worker computes a local top-K per category from its own products.
- A merge stage combines the local lists per category and takes the global top-K.
This is exact: any product in the global top-K must be in the top-K of the worker that owns it, because that worker holds its complete count. (The trick only works because each product's count lives on exactly one worker. If counts for the same product were split across workers, a local top-K could drop a product whose combined count is large.)
With 20,000 categories × 100 entries per worker, the merge input stays small.
Category fan-out
A sale in "Over-ear Headphones" also counts toward "Headphones", "Audio" and "Electronics". Each event fans out to about 4 category keys. Load the product-to-category map into each worker's memory (it changes rarely) and refresh it periodically. When a product moves category, you can either recompute from hourly aggregates or accept a 24-hour transition.
Deep dive 3: approximate heavy hitters with a count-min sketch
Our product space is small enough for exact counts. But interviewers often ask, "what if there were billions of distinct keys and you had little memory per worker?" That is where approximate counting comes in, for example ranking trending search queries or URLs.
A count-min sketch is a small 2-D array of counters with d rows and w columns, plus d independent hash functions (one per row).
- Add item x: for each row i, increment
counter[i][h_i(x)]. - Estimate count of x: take the minimum of
counter[i][h_i(x)]across rows.
Different items can land in the same cell (a collision), which only ever adds extra counts. So the estimate is never below the true count; it can only overestimate. Taking the minimum across rows picks the row with the fewest collisions.
Worked example
Width w = 5, depth d = 2. Hash functions (deliberately simple):
- h1(x) = x mod 5
- h2(x) = (x div 5) mod 5, where "div" is integer division
Stream of product IDs: 11, 23, 11, 37, 42, 11, 58, 23, 11, 42, 37, 11.
True counts: 11 → 5, 23 → 2, 37 → 2, 42 → 2, 58 → 1.
Hash positions:
| ID | h1 | h2 |
|---|---|---|
| 11 | 1 | 2 |
| 23 | 3 | 4 |
| 37 | 2 | 2 |
| 42 | 2 | 3 |
| 58 | 3 | 1 |
Counters after processing all 12 items:
col 0 col 1 col 2 col 3 col 4
row 1 (h1): 0 5 4 3 0
row 2 (h2): 0 1 7 2 2
Row 1, column 2 is 4 because 37 (2) and 42 (2) collide there. Row 2, column 2 is 7 because 11 (5) and 37 (2) collide.
Estimates (minimum of the two cells):
| ID | Row 1 cell | Row 2 cell | Estimate | True |
|---|---|---|---|---|
| 11 | 5 | 7 | 5 | 5 |
| 23 | 3 | 2 | 2 | 2 |
| 37 | 4 | 7 | 4 | 2 |
| 42 | 4 | 2 | 2 | 2 |
| 58 | 3 | 1 | 1 | 1 |
Four estimates are exact. Product 37 is overestimated (4 instead of 2) because it collided in both rows. Real sketches use far wider rows, so such double collisions are rare.
Sizing a sketch
The standard guarantee: with w = ⌈e/ε⌉ and d = ⌈ln(1/δ)⌉, every estimate exceeds the true count by at most ε × N with probability at least 1 − δ, where N is the total number of events and e ≈ 2.718.
For N = 45 million events per day, ε = 0.00001 and δ = 0.01:
- w = ⌈e / 0.00001⌉ = ⌈271,828.2⌉ = 271,829 columns
- d = ⌈ln 100⌉ = ⌈4.61⌉ = 5 rows
- Memory = 5 × 271,829 × 4 bytes ≈ 5.4 MB
- Maximum overcount ≈ 0.00001 × 45M = 450 units (with 99% probability)
An error of 450 is negligible for a product selling tens of thousands, but large for a niche category whose leader sells 50 a day. This is why the sketch suits global heavy hitters, not small categories.
Sketch plus heap
A sketch only answers "how many for x?". To find the top items, pair it with a min-heap of size k: after adding each item, estimate its count and, if it beats the heap minimum, insert or update it. The heap holds candidate IDs; the sketch holds approximate counts for everything.
For a sliding window, keep one sketch per hour and add them cell by cell (sketches with the same width, depth and hash functions can be merged by element-wise addition). Expiring an hour means dropping its sketch.
Interview tip
Say "exact counts fit in about 1 GB here, so I would count exactly, and reach for a count-min sketch only when the key space or memory budget makes exact counting impractical." Reaching for a sketch you do not need is a common over-engineering signal.
Deep dive 4: batch versus streaming, and freshness versus cost
Batch with MapReduce-style jobs
MapReduce is a model for processing large datasets in parallel: a map step transforms each record into key-value pairs, the framework shuffles (groups) pairs by key across machines, and a reduce step combines each group.
Hourly job, reading the last hour of archived events:
map(event):
if event.status is valid:
emit((event.hour, event.product_id), event.quantity)
reduce((hour, product_id), quantities):
emit(hour, product_id, sum(quantities)) -> sales_hourly
Then a second job per publish:
map(row in sales_hourly where hour in last 24):
for category in categories_of(row.product_id):
emit(category, (row.product_id, row.units))
reduce(category, pairs):
totals = sum units per product
emit(category, top_k(totals, 100)) -> ranking store
Pros: simple, easy to reason about, exact, and easy to rerun after bugs (just recompute from raw events). Cons: freshness is bounded by the schedule plus the run time, typically tens of minutes to an hour or more.
Streaming
A stream processor (Flink, Kafka Streams, Spark Structured Streaming and similar) keeps counters in memory, updates them per event, and checkpoints state to durable storage so a crashed worker can resume. Rankings can be seconds to minutes fresh.
Cons: more moving parts (state management, watermarks for late events, exactly-once or deduplicated processing), and bugs corrupt live state, so you still want a way to recompute from the log.
Choosing
| Freshness target | Approach | Relative cost |
|---|---|---|
| Daily | One nightly batch job | Lowest |
| Hourly | Hourly batch over hourly aggregates | Low |
| A few minutes | Streaming counters, publish every few minutes | Medium |
| Seconds | Streaming with per-event updates and push to caches | Highest; rarely worth it for best-seller lists |
The best-seller list on a shopping site barely changes minute to minute. A sensible default is streaming counters with publication every 5–15 minutes, plus a nightly batch recompute that corrects any drift. If the team is small, start with the hourly batch alone.
Scaling and bottlenecks
| Bottleneck | Cause | Mitigation |
|---|---|---|
| Hot product on a flash sale | One partition receives most events | Pre-aggregate in the producer or a first stage (local combine per second), or split the key into sub-keys and sum |
| Category fan-out | Each sale updates ~4 categories | Local top-K per worker, then merge; only emit changed products |
| Merge stage for huge categories | "Electronics" has millions of products | Two-stage top-K keeps merge input at workers × K |
| Read spikes on sale days | Everyone opens best-seller pages | CDN and in-process caches; data is tiny and versioned |
| Late or out-of-order events | Mobile retries, partition lag | Event-time buckets with a grace period (watermark) |
Failure handling
- Duplicate events. At-least-once delivery means retries can deliver the same event twice. Deduplicate on
(order_id, line_item_id)with a set that expires after the window (25 hours is enough). - Stream worker crash. Restore counters from the last checkpoint and replay the log from the checkpointed offset. Because checkpoint and offset are saved together, no event is lost or double-counted.
- Ranking job fails. Keep serving the last published version. Alert if
computed_atis older than, say, 2 hours. - Bad data published (a bug ranks a test product first). Versioned rankings allow an instant rollback by moving the pointer to the previous version.
- Category service down. Use the cached product-to-category map; rankings for a newly moved product may lag.
- Fraud and manipulation. Sellers may place fake orders to climb rankings. Exclude orders flagged by fraud systems, cap the contribution of a single buyer per product, and let corrections flow through cancellation events.
Trade-offs and alternatives
| Decision | Chosen | Alternative | Why |
|---|---|---|---|
| Counting | Exact, hourly buckets | Count-min sketch | State fits in ~1 GB; exact avoids rank errors in small categories |
| Pipeline | Streaming plus nightly batch correction | Batch only, or stream only | Freshness of minutes with a safety net |
| Window | 24 hourly buckets | Exponential decay score | Buckets are easy to explain ("last 24 hours"); decay is smoother and needs one number per product |
| Top-K | Local heaps then merge | Sort everything centrally | Less data moved, O(n log k) per worker |
| Serving | Versioned lists in KV plus CDN | Query counters live | Reads are huge and data changes slowly |
| Deduplication | Order line ID set with TTL | Exactly-once processing framework | Simple and framework-independent |
An exponentially decayed score is a common alternative to a hard window: each hour, multiply every product's score by a factor such as 0.9 and add new sales. Recent sales count more, old sales fade smoothly, and you store one number per product instead of 24. The downside is that "score 412.7" is harder to explain than "412 sold in the last 24 hours".
What interviewers probe
"How do you handle a product that goes viral in one hour?" Pre-aggregate (combine counts locally per second before sending), split the hot key across sub-partitions, and rely on the window's hourly granularity; the ranking will show it within one publish cycle.
"Why not just run SELECT ... ORDER BY count DESC LIMIT 100 on the orders table?" It scans and groups millions of rows per category on every request, competing with checkout writes. Precomputing turns a heavy query into a key lookup.
"How would you rank by revenue too?" Keep a second counter (or a pair of counters per bucket) and a second set of heaps. Watch out for currency conversion and refunds.
"How do you test correctness?" Run the nightly batch from raw events and compare its result with the streaming output; alert if the top-K sets differ beyond a threshold.
"What changes for a 7-day window?" Use daily buckets for older days and hourly buckets for the latest day, or 168 hourly buckets. State grows but still fits easily.
Interview questions
Q1. Why is a rolling window harder than an all-time count?
An all-time count only increases, so one counter per product is enough. A rolling window must also subtract sales that leave the window, which means remembering when each sale happened. Hourly buckets solve this by keeping counts per hour and dropping the oldest bucket as time advances.
Q2. How do you find the top 100 of a million counts efficiently?
Use a min-heap of size 100. For each count, compare with the heap minimum and replace it if larger. This runs in O(n log k) time and O(k) memory, which is better than sorting all n counts and works on a stream.
Q3. Why is the two-stage top-K merge exact here?
Each product's full count lives on one worker because events are partitioned by product ID. Any product in the global top-K therefore must be in its own worker's local top-K. If a product's count were split across workers, local top-K lists could miss it and the merge would be approximate.
Q4. What is a count-min sketch and what error does it have?
It is a d × w array of counters with one hash function per row. You increment one cell per row on insert and take the minimum of those cells on query. It never underestimates; with w = ⌈e/ε⌉ and d = ⌈ln(1/δ)⌉, the overestimate is at most εN with probability at least 1 − δ.
Q5. When would you choose a sketch over exact counts?
When the number of distinct keys is huge (billions of search queries or URLs) or memory per worker is tight, and small errors in mid-tail counts are acceptable. For 50 million products, exact counts fit in about a gigabyte, so exact counting is the better choice.
Q6. Batch or streaming for best-seller lists?
It depends on the freshness requirement. Hourly freshness is easily met by batch jobs over hourly aggregates. Minute-level freshness needs streaming. Many systems run streaming for freshness and a nightly batch to recompute and correct drift.
Q7. How do you avoid double counting when events are retried?
Give each line item a stable ID and deduplicate on it, keeping the seen set for slightly longer than the window. Also checkpoint consumer offsets together with state so a restart neither skips nor replays counted events.
Q8. How do cancellations affect rankings?
A cancellation event carries the original order time, and the pipeline subtracts the quantity from that hour's bucket if it is still in the window. Rankings update at the next publish. Cancellations older than the window need no action.
Q9. How do you serve 35,000 ranking reads per second?
The full ranking dataset is about 32 MB and changes every few minutes, so cache it everywhere: in each API process's memory, behind a CDN with a short TTL, and in a key-value store as the source. Reads never touch the counting pipeline.
Q10. How do you prevent sellers from gaming the ranking?
Exclude orders flagged as fraudulent, cap how much one buyer can contribute to a product's count, subtract cancellations and refunds, and monitor for sudden spikes from new accounts. Ranking by units over a window also makes short bursts age out.
Q11. Why bucket by event time instead of processing time?
Processing time depends on pipeline delays and replays, so the same data would produce different windows on different runs. Event time ties each sale to when it really happened, making results reproducible. A watermark (grace period) handles late arrivals.
Q12. How would you publish rankings safely?
Write each new ranking as a new version, then atomically move a pointer to it. Readers always see a complete list, and rollback is just moving the pointer back.
Key takeaways
- A sales ranking is top-K per category over a rolling window of order events; the data is small, but windowing, deduplication and freshness make it interesting.
- Hourly (or finer) buckets implement a sliding window: add new sales, subtract the bucket that ages out.
- A min-heap of size k finds the top k in O(n log k); compute local top-K per worker and merge, which is exact when each product lives on one worker.
- Count-min sketches give bounded-error counts in tiny memory; use them for huge key spaces, not when exact counts fit easily.
- Batch (MapReduce-style) is simple and exact but slower; streaming is fresh but more complex; a nightly batch correcting a streaming path is a common middle ground.
- Bucket by event time, deduplicate by line-item ID, and checkpoint state with offsets to survive retries and crashes.
- Serve versioned, precomputed rankings from caches and a CDN; never compute rankings on the read path.
Next lesson
Continue with Design a Video Streaming Platform.

