Elective Distributed and Object Oriented Database

Distributed and Object Oriented DatabaseUnit 16 min read

Distributed DB: Definitions, Architectures, Challenges & Real Systems

Unit 1 of Distributed and Object Oriented Database introduces core concepts of distributed databases—what they are, how they differ from centralized systems, their architectures (homogeneous/heterogeneous), key challenges (fault tolerance, replication, concurrency), and real-world implementations like eSewa’s transacti

What is a Distributed Database?

A distributed database (DDB) is a collection of multiple interconnected databases spread across different physical locations, appearing as a single logical database to users. Unlike centralized databases, data is stored on multiple nodes (computers or servers) connected via a network, enabling decentralized data management.

Key Characteristics

mindmap
  root((Distributed Database))
    Characteristics
      Data Distribution: Sharded or replicated across nodes
      Logical Unity: Appears as one system to users
      Autonomy: Local control at each site
      Transparency: Hides distribution details
    Components
      Nodes: Physical storage sites
      Network: Connects nodes (LAN/WAN)
      DDBMS: Software managing distribution
    Types
      Homogeneous: Same DBMS at all nodes
      Heterogeneous: Mixed DBMS (e.g., Oracle + PostgreSQL)

Why distribute data?

  • Scalability: Handle large volumes (e.g., Ncell’s subscriber data).
  • Fault Tolerance: If one node fails, others continue (e.g., eSewa’s payment system).
  • Performance: Local processing reduces latency (e.g., Daraz’s regional warehouses).
  • Geographic Reach: Serve global users (e.g., Google’s distributed search indexes).

Centralized vs. Distributed Databases

Feature Centralized Database Distributed Database
Data Location Single site Multiple sites
Scalability Limited by hardware Scales horizontally
Fault Tolerance Single point of failure Redundancy across nodes
Network Dependency None Critical for communication
Complexity Low High (distribution, synchronization)
Example Bank’s local ATM system NEPSE’s stock trading across exchanges

Architectures of Distributed Databases

1. Homogeneous DDB

All nodes use the same DBMS (e.g., PostgreSQL at all sites). Simplifies management but limits flexibility. Example: NTC’s fiber-optic network management uses homogeneous databases for uniform monitoring.

2. Heterogeneous DDB

Nodes use different DBMS (e.g., Oracle for transactions, MongoDB for logs). Requires middleware for interoperability. Example: Pathao’s backend mixes SQL (user data) and NoSQL (ride history) databases.


How Data is Distributed: Partitioning and Replication

Partitioning (Sharding)

Data is horizontally split by rows (e.g., users by region) or vertically by columns (e.g., customer details vs. orders).

Worked Example: Daraz’s Order Queue

  • Problem: Daraz needs to process orders from Kathmandu and Pokhara without overloading a single server.
  • Solution: Orders are partitioned by warehouse location.
    • Kathmandu orders → Node 1 (Kathmandu server).
    • Pokhara orders → Node 2 (Pokhara server).
  • Query: "Show orders for user ID 12345 from Kathmandu."
    • The DDBMS routes the query to Node 1 (localization).

Replication

Copies of data are stored at multiple nodes to improve read performance or ensure availability.

Example: WhatsApp replicates chat histories across data centers to ensure messages are never lost during outages.


Challenges in Distributed Databases

1. Data Consistency

Ensuring all nodes have the same data after updates. Solutions:

  • Strong Consistency: Wait for all replicas to confirm (slow but accurate).
  • Eventual Consistency: Replicas sync over time (faster, used by Facebook).
  • CAP Theorem: Choose 2 out of 3:
    • Consistency (all nodes agree),
    • Availability (no node failures),
    • Partition Tolerance (network issues).

2. Fault Tolerance

  • Node Failure: If a node crashes, the system must recover (e.g., automatic failover in Ncell’s billing).
  • Network Partition: Split-brain scenarios require quorum-based writes (e.g., 2/3 nodes must agree).

3. Concurrency Control

Multiple transactions accessing shared data need locking or optimistic concurrency (e.g., bank transfers between eSewa and Khalti).

4. Security and Access Control

  • Authentication: Users must prove identity (e.g., OTP for eSewa payments).
  • Authorization: Role-based access (e.g., NEPSE traders vs. investors).

Real-World Applications

1. eSewa’s Payment System

  • Problem: Millions of transactions daily across Nepal.
  • Solution:
    • Partitioning: Transactions split by district (e.g., Kathmandu vs. Pokhara nodes).
    • Replication: Critical transaction logs replicated to backup servers.
    • Consistency: Strong consistency for fund transfers (CAP prioritizes consistency over availability).

2. Ncell’s Billing Database

  • Challenge: 20M+ subscribers; billing data must be fault-tolerant.
  • Design:
    • Homogeneous DDB: Oracle at all regional data centers.
    • Replication: Daily backups to a cloud-based replica.
    • Partitioning: Users partitioned by SIM prefix (e.g., 98XXXX vs. 97XXXX).

3. Google’s Search Index

  • Scale: Billions of web pages indexed globally.
  • Approach:
    • Sharding: Index partitioned by URL hash (e.g., a-f on Node 1, g-l on Node 2).
    • Replication: Multiple copies for low-latency searches.
    • Eventual Consistency: Updates propagate within seconds.

Exam Tip

  1. Define Clearly: Distinguish between distributed DB (logical unity) and federated DB (loose coupling).
  2. Draw Architectures: Sketch partitioning/replication diagrams in exams (e.g., Daraz’s warehouse-based sharding).
  3. CAP Trade-offs: Always justify why a system chooses CA, CP, or AP (e.g., eSewa = CP; WhatsApp = AP).
  4. Real Examples: Link concepts to eSewa (consistency), Ncell (replication), or Daraz (partitioning).
  5. Shortcuts:
    • Homogeneous vs. Heterogeneous: Remember "same DBMS = homogeneous."
    • Partitioning: "Rows = horizontal; columns = vertical."

database server rack in data centerReal photo of a server rack with labeled components (switches, RAID arrays, cooling). (Image: Aaron Hall, CC BY-SA 2.0, via Wikimedia Commons)

Based on the TU BSc CSIT syllabus for Distributed and Object Oriented Database, unit 1.

Discussion

Loading…