System Design: Counting Things at Large Scale
Key Concepts: Requirements Clarification (Users, Scale, Performance, Cost), Functional & Non-Functional Requirements, Data Model (Raw Events vs. Aggregated Data), SQL vs. NoSQL Databases (Cassandra), Data Partitioning, Data Replication, Consistency (Eventual, Tunable), Data Normalization vs. Denormalization, Data Processing (Push vs. Pull, Checkpointing), Stream Processing, Batch Processing, Data Aggregation, Message Queues, Dead Letter Queue, Data Enrichment, State Management, Service Discovery (Client-Side, Server-Side), Load Balancing, Message Formats (Textual vs. Binary), Caching, Rollup, Hot Storage vs. Cold Storage, Performance Testing, Health Monitoring, Audit Systems, Hot Partitions, Technology Stack.
1. Requirements Clarification
- Main Topic: Importance of asking clarifying questions to the interviewer to define the scope and requirements of the system.
- Key Points:
- Interviewers assess how candidates handle ambiguity.
- Clarification helps determine the appropriate technologies and building blocks.
- Focus on four categories: Users, Scale, Performance, and Cost.
- Specific Questions:
- Users: Who will use the system? How will they use the data? (e.g., Youtube viewers, video owners, machine learning models, marketing department).
- Scale: How many read queries per second? How much data is queried per request? How many video views per second? Should we deal with traffic spikes?
- Performance: How fast must the system be? (e.g., real-time vs. batch processing). How fast data must be retrieved from the system?
- Cost: Minimize development cost? Minimize future maintenance cost?
- Advice: Think along the four categories, focus on data flow, and don't worry too much about time during clarification.
- Functional Requirements: System behavior, specifically APIs (e.g.,
countViewEvents(video)). Generalize APIs by introducing parameters likeeventTypeand supporting functions likesumandaverage. - Non-Functional Requirements: System qualities (e.g., fast, fault-tolerant, secure). Prioritize scalability, performance, and availability. Consider consistency (CAP theorem) and cost minimization.
2. High-Level Architecture
- Main Topic: Initial architectural design and the importance of defining the data model early.
- Key Points:
- Start with a simple architecture: Database, Processing Service, Query Service.
- Focus on data model: What data to store and how.
- Data Model Options:
- Store Raw Events: Capture all attributes (video ID, timestamp, user info).
- Pros: Fast writes, flexible queries, ability to recalculate.
- Cons: Slow reads, high storage costs.
- Aggregate Data: Calculate views on the fly and store aggregated data (e.g., total count per minute).
- Pros: Fast reads, real-time decision making.
- Cons: Limited query flexibility, complex aggregation pipeline, difficult to fix errors.
- Store Raw Events: Capture all attributes (video ID, timestamp, user info).
- Decision Making: Ask the interviewer about expected data delay to choose between stream (real-time) and batch processing.
- Combined Approach: Store raw events for a limited time and aggregate data in real-time for immediate statistics.
3. Database Selection (SQL vs. NoSQL)
- Main Topic: Evaluating SQL and NoSQL databases based on non-functional requirements.
- Key Points:
- Both SQL and NoSQL can scale and perform well.
- Evaluate against scalability, performance, availability, consistency, data recovery, security, data model changes, hardware, and cost.
- SQL Database Scaling (Example: MySQL with Vitess):
- Sharding (horizontal partitioning) to split data across multiple machines.
- Cluster Proxy to route traffic to the correct shard.
- Configuration Service to monitor shard health.
- Shard Proxy for caching, health monitoring, and query termination.
- Data Replication (Master/Read Replica) for availability and data center redundancy.
- NoSQL Database Scaling (Example: Apache Cassandra):
- Each shard (node) is equal.
- Gossip Protocol for nodes to exchange state information.
- Coordinator Node: Any node can act as a coordinator to route requests.
- Consistent Hashing to pick the node for data storage.
- Quorum Writes/Reads for data consistency.
- Tunable Consistency: Choose availability over consistency (eventual consistency).
4. Data Modeling (SQL vs. NoSQL)
- Main Topic: Contrasting data modeling approaches for SQL and NoSQL databases.
- SQL (Relational Databases):
- Start with defining nouns (entities) in the system.
- Convert nouns into tables and use foreign keys for relationships.
- Data Normalization: Minimize data duplication.
- NoSQL (Cassandra):
- Think in terms of queries, not nouns.
- Data Denormalization: Store everything required for a report together.
- Example: Store video info, hourly views, and channel info in a single table for a report.
5. Data Processing
- Main Topic: Designing a scalable, reliable, and fast data processing service.
- Key Points:
- Partitioning for scalability.
- Replication for data loss prevention.
- In-memory processing for speed.
- Pre-Aggregation: Accumulate data in the processing service memory before updating counters in the database.
- Push vs. Pull: Pull events from a temporary storage (queue) for better fault tolerance and scalability.
- Checkpointing: Store the current offset in a persistent storage to resume processing after a failure.
- Processing Service Components:
- Consumer: Reads events from the partition, deserializes them, and eliminates duplicates using a distributed cache.
- Aggregator: Accumulates data in memory (hash table) for a period of time.
- Internal Queue: Decouples consumption and processing, allowing multi-threaded processing.
- Database Writer: Stores pre-aggregated views count in the database.
- Dead Letter Queue: Stores undelivered messages due to database issues.
- Data Enrichment: Retrieves additional attributes (video title, channel name) from an embedded database (e.g., RocksDB).
- State Management: Periodically save the entire in-memory data to a durable storage for recovery.
6. Data Ingestion Pipeline
- Main Topic: End-to-end data flow from event generation to storage.
- Components:
- API Gateway: Single entry point for client requests.
- Partitioner Service Client: Batches events and sends them to the partitioner service.
- Load Balancer: Distributes events across partitioner service machines.
- Partitioner Service: Routes events to partitions based on a partition strategy.
- Partitions: Store events on disk in the form of append-only log files.
- Key Concepts:
- Blocking vs. Non-Blocking I/O: Non-blocking I/O for higher throughput.
- Buffering and Batching: Combine events to reduce overhead.
- Timeouts: Connection timeout and request timeout.
- Retries: Exponential backoff and jitter to prevent retry storms.
- Circuit Breaker: Stop repeated attempts to failing operations.
- Load Balancing: Hardware vs. Software, TCP vs. HTTP, Round Robin, Least Connections, Least Response Time, Hash-Based.
- Service Discovery: Server-Side (Load Balancer), Client-Side (Zookeeper).
- Replication: Single Leader, Multi Leader, Leaderless.
- Message Formats: Textual (JSON) vs. Binary (Thrift, Protocol Buffers, Avro).
- Hot Partitions: Spread events for popular videos across several partitions.
7. Data Retrieval Path
- Main Topic: Retrieving and serving video statistics to users.
- Components:
- API Gateway: Routes requests to the Query Service.
- Query Service: Retrieves statistics from the database and object storage.
- Database: Stores recent statistics.
- Object Storage (e.g., AWS S3): Stores older, rolled-up statistics.
- Distributed Cache: Stores query results for faster access.
- Rollup: Aggregate data over time to reduce storage costs (e.g., per-minute -> per-hour -> per-day).
- Hot Storage vs. Cold Storage: Store frequently accessed data in a fast database and less frequently accessed data in object storage.
- Data Federation: Query service retrieves data from multiple storages and stitches it together.
8. Technology Stack
- Main Topic: Discussing relevant technologies for building the system.
- Examples:
- Networking: Netty (non-blocking IO framework).
- Client-Side Resilience: Hystrix, Polly (timeouts, retries, circuit breaker).
- Load Balancing: Citrix Netscaler (hardware), NGINX (software), AWS ELB (cloud).
- Message Queue: Apache Kafka, Amazon Kinesis.
- Stream Processing: Apache Spark, Apache Flink, Kinesis Data Analytics.
- Wide Column Databases: Apache Cassandra, Apache HBase.
- Time Series Databases: InfluxDB.
- Data Warehousing: Apache Hadoop, AWS Redshift.
- Object Storage: AWS S3.
- MySQL Scaling: Vitess.
- Distributed Cache: Redis.
- Message Broker: RabbitMQ, Amazon SQS.
- Embedded Database: RocksDB.
- Service Discovery: Apache Zookeeper, Netflix Eureka.
- Monitoring: AWS CloudWatch, Elasticsearch/Logstash/Kibana (ELK).
- Hashing: MurmurHash.
9. Bottlenecks and Tradeoffs
- Main Topic: Identifying and addressing potential bottlenecks in the system.
- Performance Testing: Load testing, stress testing, soak testing.
- Health Monitoring: Metrics, dashboards, alerts (latency, traffic, errors, saturation).
- Audit Systems: Weak (continuous end-to-end test) vs. Strong (parallel system using a different path, e.g., Lambda Architecture).
- Hot Partitions: Spread events for popular videos across several partitions.
- Overload Handling: Batch events and store them in object storage, then process them asynchronously using a cluster of machines.
10. Synthesis/Conclusion
- Main Takeaways:
- Knowledge of system design concepts is crucial for success.
- Requirements clarification is essential for defining the scope and choosing the right technologies.
- Data modeling is a key decision that impacts performance and flexibility.
- Scalability, performance, and availability are primary non-functional requirements.
- There are tradeoffs in every design decision.
- The same ideas can be applied to various counting and monitoring problems.
- Continuous learning and exploration of new technologies are essential.
AI summaries can miss context or contain errors. Check important details against the original video.