Distributed SystemUnit 1312 min read
Scalability & Performance: Load Balancing, Caching, Replication, and Bottlenecks
Unit 13 of Distributed System: Explores how distributed systems scale horizontally/vertically, optimize performance via load balancing, caching, replication, and network efficiency, while identifying bottlenecks and trade-offs in latency, throughput, and resource utilization.
TAKEAWAYS:
- Scalability is achieved via horizontal scaling (adding nodes) or vertical scaling (upgrading hardware), with load balancing distributing requests to avoid overload.
- Caching (local, distributed, CDNs) reduces latency by storing frequently accessed data closer to users, but requires cache coherence protocols like write-through or write-back.
- Replication improves availability and performance but introduces consistency trade-offs (CAP theorem) and stale data risks; quorum-based and eventual consistency models are common.
- Network bottlenecks (bandwidth, latency, packet loss) are mitigated by asynchronous communication, pipelining, and compression, while congestion control (e.g., TCP’s ACK/NACK) prevents network collapse.
- Performance metrics (throughput, response time, utilization) are optimized via benchmarking tools (e.g., JMeter, Apache Bench) and autoscaling (e.g., Kubernetes HPA).
- Real-world trade-offs: Daraz’s order processing uses replication + caching for fast checkout but risks inconsistent inventory; NTC’s CDN caches popular videos to reduce server load but may show stale content.
1. Scalability in Distributed Systems
Scalability refers to a system’s ability to handle increased load by adding resources (nodes, bandwidth, or processing power) without degrading performance. Distributed systems scale via:
1.1 Horizontal vs. Vertical Scaling
| Type | Definition | Example | Pros | Cons |
|---|---|---|---|---|
| Vertical | Upgrading a single node (e.g., more CPU/RAM). | Upgrading a Ncell server to handle more mobile users. | Simple, no reconfiguration. | Expensive, single point of failure. |
| Horizontal | Adding more nodes (e.g., servers, databases). | Daraz adding more web servers during festival sales. | Cost-effective, fault-tolerant. | Complex coordination (e.g., load balancing). |
Visual: Horizontal vs. Vertical Scaling
graph TD
A["Single Server"] -->|"Vertical"| B[Upgraded Server: More CPU/RAM
(Example: 16-core → 32-core, 32GB RAM → 64GB RAM)]
A -->|"Horizontal"| C["Cluster of Servers: 3 Nodes"]
C --> D["Load Balancer: Round Robin"]
D --> E["Server 1\nHandles 1/3 Traffic\n(1000 req/s)"]
D --> F["Server 2\nHandles 1/3 Traffic\n(1000 req/s)"]
D --> G["Server 3\nHandles 1/3 Traffic\n(1000 req/s)"]
B --> H["Single Point of Failure\n(No redundancy)"]
C --> I["Fault-Tolerant\n(If one node fails, others handle load)"]1.2 Load Balancing
Distributes incoming requests across multiple servers to prevent overload and maximize throughput. Methods:
- Static Load Balancing: Predefined distribution (e.g., round-robin).
- Dynamic Load Balancing: Monitors server load (e.g., NTC’s CDN routes users to the least busy server).
- Global Server Load Balancing (GSLB): Routes users to the nearest geographic server (e.g., Google’s Anycast DNS).
Worked Example: Pathao’s Ride Allocation Pathao uses dynamic load balancing to assign drivers to zones based on real-time demand. If Zone A has high demand, the system routes more requests to drivers in Zone A while reducing requests to Zone B to avoid driver burnout.
2. Performance Optimization Techniques
2.1 Caching
Stores frequently accessed data closer to users to reduce latency. Types:
- Local Caching: Client-side (e.g., browser cache for YouTube videos).
- Distributed Caching: Shared cache (e.g., Redis in Khalti’s payment gateway).
- Content Delivery Networks (CDNs): Edge caching (e.g., Cloudflare for eSewa’s website).
Cache Coherence Protocols:
- Write-Through: Update cache and disk simultaneously (strong consistency, higher latency).
- Write-Back: Update cache first, then disk (faster but risks stale data).
Visual: Cache Hit/Miss
graph TD
A["User Request"] --> B{"Cache Hit?"}
B -->|"Yes"| C["Data from Cache<br/>(0 ms latency)"]
B -->|"No"| D["Fetch from Disk/DB<br/>(100 ms latency)"]
D --> E["Update Cache"]2.2 Replication
Copies data across multiple nodes to:
- Improve read performance (e.g., NEPSE’s stock data replication).
- Ensure high availability (e.g., Ncell’s backup databases).
Replication Strategies:
| Model | Description | Example | Consistency |
|---|---|---|---|
| Primary-Backup | One primary node handles writes; backups replicate. | Bank transaction logs. | Strong (eventual). |
| Multi-Master | All nodes can write (conflict resolution needed). | Wikipedia edits. | Weak (eventual). |
| Quorum-Based | Writes require majority approval (e.g., DynamoDB’s Paxos). | Khalti’s payment processing. | Strong. |
Worked Example: Daraz’s Inventory Replication Daraz replicates inventory data across 3 data centers (Kathmandu, Pokhara, Dhaka). When a user checks stock for a product:
- Request hits Cache (CDN) → Hit: Data served in 50ms.
- Miss: Query Primary Node (Kathmandu) → If unavailable, Backup Node (Pokhara) responds in 150ms.
- Conflict: If two users buy the same item simultaneously, quorum voting ensures only one transaction succeeds.
3. Network Bottlenecks and Mitigations
3.1 Common Bottlenecks
| Bottleneck | Cause | Impact | Mitigation |
|---|---|---|---|
| Bandwidth | Limited network capacity. | Slow data transfer (e.g., NTC’s slow rural upload). | Compression (e.g., gzip), CDNs. |
| Latency | Distance between nodes. | Delayed responses (e.g., WhatsApp calls to Europe). | Edge computing, Anycast DNS. |
| Packet Loss | Network congestion. | Retransmissions (e.g., YouTube buffering). | TCP retransmission, UDP for real-time. |
| CPU/Memory | Overloaded servers. | Slow processing (e.g., Ncell’s call drops). | Load balancing, auto-scaling. |
Visual: Network Latency vs. Throughput
3.2 Congestion Control
Prevents network collapse by:
- TCP’s ACK/NACK: Slows down if packet loss > threshold.
- Asynchronous Communication: Decouples senders/receivers (e.g., Kafka for eSewa’s transaction logs).
Worked Example: NTC’s Video Streaming NTC uses adaptive bitrate streaming (ABR) to adjust video quality based on:
- Network latency: If latency > 200ms, switch to 720p.
- Bandwidth: If bandwidth drops below 5 Mbps, reduce to 480p.
- Buffer: If buffer < 5s, pause playback to avoid stalling.
4. Performance Metrics and Benchmarking
| Metric | Definition | How to Measure | Example |
|---|---|---|---|
| Throughput | Data processed per unit time (e.g., transactions/sec). | Benchmark tools (e.g., JMeter). | Khalti: 10,000 transactions/sec. |
| Response Time | Time from request to response. | Ping, latency tests. | eSewa: 300ms avg. |
| Utilization | % of resources (CPU, memory) in use. | top command, Prometheus. |
Ncell’s server: 70% CPU usage. |
| Availability | % uptime (e.g., 99.9% SLA). | Uptime monitoring (e.g., Nagios). | NEPSE: 99.95% uptime. |
Visual: Performance Metrics Trade-off
5. Autoscaling and Elasticity
Automatically adjusts resources based on demand:
- Vertical Autoscaling: Upgrades a single node (e.g., Ncell’s peak-hour server upgrade).
- Horizontal Autoscaling: Adds/removes nodes (e.g., Daraz’s festival-season scaling).
Example: Kubernetes Horizontal Pod Autoscaler (HPA)
- Trigger: CPU usage > 70% for 5 minutes.
- Action: Spins up 3 new pods to handle load.
- Result: Response time drops from 800ms → 200ms.
6. Case Study: Google’s Distributed Scalability
Google uses:
- Global Load Balancing: Routes users to the nearest Google Cloud region.
- Borg System: Manages millions of containers across data centers.
- Colossus: Handles petabytes of video data (e.g., YouTube).
- Spanner: Globally distributed database with strong consistency.
Why It Works:
- Horizontal scaling: Adds 1000s of machines during peak traffic.
- Caching: Memcached reduces database load by 90%.
- Replication: 3 copies of every piece of data across regions.
In the Real World
Daraz’s Order Processing
- Idea: Replication + Caching for fast checkout.
- How: Orders are cached in Redis (50ms response), but inventory is eventually consistent (risk: overselling).
- Trade-off: Users get fast checkout, but rare cases of "out of stock" errors occur.
NTC’s Video Streaming (YouTube Nepal)
- Idea: CDN Caching + Adaptive Bitrate.
- How: Videos are cached in NTC’s edge servers (e.g., in Pokhara), but quality drops if bandwidth < 3 Mbps.
- Real Impact: Rural users see 480p videos instead of 1080p.
Ncell’s 5G Network
- Idea: Load Balancing + Edge Computing.
- How: 5G base stations dynamically allocate spectrum to users, reducing congestion in busy areas (e.g., Thamel).
- Result: Call drops drop from 12% → 2% during peak hours.
Exam Tip
- Focus on trade-offs: Always discuss pros/cons of techniques (e.g., "Replication improves availability but risks stale data").
- Draw diagrams: Include load balancer flows, cache hit/miss, and replication models (primary-backup, multi-master).
- Real-world links: Connect to Nepali examples (e.g., "How does Khalti use caching to handle 100,000 transactions/sec?").
- Math in benchmarks: Know how to calculate throughput (e.g., "If 1000 requests take 2 sec, throughput = 500 req/sec").
- CAP Theorem: Always mention it when discussing replication consistency (e.g., "DynamoDB uses AP for scalability").
Common Pitfalls:
- ❌ Saying "Replication always ensures strong consistency" → Correction: It depends on the model (e.g., Paxos vs. eventual consistency).
- ❌ Ignoring network bottlenecks → Always mention latency/bandwidth as a limiting factor.
- ❌ Not comparing horizontal vs. vertical scaling → Highlight cost vs. flexibility.
Based on the TU BCA syllabus for Distributed System (CACS352), unit 13.
Discussion
Loading…