CACS352 Distributed System

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).
ClientRequestLoad BalancerRound RobinWeb ServersHandles 1/4 TrafficDatabaseData Storage
How **Daraz** distributes traffic across 4 web servers during Black Friday.

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:

  1. Request hits Cache (CDN) → Hit: Data served in 50ms.
  2. Miss: Query Primary Node (Kathmandu) → If unavailable, Backup Node (Pokhara) responds in 150ms.
  3. 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

Increased resourcesOptimized protocolsLow LatencyHigh ThroughputBalanced
Optimal trade-off: Balanced latency and throughput in distributed systems (e.g., **Daraz** during peak sales).

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).
Time tPacket sent (TCPsequence 100)t+100msACK lost →Retransmissiont+200msACK received →Congestion window halvt+500msCongestion windowrecovers (cubic algori
TCP congestion control in **YouTube** video streaming (real-world packet loss handling).

Worked Example: NTC’s Video Streaming NTC uses adaptive bitrate streaming (ABR) to adjust video quality based on:

  1. Network latency: If latency > 200ms, switch to 720p.
  2. Bandwidth: If bandwidth drops below 5 Mbps, reduce to 480p.
  3. 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

High Throughput (10,000 req/sec) (60%)Low Latency (100ms) (20%)Balanced (7,000 req/sec, 150ms) (20%)
Trade-off between throughput and latency in **eSewa**’s payment system (real-world example).

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:

  1. Global Load Balancing: Routes users to the nearest Google Cloud region.
  2. Borg System: Manages millions of containers across data centers.
  3. Colossus: Handles petabytes of video data (e.g., YouTube).
  4. 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

  1. 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.
  2. 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.
  3. 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…