The Core Insight: Separate Read Paths from Write Paths
When you press Play, here's what actually happens:
- Your request hits a CDN edge node (99% of the time, the video is already cached near you)
- Your watch history and preferences are fetched from a read replica — not the primary DB
- Playback events (pauses, seeks, buffer stats) are queued asynchronously, not written immediately
- Only after you stop watching does a batch job sync your actual watch state
The Secret Ingredient: Apache Cassandra
Netflix uses Cassandra for user activity data. Here's why:
- Designed for write-heavy, distributed workloads
- Uses eventual consistency your "continue watching" list might be 2 seconds stale, but nobody cares
- Scales horizontally across data centers without complex sharding logic
The Tradeoff They Made
Strong consistency was traded for performance.
Your "recently watched" might not update instantly. But 270 million people can stream simultaneously without a database bottleneck.
What This Means for You
Not every data write needs to be synchronous. Not every read needs to come from primary.
Understanding where you can afford eventual consistency is one of the most underrated distributed systems skills.
What's a system design tradeoff you've had to make at work? Let's discuss in the comments 👇

Top comments (0)