Elective Distributed Networking

Distributed NetworkingUnit 310 min read

Communication Paradigms: RPC, Message Passing, Remote Objects & Middleware

Unit 3 of Distributed Networking explores how processes communicate across networks—Remote Procedure Calls (RPC), message-passing models, remote objects, and middleware architectures. It covers mechanisms, protocols, and real-world trade-offs in distributed systems like eSewa’s transaction flows or Pathao’s ride-matchi

TAKEAWAYS:

  • RPC lets a client call a remote procedure as if it were local, hiding network complexity via stubs and skeletons.
  • Message-passing (synchronous/asynchronous) is the foundation for distributed coordination, with queues (e.g., Pathao’s ride requests) ensuring reliability.
  • Remote objects (e.g., CORBA, Java RMI) abstract location transparency, but require serialization and marshaling.
  • Middleware (e.g., message brokers like RabbitMQ) decouples senders/receivers, enabling scalability but adding latency.
  • Trade-offs: RPC is simple but fragile; message-passing is robust but complex; remote objects offer transparency at a cost.
  • Real-world tie: eSewa uses asynchronous message-passing for payment confirmations to avoid blocking users during high traffic.

1. Communication Paradigms: Definitions and Models

Distributed systems rely on interprocess communication (IPC) across machines. The three primary paradigms are:

A. Remote Procedure Call (RPC)

  • Definition: A client invokes a procedure on a remote server as if it were local. The client stub marshals arguments, sends them over the network, and the server stub unmarshals them to call the actual procedure.
  • Key Components:
    • Client Stub: Local proxy that packages calls into network messages.
    • Server Stub: Receives messages and invokes the real procedure.
    • Transport Layer: Handles data transmission (TCP/UDP).
    • Marshalling/Unmarshalling: Converts data structures to/from byte streams.
sequenceDiagram
    participant Client
    participant Server
    Client->>Server: RPC_Call(params)
    Server->>Server: Execute procedure
    Server-->>Client: Return result
Typical RPC request‑response flow

B. Message-Passing

  • Definition: Processes exchange discrete messages (e.g., JSON, XML) via queues or direct links. No shared memory; explicit send/receive operations.
  • Types:
    • Synchronous: Sender waits for acknowledgment (e.g., HTTP requests).
    • Asynchronous: Sender continues without waiting (e.g., email, Pathao ride alerts).
    • Group Communication: One-to-many (e.g., stock price broadcasts).

Mermaid Diagram: Message-Passing Models

flowchart TD
    Sync["Synchronous (request‑response)"] -->|"waits for ACK"| Receiver["Receiver"]
    Async["Asynchronous (fire‑and‑forget)"] -->|"no wait"| Receiver
    Group["Group Communication (publish‑subscribe)"] -->|"one‑to‑many"| Subscribers["Multiple Receivers"]

C. Remote Objects

  • Definition: Objects on different machines communicate as if they were local, using location transparency. Examples: CORBA, Java RMI, .NET Remoting.
  • Mechanism:
    • Object Request Broker (ORB): Middleware that routes calls (e.g., CORBA’s IIOP protocol).
    • Interface Definition Language (IDL): Defines object interfaces (e.g., interface BankAccount { void deposit(float amount); }).

Worked Example: eSewa’s Payment Flow

  1. User requests payment via eSewa app (client stub).
  2. App marshals PaymentRequest (amount, recipient) into JSON.
  3. RPC call to eSewa server stub → server validates and processes payment.
  4. Server stub returns PaymentConfirmation to client stub.
  5. Real-world twist: During Diwali, eSewa uses asynchronous queues to buffer high-volume transactions, preventing client timeouts.

2. How RPC Works: Step-by-Step Trace

Assume a client calls balance() on a remote bank account:

  1. Client Stub:

    • Marshals balance() into a message: {method: "balance", args: [], clientID: "user123"}.
    • Sends over TCP to server IP 203.123.45.67:8080.
  2. Network Transport:

    • TCP ensures reliable delivery (ACKs, retransmissions).
    • Firewall rules allow port 8080 (e.g., Ncell’s API gateway).
  3. Server Stub:

    • Unmarshals message → invokes BankAccount.balance().
    • Returns {balance: 5000}.
  4. Response Path:

    • Server stub marshals response → sends back to client.
    • Client stub unmarshals → displays balance.

Potential Failures:

  • Network Partition: Client stub times out (retry after 5s).
  • Server Crash: Stub detects no response → throws RemoteException.

3. Message-Passing: Queues and Protocols

A. Direct vs. Indirect Communication

Model Example Pros Cons
Direct HTTP requests Low latency Tight coupling
Indirect RabbitMQ, Kafka Decoupled, scalable Higher complexity

B. Protocols

  • HTTP/REST: Stateless, JSON/XML payloads (e.g., Daraz’s order API).
  • AMQP (Advanced Message Queuing Protocol): Used by RabbitMQ for reliable queues.
  • WebSockets: Full-duplex (e.g., Pathao’s real-time ride tracking).

