Distributed SystemUnit 1220 min read
Networking & Protocols: OSI, TCP/IP, RPC, and Distributed Messaging
Unit 12 of Distributed System: Explores how distributed systems communicate via protocols (OSI, TCP/IP), message passing (RPC, message queues), and synchronization tools (time protocols, clocks), with real-world examples from eSewa, WhatsApp, and NTC.
TAKEAWAYS:
- Distributed systems rely on standardized protocols (OSI, TCP/IP) to ensure reliable communication across heterogeneous networks.
- RPC (Remote Procedure Call) abstracts distributed calls like local function calls, while message-oriented middleware decouples senders and receivers (e.g., Kafka, RabbitMQ).
- Clock synchronization (NTP, Lamport timestamps) is critical for consistency in distributed transactions (e.g., bank transfers, Daraz order processing).
- Network topologies (bus, star, mesh) and routing protocols (OSPF, BGP) determine how data travels in real-world networks like NTC’s fiber backbone.
- Failure handling (timeouts, retries) and recovery mechanisms (checkpointing) mirror how eSewa handles payment retries or Pathao recovers from driver unavailability.
- Security protocols (TLS, SSH) protect data in transit, just as Ncell encrypts user calls over its 5G network.
1. Introduction to Networking in Distributed Systems
Distributed systems communicate over networks, which require protocols—rules governing data exchange. These protocols define:
- Syntax: Format of messages (e.g., IP addresses, port numbers).
- Semantics: Meaning of messages (e.g., "ACK," "NACK").
- Timing: When to send/receive (e.g., handshakes, timeouts).
1.1 Core Concepts: Name, Address, Identifier
| Term | Definition | Example in Distributed Systems |
|---|---|---|
| Name | Human-readable identifier (e.g., user123, server-A). |
eSewa account name: "john.doe@esewa.com" |
| Address | Network-locator (e.g., IP 192.168.1.1, MAC 00:1A:2B:3C:4D:5E). |
NTC router’s IP: 203.101.152.1 |
| Identifier | Unique handle for processes (e.g., PID, UUID). | WhatsApp message ID: abc123-xyz456 |
Why it matters:
- A name helps users (or apps) reference resources (e.g.,
Daraz.com). - An address routes packets to the correct machine (e.g., Ncell’s cell tower).
- An identifier ensures uniqueness for processes (e.g., a Pathao driver’s session token).
1.2 Types of Communication in Distributed Systems
Distributed systems use four primary communication models, each with trade-offs:
graph TD
A["Communication Models"] --> B["Synchronous"]
A --> C["Asynchronous"]
A --> D["Transient"]
A --> E["Persistent"]
B --> B1["Request-Reply (RPC)"]
B --> B2["Streaming (WebSockets)"]
C --> C1["Event-Driven (Kafka)"]
C --> C2["Message Queues (RabbitMQ)"]
D --> D1["Temporary (e.g., chat messages)"]
E --> E1["Permanent (e.g., database logs)"]| Model | Description | Example Use Case | Advantages | Disadvantages |
|---|---|---|---|---|
| Synchronous | Sender waits for response (e.g., HTTP GET). | eSewa transaction confirmation. | Simple, real-time feedback. | Blocking, high latency risk. |
| Asynchronous | Sender doesn’t wait (e.g., email, Kafka). | Pathao ride request (driver accepts later). | Scalable, decoupled. | Complex error handling. |
| Transient | Messages discarded if not delivered immediately. | WhatsApp message (deleted after 7 days). | Low storage overhead. | No replayability. |
| Persistent | Messages stored until processed. | Bank transaction logs. | Reliable, auditable. | Higher storage cost. |
Worked Example: eSewa’s Synchronous Payment
- User clicks "Pay" → eSewa sends a synchronous request to the bank.
- Bank processes payment and sends an ACK (success) or NACK (failure).
- eSewa updates the UI immediately (real-time feedback).
Trace:
User → [eSewa] → [Bank: Synchronous Request] → [Bank: ACK/NACK] → [eSewa: Update UI]
2. Networking Models: OSI vs. TCP/IP
Distributed systems use layered models to organize communication. The two dominant models are:
2.1 OSI Model (7 Layers)
Layers (Top to Bottom):
- Application (HTTP, FTP, RPC)
- Presentation (Encryption: TLS, Compression)
- Session (NetBIOS, RPC)
- Transport (TCP, UDP)
- Network (IP, Routing)
- Data Link (Ethernet, Wi-Fi)
- Physical (Cables, Signals)
Why OSI?
- Standardized for interoperability (e.g., a Nepali bank’s system talking to a global payment gateway).
- Helps troubleshoot by isolating issues (e.g., "Is it a transport layer problem (TCP) or a network layer problem (IP)?").
2.2 TCP/IP Model (4 Layers)
Layers (Top to Bottom):
- Application (HTTP, DNS, SMTP)
- Transport (TCP, UDP)
- Internet (IP, ICMP)
- Network Access (Ethernet, PPP)
Key Differences:
| Feature | OSI Model | TCP/IP Model |
|---|---|---|
| Layers | 7 (detailed) | 4 (simplified) |
| Use Case | Academic/theoretical | Practical (used by all modern OS) |
| Session | Explicit layer | Merged into Application/Transport |
Real-World Example: YouTube Video Stream
- Application Layer: HTTP requests video chunks.
- Transport Layer: TCP ensures reliable delivery (retries lost packets).
- Internet Layer: IP routes packets via NTC’s fiber network.
- Network Access: Ethernet/Wi-Fi delivers to your phone.
3. Communication Protocols
Protocols define how data is exchanged. Key protocols in distributed systems:
3.1 Remote Procedure Call (RPC)
RPC lets a program call a procedure on another machine as if it were local.
How RPC Works:
- Client encodes arguments into a message.
- Stub (client-side) sends the message over the network.
- Server decodes the message, executes the procedure.
- Stub (server-side) encodes the return value.
- Client receives and decodes the response.
sequenceDiagram
participant Client
participant StubClient
participant Network
participant StubServer
participant Server
Client->>StubClient: Call getWeather("Kathmandu")
StubClient->>Network: RPC Message (encoded args)
Network->>StubServer: Deliver message
StubServer->>Server: Decode & execute
Server-->>StubServer: Return 25°C
StubServer->>Network: Encode response
Network->>StubClient: Send back
StubClient-->>Client: Return 25°CExample: Daraz’s Inventory Check
- Client: Your phone calls
checkStock("Laptop"). - Server: Daraz’s warehouse RPC server checks inventory.
- Response:
{"status": "in_stock", "quantity": 10}.
Advantages:
- Transparency: Developers write code as if functions were local.
- Scalability: Load balances across servers (e.g., Ncell’s call centers).
Disadvantages:
- Complexity: Network latency, serialization overhead.
- Fault Tolerance: Requires timeouts/retry logic (e.g., if Daraz’s server crashes).
3.2 Message-Oriented Communication
Unlike RPC (synchronous), message queues (e.g., RabbitMQ, Kafka) decouple senders and receivers.
Key Concepts:
- Producer: Sends messages (e.g., Pathao app sends ride requests).
- Queue: Buffers messages (e.g., Pathao’s driver assignment queue).
- Consumer: Processes messages (e.g., a driver’s phone).
sequenceDiagram
participant PathaoApp
participant RabbitMQ
participant DriverPhone
PathaoApp->>RabbitMQ: Publish("New Ride Request: Kathmandu")
RabbitMQ->>DriverPhone: Deliver to available drivers
DriverPhone-->>RabbitMQ: ACK (accepted) or NACK (rejected)Example: WhatsApp Messages
- Producer: You send a message to a group.
- Queue: WhatsApp’s servers buffer messages until recipients are online.
- Consumer: Recipients’ phones receive messages asynchronously.
Advantages:
- Decoupling: Producers/consumers don’t need to know each other.
- Scalability: Handles spikes (e.g., Ncell during New Year calls).
Disadvantages:
- Ordering: Messages may arrive out of order (unless prioritized).
- Storage: Queues consume disk space (e.g., WhatsApp’s message history).
3.3 Time and Clock Synchronization
Distributed systems need consistent time for:
- Transaction ordering (e.g., bank transfers).
- Debugging (e.g., "When did error X occur?").
Clock Types:
| Type | Description | Example Use Case |
|---|---|---|
| Physical | Real-world time (e.g., atomic clocks). | NTP servers (e.g., pool.ntp.org). |
| Logical | Virtual time (e.g., Lamport timestamps). | Distributed logs (e.g., Kafka). |
3.3.1 Network Time Protocol (NTP)
NTP synchronizes clocks across machines using time servers.
How NTP Works:
- Client requests time from an NTP server.
- Server responds with its time + delay/offset.
- Client adjusts its clock.
sequenceDiagram
participant Client
participant NTPServer
Client->>NTPServer: Request time
NTPServer-->>Client: Time + delay/offset
Client->>NTPServer: Sync request (if needed)
NTPServer-->>Client: Final adjustmentExample: NEPSE Stock Market
- All brokers’ systems must sync to Nepal Standard Time (NST) to avoid:
- Double-spending: Two trades processed at the same "time."
- Audit issues: "When did this trade occur?"
3.3.2 Lamport Timestamps
For logical clocks, Lamport timestamps ensure causal ordering of events.
Rules:
- Each event gets a timestamp
T. - If
A → B(A causes B), thenT(A) < T(B). - If
AandBare concurrent,T(A)andT(B)can be equal or ordered arbitrarily.
Example: Bank Transaction Log
Event A: User clicks "Transfer Rs. 1000"
Event B: Bank debits Rs. 1000 from account X
Event C: Bank credits Rs. 1000 to account Y
T(A) = 1,T(B) = 2,T(C) = 3(causal order).- Two parallel transactions can have
T = 3andT = 4.
4. Network Topologies and Routing
Distributed systems use topologies to connect nodes. Common types:
| Topology | Description | Example in Distributed Systems |
|---|---|---|
| Bus | All nodes share a single communication line. | Ethernet cables in a college lab. |
| Star | Central node (hub) connects to all others. | NTC’s cell towers connected to a core router. |
| Mesh | Every node connects to every other node. | Bitcoin’s peer-to-peer network. |
| Ring | Nodes form a closed loop. | Token-passing in old LANs. |
4.1 Routing Protocols
Routing protocols determine how data travels from source to destination.
| Protocol | Type | Description | Example Use Case |
|---|---|---|---|
| OSPF | Interior | Link-state routing (fast convergence). | NTC’s internal network. |
| BGP | Exterior | Path-vector routing (internet-scale). | Global internet (e.g., Google → Ncell) |
| RIP | Interior | Distance-vector (slow convergence). | Small office networks. |
Example: Google Search Request
- Your phone sends a DNS request to
8.8.8.8(Google’s DNS). - DNS resolves
google.comto an IP (e.g.,142.250.190.46). - BGP routes the request through ISPs (Ncell → NTC → Global backbone).
- HTTP delivers the search results.
5. Security in Networking
Distributed systems must protect data in transit and at rest.
5.1 TLS (Transport Layer Security)
Encrypts data between client and server (e.g., HTTPS).
How TLS Works:
- Client and server agree on a cipher suite (e.g., AES-256).
- Server sends its public key (in a certificate).
- Client generates a pre-master secret, encrypts it with the server’s public key.
- Both sides derive a symmetric key for encryption.
sequenceDiagram
participant Client
participant Server
Client->>Server: TLS Handshake (Hello)
Server-->>Client: Certificate + Public Key
Client->>Server: Encrypted Pre-Master Secret
Client->>Server: Finished (symmetric key derived)
Server-->>Client: FinishedExample: eSewa’s Secure Payments
- When you pay via eSewa, TLS 1.3 encrypts your card details before sending them to the bank.
5.2 SSH (Secure Shell)
Encrypts remote login and command execution (e.g., accessing a server).
Example: Ncell’s Network Monitoring
- Engineers use SSH to log into Ncell’s routers to check traffic:
ssh admin@router-ktm "show interface status"
6. Scalability and Performance
Distributed systems must handle growth without performance degradation.
6.1 Load Balancing
Distributes traffic across multiple servers (e.g., Ncell’s call centers).
Example: Daraz During Festive Sales
- Round-robin: Requests alternate between servers.
- Least connections: Directs traffic to the least busy server.
graph TD
A["Load Balancer"] --> B["Server 1: CPU 80%"]
A --> C["Server 2: CPU 40%"]
A --> D["Server 3: CPU 20%"]
User --> A: HTTP Request
A -->|"Round-robin"| B
A -->|"Least connections"| C6.2 Caching
Stores frequent data closer to users (e.g., CDNs like Cloudflare).
Example: YouTube Video Caching
- You request a video from
youtube.com. - DNS resolves to a nearby CDN (e.g.,
cdn.youtube.netin Kathmandu). - CDN serves the video from cache (if available) or fetches from origin.
7. Fault Tolerance and Recovery
Distributed systems must recover from failures.
7.1 Types of Failures
| Failure Type | Description | Example |
|---|---|---|
| Node Failure | A machine crashes (e.g., server reboot). | Ncell’s cell tower goes offline. |
| Network Failure | Link breaks (e.g., fiber cut). | NTC’s Kathmandu-Pokhara link down. |
| Software Failure | Bug or crash in an application (e.g., Pathao app freezes). |
7.2 Recovery Mechanisms
| Mechanism | Description | Example |
|---|---|---|
| Checkpointing | Save state periodically. | Bank transaction logs. |
| Replication | Run multiple copies of a service (e.g., Ncell’s call center redundancy). | |
| Timeouts | Abort stalled requests (e.g., eSewa’s 30-second payment timeout). |
7.3 Commit Protocols
Ensure all nodes agree on a transaction (e.g., bank transfer).
Two-Phase Commit (2PC):
- Prepare Phase: Coordinator asks all participants if they can commit.
- Commit Phase: If all say "Yes," coordinator sends "Commit."
Disadvantages:
- Blocking: Coordinator waits for all participants.
- Single Point of Failure: Coordinator crash causes deadlock.
One-Phase Commit (1PC):
- Simpler but no rollback if a participant fails.
In the Real World
eSewa’s RPC for Payments
- Idea: Uses RPC to call the bank’s
processPayment()method. - How: Your phone sends an RPC request to eSewa’s server, which forwards it to the bank’s RPC endpoint.
- Real Impact: Ensures atomic transactions (either debit + credit, or neither).
- Idea: Uses RPC to call the bank’s
NTC’s BGP for Internet Routing
- Idea: Uses BGP to route traffic between Nepal and the global internet.
- How: NTC’s routers exchange path information with ISPs (e.g., Nepal Telecom → Reliance Jio → Google).
- Real Impact: When you load YouTube, BGP ensures your request takes the fastest path.
WhatsApp’s Message Queues
- Idea: Uses asynchronous message queues to deliver messages even if recipients are offline.
- How: Your message is stored in WhatsApp’s queue until the recipient connects.
- Real Impact: "Seen" status updates only when the recipient is online.
NEPSE’s NTP for Trade Timing
- Idea: Uses NTP to synchronize brokers’ clocks to Nepal Standard Time (NST).
- How: All trading systems (e.g., NEPSE, local brokers) sync to
pool.ntp.org. - Real Impact: Prevents double-trading or fraud (e.g., "What if two trades happened at the same millisecond?").
Exam Tip
Diagrams Are Your Friend
- Always draw sequence diagrams for RPC, NTP, or 2PC.
- Example: For RPC, show the client-stub → network → server-stub flow.
- For NTP, show the client → server → adjustment steps.
Compare OSI vs. TCP/IP
- Highlight layers they share (Transport, Network) and differences (OSI’s Session/Presentation).
- Example:
Layer OSI (7) TCP/IP (4) Top Application Application Middle Presentation, Session (Merged into App) Bottom Transport, Network Transport, Internet
Clock Synchronization Tricks
- Know Lamport timestamps for causal ordering.
- Know NTP for real-world time sync (mention
pool.ntp.org). - Example question: "Why does NEPSE need NTP?" → Answer: To prevent double-spending or audit conflicts.
Failure Handling
- For 2PC vs. 1PC, explain:
- 2PC: Atomicity but blocking.
- 1PC: Faster but no rollback.
- Example: "How does eSewa recover if the bank’s RPC server crashes?" → Use timeouts + retries.
- For 2PC vs. 1PC, explain:
Real-World Mapping
- Link concepts to Nepali apps:
- RPC: eSewa/bank transactions.
- Message Queues: Pathao ride assignments.
- NTP: NEPSE stock trading.
- Example: "How does Daraz use message queues during Black Friday?" → Decouples order processing from inventory checks.
- Link concepts to Nepali apps:
Avoid Common Pitfalls
- Don’t confuse address (IP) with name (domain).
- Don’t mix up synchronous (RPC) and asynchronous (Kafka) communication.
- For clocks, distinguish physical (atomic) vs. logical (Lamport).
Final Note: This unit is heavily diagram-based. Practice drawing:
- RPC flow (client → stub → network → server).
- NTP handshake (request → response → adjustment).
- 2PC protocol (prepare → commit phases).
- OSI vs. TCP/IP layer comparison.
Based on the TU BCA syllabus for Distributed System (CACS352), unit 12.
Discussion
Loading…