Price Drop Notification System
Design a data pipeline for an e commerce platform that monitors product prices across millions of SKUs, detects meaningful price drops, notifies subscribed users across multiple channels (email, push, SMS), and measur...
Design a data pipeline for an e-commerce platform that monitors product prices across millions of SKUs, detects meaningful price drops, notifies subscribed users across multiple channels (email, push, SMS), and measures notification effectiveness, click-through rates, conversions, and revenue attribution. The platform has 50 million registered users, 10 million products with watchlists, and prices update roughly 5 times per product per day. How would you design this end to end?
How to Approach This Problem
What Makes This Problem Unique At first glance this sounds simple: detect a price drop, send a notification. The trap is that candidates design for the happy path and miss the three hard problems that actually define this system. Hard problem 1: Fan out asymmetry A single price event on a popular product triggers millions of notifications. Your architecture must handle a 1,000,000:1 amplification ratio from one Kafka message to one million deliveries without creating a thundering herd on downstream email and push services. Hard problem 2: Price oscillation and deduplication Prices bounce. A product can drop, partially recover, and drop again multiple times a day. Without a dedup strategy, us
Clarifying Questions to Ask the Interviewer
Functional Requirements What constitutes a "price drop"? Any decrease? Percentage based ( 5%)? Absolute amount ( $10)? User set threshold ("notify me when below $50")? How do we receive price data? Real time seller API push? Periodic scraping? Marketplace feed (like Amazon MWS)? What notification channels are supported? Email, push notification, SMS, in app? All or a subset? Do users explicitly subscribe to products (watchlist), or do we also infer interest from browsing/purchase history? Should we track notification effectiveness? (open rates, click through rates, purchases attributed to notifications) Do we support "price target" alerts? (notify me when Product X drops below $50) Should we
Envelope Estimation & Capacity Planning
Throughput Math Metric Value Calculation Products in catalog 10,000,000 Given Price updates per product per day 5 Given (seller feeds + scrapes) Total price events per day 50,000,000 10M × 5 Price events per second (avg) ~580 50M ÷ 86,400 Price events per second (peak) ~5,800 10× peak (Black Friday, flash sales) Average price event size ~500 bytes product id, seller id, price, currency, timestamp Daily raw ingestion ~25 GB/day 50M × 500 bytes Fan Out Math (The Critical Bottleneck) Metric Value Calculation Users with active watchlists 10,000,000 Given Avg products per watchlist 20 Assumption Total watchlist entries 200,000,000 10M × 20 Avg subscribers per product 20 200M ÷ 10M products Price
Architecture Walkthrough, End-to-End Data Flow
Why This Isn't a Simple CRUD Problem At first glance, "detect price drops and send notifications" sounds easy, poll prices, compare, send email. But at 10M products × 10M watchlist users, the naive approach (for each product, query all subscribers) does a full scan of 200M watchlist entries 50M times per day. The real challenge is the fan out : a single event (iPhone price drops) must efficiently reach 1M subscribers within 5 minutes without overloading any downstream service. The architecture separates three concerns into independent, decoupled stages: 1. Detection (streaming), What changed? 2. Matching (lookup + fan out), Who cares? 3. Delivery (multi channel dispatch), How to tell them? T
Component Deep Dive
Price Change Detector, Flink Keyed State "The detection layer uses Flink's keyed state to hold the last known price per product in memory, this gives us O(1) comparisons without a database round trip on every event." Why Flink keyed state? We receive 580 price events/sec (5,800 at peak). If we looked up the last price from a database for each event, that's 580 DB queries/sec, feasible but adds latency and creates a DB bottleneck during spikes. Flink's keyed state stores the last price per product id in local RocksDB, co located with the processing, zero network overhead. State schema per product: Detection pseudocode: State management: State backend: RocksDB (disk backed, handles 10M product
Data Modeling & Schema Design
Core Tables Fact Table: Price History (Data Lake, Source of Truth) Partitioning rationale: DATE(event timestamp) gives ~50M rows per partition (one day). Queries almost always filter by date range. Clustering by product id speeds up per product lookups for backfill and analysis. Dimension Table: User Watchlist Fact Table: Notification Log Fact Table: Notification Effectiveness Gold Table: Daily Notification Metrics dbt Model: Effectiveness Join Data Lake Layout File format: Parquet (columnar, compressed, predicate pushdown) Table format: Delta Lake or Apache Iceberg (ACID, time travel, schema evolution, efficient MERGE for dedup)