Mermaid Diagram: AMQP Message Flow

sequenceDiagram
    participant Client
    participant Broker as RabbitMQ
    participant Server
    Client->>Broker: Publish("ride_request", {user: "A", dest: "KTM"})
    Broker->>Server: Deliver to "DriverQueue"
    Server->>Broker: ACK
    Broker-->>Client: Confirmation

C. Real-World Example: Pathao’s Ride Matching

  1. User Request: App sends ride_request message to Pathao’s queue.
  2. Driver Assignment: Pathao’s server dequeues → matches with nearest driver via geohashing.
  3. Asynchronous Update: Driver’s app receives ride_assigned message via WebSocket.

4. Remote Objects: CORBA vs. Java RMI

Feature CORBA Java RMI
Language Multi-language (IDL) Java-only
Protocol IIOP (TCP) Java-specific (serialization)
Use Case Enterprise (e.g., NEPSE trading) Internal Java apps

Example: NEPSE’s Trading System

  • Uses CORBA for location-transparent stock quotes.
  • Brokers invoke getPrice("NEPSE") on a remote object without knowing its host.

5. Middleware: The Glue of Distributed Systems

Middleware abstracts low-level details (networking, serialization) and provides:

  • Transparency: Location, migration, replication.
  • Fault Tolerance: Retries, checkpoints (e.g., banks’ transaction logs).
  • Scalability: Load balancing (e.g., Daraz’s CDN).
ApplicationApp codeMiddlewareMiddleware servicesOperating SystemOS services
Middleware sits between applications and the OS

Types of Middleware:

  1. Message-Oriented: RabbitMQ, Kafka (e.g., NTC’s SMS gateway).
  2. Object Request Brokers: CORBA, .NET Remoting.
  3. Transaction Processors: IBM MQ for financial settlements.

6. Trade-offs and Challenges

Paradigm Advantages Disadvantages When to Use
RPC Simple, familiar syntax Brittle (network failures break calls) Internal microservices
Message-Passing Decoupled, scalable Complex error handling High-throughput systems (e.g., Pathao)
Remote Objects Location transparency Performance overhead (serialization) Enterprise legacy systems

Key Challenges:

  • Latency: RPC adds ~50–200ms RTT (critical for stock trading).
  • Fault Tolerance: Message queues must persist unacknowledged messages (e.g., Ncell’s SMS retries).
  • Security: RPC requires authentication (e.g., eSewa’s OAuth tokens).

7. Current Developments

  • gRPC: Google’s RPC framework using Protocol Buffers (faster than JSON).
  • Serverless Messaging: AWS SQS, Azure Service Bus (pay-per-use queues).
  • Blockchain: Decentralized message-passing (e.g., Ethereum smart contracts).

Example: Daraz’s Order Fulfillment

  1. Synchronous RPC: User clicks "Order" → Daraz’s cart service calls inventory RPC.
  2. Asynchronous Queue: Inventory update → shipping service via Kafka.
  3. Webhook: Customer receives order_shipped via WebSocket.

In the Real World

  1. eSewa’s Payment System

    • Uses asynchronous message-passing (RabbitMQ) to decouple payment initiation from bank confirmation.
    • During festivals, queues buffer 10,000+ transactions/sec to prevent client timeouts.
  2. Pathao’s Ride Matching

    • Remote objects: Driver locations are exposed as remote objects (e.g., Driver.getLocation()).
    • Message queues: Ride requests are published to a topic; nearest driver subscribes via geofencing.
  3. NEPSE’s Trading Platform

    • CORBA-based remote objects for real-time stock price updates across brokers.
    • Synchronous RPC for order execution (low latency required).

Exam Tip

  1. Diagrams Are Mandatory:

    • Draw RPC architecture (client stub → network → server stub).
    • Sketch message-passing flows (e.g., producer → queue → consumer).
    • Label middleware layers (e.g., API gateway → message broker → DB).
  2. Compare Paradigms:

    • Contrast RPC’s simplicity with message-passing’s reliability.
    • Explain why Java RMI is less flexible than CORBA for multi-language systems.
  3. Real-World Scenarios:

    • Relate eSewa’s queues to asynchronous message-passing.
    • Link Pathao’s geohashing to remote object location transparency.
  4. Common Pitfalls:

    • Marshalling errors: Forgetting to handle NotSerializableException in Java RMI.
    • Network partitions: RPC fails silently; message-passing requires explicit retries.
  5. Short-Answer Tricks:

    • RPC: "Client stub marshals → network → server stub unmarshals."
    • Message-passing: "Decoupled via queues; sender/receiver unaware of each other."
    • Middleware: "Provides transparency, fault tolerance, and scalability."

Based on the TU BSc CSIT syllabus for Distributed Networking, unit 3.

Discussion

Loading…