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-fon Node 1,g-lon Node 2). - Replication: Multiple copies for low-latency searches.
- Eventual Consistency: Updates propagate within seconds.
- Sharding: Index partitioned by URL hash (e.g.,
Exam Tip
- Define Clearly: Distinguish between distributed DB (logical unity) and federated DB (loose coupling).
- Draw Architectures: Sketch partitioning/replication diagrams in exams (e.g., Daraz’s warehouse-based sharding).
- CAP Trade-offs: Always justify why a system chooses CA, CP, or AP (e.g., eSewa = CP; WhatsApp = AP).
- Real Examples: Link concepts to eSewa (consistency), Ncell (replication), or Daraz (partitioning).
- Shortcuts:
- Homogeneous vs. Heterogeneous: Remember "same DBMS = homogeneous."
- Partitioning: "Rows = horizontal; columns = vertical."
Real 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…