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:
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
Customerstable 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
Orderstable intoOrder_Details(columns: order_id, product_id, quantity) andOrder_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_Summaryfragment 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:
- One node is the primary copy (master), others are backups.
- 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 dataLocation 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:
- Fragmentation: Decide how to split data (horizontal/vertical/derived).
- Allocation: Assign fragments to nodes (e.g., by proximity to users).
- Replication: Choose replication strategy (primary, full, selective).
- 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.
Fragmentation:
- Horizontal: Split
Ordersby 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) fromCustomer_Info.
- Horizontal: Split
Allocation:
- Store Kathmandu orders on Server A, Pokhara on Server B, etc.
- Replicate
Customer_Infofully for fast login.
Replication:
- Use selective replication for high-demand products (e.g., recharge cards).
- Use primary copy for inventory updates (only one node updates stock).
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.
- If a user in Pokhara orders a product in Kathmandu:
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
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.
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.
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.
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
- Define Clearly: Start with precise definitions (e.g., "A distributed database is a collection of interconnected databases...").
- Compare Tables: Use comparison tables (e.g., centralized vs. distributed, fragmentation types) to score easy marks.
- Real-World Links: Always tie examples to Nepali companies (Ncell, eSewa, Daraz, NEPSE) or global tech (Google, WhatsApp). Examiners love context!
- Diagrams > Text: Draw mermaid diagrams for:
- Distributed database architectures (client-server vs. P2P).
- Data fragmentation (horizontal/vertical).
- Replication strategies (primary copy vs. full).
- ACID in Distributed Systems: Explain how distributed transactions maintain atomicity/consistency (e.g., two-phase commit protocol).
- Transparency Types: Memorize the 7 transparencies (location, failure, etc.) and give one example each.
- Worked Examples: Solve fragmentation/allocation problems step-by-step (e.g., "Fragment the
Employeestable by department"). - 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 : synchronizessequenceDiagram
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 resultsIn 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…