CSC461 Advanced Database

Advanced DatabaseUnit 516 min read

Distributed Databases: Architectures, Fragmentation & Transparency

Unit 5 of Advanced Database explores distributed database systems, covering architectures (client-server, peer-to-peer, homogeneous/heterogeneous), data fragmentation (horizontal/vertical/derivation), replication techniques, transparency types (location, failure, migration), and real-world applications like Ncell’s cus

TAKEAWAYS:

  • Distributed databases split data across multiple nodes to improve scalability, reliability, and performance, unlike centralized systems that become bottlenecks.
  • Fragmentation (horizontal/vertical/derivation) divides data logically or physically, while replication copies data to multiple sites for fault tolerance.
  • Transparency (location, failure, migration, etc.) hides complexity from users, making distributed systems appear as a single logical database.
  • Architectures like client-server, peer-to-peer, and homogeneous/heterogeneous determine how nodes communicate and share data.
  • Consistency models (strong, eventual, causal) trade off between data accuracy and system responsiveness in distributed environments.
  • Real-world use: Ncell distributes customer records across regional servers for faster access, while eSewa replicates transaction logs to prevent data loss.

1. What is a Distributed Database?

A distributed database (DDB) is a collection of multiple interconnected databases spread across different physical locations (nodes) that appear as a single logical database to users. Unlike centralized databases, DDBs distribute data, processing, and control to improve:

  • Scalability: Handle growing data without a single server bottleneck.
  • Reliability: No single point of failure (e.g., if one server crashes, others continue).
  • Performance: Data stored closer to users reduces latency (e.g., Ncell’s regional servers).
  • Availability: Systems remain operational even during partial failures.

Key Difference from Centralized Databases:

Feature Centralized Database Distributed Database
Data Location Single site Multiple sites (nodes)
Bottleneck High (single server) Low (load balanced)
Fault Tolerance Low (single failure = crash) High (redundancy)
Scalability Limited by hardware Scales horizontally (add nodes)
Complexity Low High (network, synchronization)

2. Why Use Distributed Databases?

Real-World Example 1: Ncell’s Customer Data Distribution Ncell, Nepal’s largest telecom provider, stores customer records (subscriber details, call logs, billing) across regional data centers (Kathmandu, Pokhara, Biratnagar). This design:

  • Reduces latency: Users in Pokhara access data from the local server instead of Kathmandu.
  • Improves reliability: If the Kathmandu server fails, Pokhara users still get service.
  • Supports scalability: New regions add servers without overloading existing ones.

Real-World Example 2: eSewa’s Transaction Processing eSewa, Nepal’s leading digital payment platform, uses a distributed database to:

  • Replicate transaction logs across multiple servers to prevent data loss.
  • Fragment billing data by user region (e.g., all Kathmandu transactions on Server A, Pokhara on Server B).
  • Ensure high availability: If one server goes down, another takes over processing payments.

3. Distributed Database Architectures

Distributed databases are classified based on how nodes communicate and share data. Three primary architectures:

RequestRequestSyncSyncClient 1Client 2Server AServer BServer C
Client-Server architecture with replication between servers (e.g., Ncell’s regional servers)

A. Client-Server Architecture

  • Structure: One or more central servers store and manage data, while clients (users/applications) request services.
  • Example: Bank ATMs querying a central database.
  • Pros: Simple to implement, centralized control.
  • Cons: Server becomes a bottleneck, single point of failure.

B. Peer-to-Peer (P2P) Architecture

  • Structure: All nodes (peers) are equal; no central server. Data is distributed among peers.
  • Example: Blockchain networks (e.g., Bitcoin), where every node validates transactions.
  • Pros: No single point of failure, highly scalable.
  • Cons: Complex to manage consistency, security challenges.

C. Homogeneous vs. Heterogeneous Architectures

Type Description Example
Homogeneous All nodes use the same DBMS (e.g., all Oracle or PostgreSQL). Ncell’s regional servers (all PostgreSQL).
Heterogeneous Nodes use different DBMS (e.g., SQL Server + MySQL). eSewa might use MongoDB for transactions and PostgreSQL for analytics.

4. Data Fragmentation in Distributed Databases

Fragmentation divides a logical database into smaller physical pieces stored at different nodes. Three types:

