Elective Distributed and Object Oriented Database

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.
MetadataMetadataMetadataQuery FragmentSite 1Site 2Site NGlobal Catalog
Basic distributed DB architecture with a global catalog for fragmentation metadata.

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:

Fragment 1 (WHERE city='Kathmandu')Fragment 2 (WHERE city='Pokhara')Horizontal FragmentationFragment A (customer_id, name)Fragment B (balance, city)Vertical FragmentationFragment X (WHERE city='Kathmandu' AND balance>10000)Mixed FragmentationQuery
Comparison of fragmentation types for query decomposition.

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

CUSTOMER (Original)FRAGMENT_KATHMANDU(city='Kathmandu')FRAGMENT_POKHARA(city='Pokhara')
Horizontal fragmentation splits a table by rows based on a condition (e.g., city).

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:

  1. Localizes the query to the nearest server (e.g., Kathmandu node).
  2. Decomposes it into subqueries (e.g., check balance, deduct amount).
  3. 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:

  1. 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
    
  2. Execute locally at Site 2 (no data transfer from Site 1).
  3. 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 results

6. 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_id to 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.

Query (city='Pokhara', balance>10000)UserQuery ProcessorSite 1 (Kathmandu)Site 2 (Pokhara)
Localization in action: Query executed only at Site 2, minimizing data transfer.

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

Discussion

Loading…