Distributed and Object Oriented DatabaseUnit 39 min read
Query Processing, Decomposition & Localization in Distributed DBs
Unit 3 of Distributed and Object Oriented Database covers how queries are processed across distributed databases, how they are decomposed into subqueries for local sites, and how localization techniques optimize performance. Learn fragmentation strategies, query optimization, and real-world applications like eSewa’s tr
TAKEAWAYS:
- Queries in distributed databases are decomposed into subqueries for local processing using horizontal/vertical fragmentation and replication.
- Localization reduces network traffic by processing queries at the nearest site, improving efficiency.
- Query optimization techniques like predicate pushdown and join ordering minimize data transfer.
- Fragmentation strategies (e.g., derived, allocated) determine how data is split across nodes.
- Distributed query processing involves parsing, optimization, and execution across multiple sites.
- Real-world examples include eSewa’s transaction processing (localization) and Ncell’s billing system (fragmentation).
1. Introduction to Distributed Query Processing
Distributed databases process queries across multiple nodes (sites) connected via a network. Unlike centralized databases, queries must be decomposed into smaller subqueries, executed locally, and then reintegrated to produce results. This process involves:
- Query parsing: Breaking down a global query into local subqueries.
- Query optimization: Minimizing data transfer and computation time.
- Query execution: Running subqueries at local sites and combining results.
Why Decompose Queries?
- Network efficiency: Avoids transferring large datasets across sites.
- Load balancing: Distributes processing across multiple nodes.
- Fault tolerance: If one site fails, others can still process parts of the query.
2. Query Decomposition Techniques
Queries are decomposed based on fragmentation schemes (how data is split across sites). Common methods include:
A. Horizontal Fragmentation
Data is split row-wise (e.g., customers in Kathmandu vs. Pokhara). Example:
-- Global table: CUSTOMER (customer_id, name, city, balance)
-- Fragment 1 (Kathmandu): WHERE city = 'Kathmandu'
-- Fragment 2 (Pokhara): WHERE city = 'Pokhara'
Mermaid Diagram: Horizontal Fragmentation
B. Vertical Fragmentation
Data is split column-wise (e.g., customer details vs. transaction history). Example:
-- Global table: CUSTOMER (customer_id, name, balance, transaction_id)
-- Fragment 1: (customer_id, name, balance)
-- Fragment 2: (customer_id, transaction_id)
C. Mixed Fragmentation
Combines horizontal and vertical fragmentation for complex splits.
3. Localization: Reducing Data Transfer
Localization ensures queries are processed at the nearest site to minimize network traffic. Techniques include:
- Predicate pushdown: Applying filters (WHERE clauses) early to reduce data transfer.
- Join ordering: Executing joins at the site with the most relevant data.
- Semijoin optimization: Sending only necessary attributes for joins.
Example: eSewa Transaction Processing When a user pays a bill, eSewa:
- Localizes the query to the nearest server (e.g., Kathmandu node).
- Decomposes it into subqueries (e.g., check balance, deduct amount).
- Reintegrates results without transferring full customer records.
4. Query Optimization Strategies
Optimization reduces response time and network overhead. Key techniques:
| Technique | Description | Example |
|---|---|---|
| Predicate Pushdown | Apply filters early to reduce data transfer. | WHERE city = 'Kathmandu' before sending data. |
| Join Ordering | Execute joins at the site with the most data. | Join CUSTOMER and TRANSACTION at the customer’s local site. |
| Semijoin | Send only required attributes for joins. | Instead of full TRANSACTION table, send only transaction_id. |
| Materialized Views | Precompute frequent queries to avoid recomputation. | Store SUM(balance) for monthly reports. |
5. Worked Example: Distributed Query Execution
Scenario: Find all customers in Pokhara with a balance > 10,000. Database:
- Site 1 (Kathmandu):
CUSTOMER_KTM(city = 'Kathmandu') - Site 2 (Pokhara):
CUSTOMER_POK(city = 'Pokhara')
Steps:
- Decompose the query:
-- Global query: SELECT * FROM CUSTOMER WHERE city = 'Pokhara' AND balance > 10000 -- Local query at Site 2: SELECT * FROM CUSTOMER_POK WHERE balance > 10000 - Execute locally at Site 2 (no data transfer from Site 1).
- Return results directly to the user.
Mermaid Diagram: Query Execution Flow
sequenceDiagram
User->>+QueryProcessor: Submit query (city='Pokhara', balance>10000)
QueryProcessor->>Site2: Decompose to local query
Site2->>Site2: Execute: SELECT * FROM CUSTOMER_POK WHERE balance > 10000
Site2-->>-QueryProcessor: Return results
QueryProcessor-->>User: Display results6. Fragmentation Strategies
Data fragmentation determines how queries are processed. Common strategies:
| Strategy | Description | Use Case |
|---|---|---|
| Derived | Data is split based on a condition (e.g., city = 'Pokhara'). |
eSewa’s regional customer data. |
| Allocated | Data is pre-assigned to sites (e.g., all Pokhara customers go to Site 2). | Ncell’s regional billing systems. |
| Hybrid | Combines derived and allocated fragmentation. | Daraz’s inventory management (some items allocated, others derived). |
7. Real-World Applications
A. eSewa (Nepal)
- Localization: When you pay a bill, eSewa routes your request to the nearest server (e.g., Kathmandu or Pokhara) to minimize latency.
- Fragmentation: Customer data is horizontally fragmented by region.
B. Ncell Billing System
- Derived Fragmentation: Billing records are split by district (e.g.,
DISTRICT_KTM,DISTRICT_POK). - Query Optimization: Predicate pushdown ensures only relevant records are processed.
C. Daraz Order Processing
- Vertical Fragmentation: Order details (customer_id, product_id) are separated from payment data for security.
- Semijoin: When checking stock, Daraz sends only
product_idto the warehouse site.
8. Advantages and Disadvantages
| Advantage | Disadvantage |
|---|---|
| Faster query response (local processing) | Complex query decomposition required. |
| Reduced network traffic | Higher storage overhead (replication). |
| Fault tolerance (site failures) | Need for distributed transaction management. |
9. Exam Tip
- Focus on decomposition: Always show how a global query splits into local subqueries.
- Localization is key: Explain how predicate pushdown/semijoin reduces data transfer.
- Compare fragmentation: Know when to use horizontal vs. vertical fragmentation.
- Real-world tie-ins: Relate examples to eSewa, Ncell, or Daraz (e.g., "How would eSewa optimize a transaction query?").
- Draw diagrams: Mermaid sequence diagrams for query execution and fragmentation are high-value answers.
Based on the TU BSc CSIT syllabus for Distributed and Object Oriented Database, unit 3.
Discussion
Loading…