erDiagram
  Customers ||--o{ Orders : places
  Customers {
    int customer_id PK
    string name
    string city
  }
  Orders {
    int order_id PK
    int customer_id FK
    date order_date
  }
  Customers_1 ||--o{ Orders : "Kathmandu only"
  Customers_2 ||--o{ Orders : "Pokhara only"
  Customers_1 {
    int customer_id PK
    string name
    string city 'Kathmandu'
  }
  Customers_2 {
    int customer_id PK
    string name
    string city 'Pokhara'
  }
Horizontal fragmentation of Customers table by city (Ncell’s regional data split)

A. Horizontal Fragmentation

  • Definition: Splits tables row-wise (by tuples) based on a condition.
  • Example: Fragment Customers table by region:
    -- Fragment 1: Customers in Kathmandu
    SELECT * FROM Customers WHERE city = 'Kathmandu';
    
    -- Fragment 2: Customers in Pokhara
    SELECT * FROM Customers WHERE city = 'Pokhara';
    
  • Use Case: Ncell stores subscriber data by district to reduce query time.

B. Vertical Fragmentation

  • Definition: Splits tables column-wise (by attributes).
  • Example: Split Orders table into Order_Details (columns: order_id, product_id, quantity) and Order_Customer (columns: order_id, customer_id, date).
  • Use Case: eSewa separates transaction metadata (e.g., amount, timestamp) from user details for security.

C. Derived Fragmentation

  • Definition: Creates fragments by deriving data from other fragments (e.g., summaries, aggregates).
  • Example: A Sales_Summary fragment derived from daily transaction fragments.
  • Use Case: Daraz’s analytics team pre-computes monthly sales summaries for faster reporting.

5. Data Replication Techniques

Replication copies data to multiple nodes to improve reliability and performance. Two main strategies:

A. Primary Copy Replication

  • How it works:
    1. One node is the primary copy (master), others are backups.
    2. All updates go to the primary, which propagates changes to backups.
  • Example: WhatsApp replicates message logs across servers to ensure no messages are lost.
  • Pros: Strong consistency, simple to implement.
  • Cons: Primary becomes a bottleneck; if primary fails, system halts.

B. Full Replication

  • How it works: All nodes have a complete copy of the data. Updates are synchronized.
  • Example: Google’s global DNS system replicates DNS records across data centers.
  • Pros: High availability, no single point of failure.
  • Cons: High storage overhead, complex synchronization.

C. Selective Replication

  • How it works: Only specific fragments are replicated (not the entire database).
  • Example: NEPSE (Nepal Stock Exchange) replicates only high-frequency stock data to regional servers.
  • Pros: Saves storage, reduces network traffic.
  • Cons: Partial failures may still occur.

6. Transparency in Distributed Databases

Transparency hides the distributed nature of the database from users/applications. Seven types:

sequenceDiagram
  participant User
  participant App
  participant LocalDB
  participant RemoteDB
  User->>App: Request data (e.g., Ncell balance)
  App->>LocalDB: Check cache
  alt Not found
    App->>RemoteDB: Fetch from Kathmandu server
    RemoteDB-->>App: Return data
  end
  App-->>User: Display result
  Note right of User: User sees no difference between local/remote data
Location transparency: App fetches data from remote DB without user awareness (e.g., eSewa)
Type Definition Example
Location Users unaware of where data is stored. Querying SELECT * FROM Customers returns data from any node.
Failure System masks node failures (e.g., crashes). If Server A fails, queries route to Server B.
Migration Data can move between nodes without user knowledge. Ncell migrates a subscriber’s data to a new regional server.
Replication Users see a single copy, though data is replicated. eSewa’s transaction logs appear as one copy.
Concurrency Manages simultaneous access without conflicts. Two users booking the same Pathao ride get resolved automatically.
Performance Hides latency differences between nodes. Fast response time regardless of server location.
Scalability System appears to have unlimited capacity. Adding more nodes increases capacity seamlessly.

7. Distributed Database Design Techniques

Designing a DDB involves:

  1. Fragmentation: Decide how to split data (horizontal/vertical/derived).
  2. Allocation: Assign fragments to nodes (e.g., by proximity to users).
  3. Replication: Choose replication strategy (primary, full, selective).
  4. Transaction Management: Ensure ACID properties across nodes.

Worked Example: Daraz’s Order Processing System

Scenario: Daraz wants to distribute order data across Kathmandu, Pokhara, and Biratnagar servers.

  1. Fragmentation:

    • Horizontal: Split Orders by city:
      -- Kathmandu Orders (Fragment 1)
      SELECT * FROM Orders WHERE city = 'Kathmandu';
      
      -- Pokhara Orders (Fragment 2)
      SELECT * FROM Orders WHERE city = 'Pokhara';
      
    • Vertical: Separate Order_Details (products, quantities) from Customer_Info.
  2. Allocation:

    • Store Kathmandu orders on Server A, Pokhara on Server B, etc.
    • Replicate Customer_Info fully for fast login.
  3. Replication:

    • Use selective replication for high-demand products (e.g., recharge cards).
    • Use primary copy for inventory updates (only one node updates stock).
  4. Transaction Handling:

    • If a user in Pokhara orders a product in Kathmandu:
      • Query routes to Kathmandu server for stock.
      • Payment processed locally (Pokhara server).
      • Order confirmation sent from Pokhara server.

8. Consistency Models in Distributed Databases

Consistency ensures all nodes see the same data after updates. Three models:

Model Definition Trade-off Example
Strong All nodes see updates instantly. High latency, complex. Bank transactions (ATM withdrawals).
Eventual Updates propagate eventually (not instantly). Fast writes, eventual consistency. Social media posts (Facebook, Twitter).
Causal Preserves cause-effect order (e.g., if A updates before B, B sees A’s update). Balances speed and consistency. WhatsApp message ordering.

Example: Kathmandu Traffic Routes Database

  • Strong Consistency: Used for real-time traffic updates (e.g., Pathao’s live traffic data must reflect accidents instantly).
  • Eventual Consistency: Used for historical data (e.g., monthly traffic reports can be updated later).

9. Challenges in Distributed Databases

  • Network Latency: Delays in communication between nodes.
  • Data Consistency: Ensuring all nodes agree on data values.
  • Security: Protecting data across multiple locations (e.g., Ncell’s subscriber data).
  • Transaction Management: Maintaining ACID properties across nodes.
  • Cost: Higher infrastructure and maintenance costs.

In the Real World

  1. Ncell’s Distributed Customer Database

    • Idea Used: Horizontal fragmentation + replication.
    • How: Customer records are split by district (e.g., Kathmandu, Pokhara) and replicated across regional servers. If a server fails, another takes over, ensuring high availability.
    • Impact: Faster SIM registration, reduced call drops during peak hours.
  2. eSewa’s Transaction Processing

    • Idea Used: Primary copy replication + vertical fragmentation.
    • How: Transaction logs are stored as the primary copy on a master server, while user profiles are vertically fragmented (separate tables for User_Details, Transaction_History). Replicas ensure no data loss during server failures.
    • Impact: Millions of transactions processed daily without downtime.
  3. Daraz’s Order Fulfillment System

    • Idea Used: Selective replication + horizontal fragmentation.
    • How: Orders are horizontally fragmented by city, and best-selling products (e.g., recharge cards) are selectively replicated across all servers. This reduces latency for popular items.
    • Impact: Faster checkout, lower cart abandonment rates.
  4. NEPSE’s Stock Data Distribution

    • Idea Used: Derived fragmentation + eventual consistency.
    • How: Real-time stock prices are distributed globally, while daily summaries are derived and stored locally. Traders see eventual consistency (prices update within seconds).
    • Impact: Low-latency trading for investors across Nepal.

Exam Tip

  1. Define Clearly: Start with precise definitions (e.g., "A distributed database is a collection of interconnected databases...").
  2. Compare Tables: Use comparison tables (e.g., centralized vs. distributed, fragmentation types) to score easy marks.
  3. Real-World Links: Always tie examples to Nepali companies (Ncell, eSewa, Daraz, NEPSE) or global tech (Google, WhatsApp). Examiners love context!
  4. Diagrams > Text: Draw mermaid diagrams for:
    • Distributed database architectures (client-server vs. P2P).
    • Data fragmentation (horizontal/vertical).
    • Replication strategies (primary copy vs. full).
  5. ACID in Distributed Systems: Explain how distributed transactions maintain atomicity/consistency (e.g., two-phase commit protocol).
  6. Transparency Types: Memorize the 7 transparencies (location, failure, etc.) and give one example each.
  7. Worked Examples: Solve fragmentation/allocation problems step-by-step (e.g., "Fragment the Employees table by department").
  8. Consistency Models: Know the trade-offs between strong, eventual, and causal consistency. Use traffic systems or banking as examples.

classDiagram
    class DistributedDatabase {
        +fragmentData()
        +replicateData()
        +ensureTransparency()
    }

    class Node {
        -dataFragments
        -replicationStatus
        +processQuery()
    }

    class User {
        +submitQuery()
        +receiveResult()
    }

    User --> DistributedDatabase : queries
    DistributedDatabase --> Node : distributes
    Node --> DistributedDatabase : synchronizes
sequenceDiagram
    participant User
    participant Application
    participant DistributedDB
    participant Node1
    participant Node2

    User->>Application: Request data
    Application->>DistributedDB: Query (e.g., SELECT * FROM Orders WHERE city='Kathmandu')
    DistributedDB->>Node1: Route to Kathmandu fragment
    Node1-->>DistributedDB: Return order data
    DistributedDB-->>Application: Consolidate results
    Application-->>User: Display results

In the real world

  • Ncell’s Regional Servers: Uses horizontal fragmentation to split subscriber data by district (Kathmandu, Pokhara, Biratnagar) for faster local queries and primary copy replication to sync billing updates across regions.
  • eSewa’s Payment System: Employs vertical fragmentation (separating transaction metadata from user details) and full replication of transaction logs to prevent data loss during server failures.
  • WhatsApp’s Global Sync: Relies on primary copy replication with multiple data centers to ensure messages are delivered even if one server fails, demonstrating failure transparency.

Based on the TU BSc CSIT syllabus for Advanced Database (CSC461), unit 5.

Discussion

Loading…