A query optimization in distributed database systems
Chin‐Wan Chung · Deep Blue (University of Michigan) · 1983
This research is concerned with a model and a method of minimizing the inter-site data traffic incurred by a query in distributed relational database systems. In order to process a query which references data from multiple sites in a computer network, portions of the database at other sites have to be transferred to the user's site. The usual methodology for distributed query processing consists of reducing the referenced relations using a sequence of semijoin operations after initial local processing. The mathematical model has been developed to determine an optimal sequence of semijoins which minimizes the total inter-site data flow in processing a distributed query. The core of this model is a method which efficiently and accurately estimates the size of an intermediate result of a query. In particular, the assumption that joining attributes are independent during the processing of a query by a sequence of semijoins has been relaxed. Since the distributed query optimization problem is known to be NP-hard, a heuristic algorithm has been developed to determine a low-cost sequence of semijoins. The efficiency of the algorithm is increased by partitioning the set of joining attributes into blocks and sequencing these blocks, as well as by a straightforward, yet effective sequencing among the semijoins between the joining attributes inside a block. The algorithm decreases the cost of a query by selecting the low-cost, highly reductive semijoins first. Cost comparisons with the existing algorithms have been provided. The time complexity of the main features of the algorithm has been analytically derived. The algorithm has been implemented in PASCAL. The tests show that the scheduling time for a sequence of semijoins for a reasonable size query is less than 0.05 seconds when the program is executed by a main-frame computer.