Adaptive methods in parallel databases

Shaibal Roy · 1992

Effective parallel processing of relational databases requires an equitable distribution of workload among available resources. Unfortunately, nonuniformity (skew) and correlation in the distribution of data and queries cause inequitable allocation of workload to resources. This dissertation investigates three stages where a parallel relational database system can adapt to nonuniformity in data and queries to improve resource utilization and reduce response time. The emphasis is on intra-query parallelism in processing of large queries. Data partitioning. The performance of a horizontally partitioned relational database system depends on the semantics of data retrievals. We present and analyze a combinatorial characterization of the semantics embedded in applications, and use this characterization in evaluating the performance of partitioning strategies. Query optimization. In presence of skew and correlations in data, bottlenecks and congestion make performance of query plans unpredictable. We present and analyze a new approach of choosing among alternative query plans, which exploits partial evaluation of several alternative plans. Dynamic load balancing. Data skew is known to cause imbalance in parallel hash-join computation. We present an algorithm that achieves load balancing by adapting to data skew dynamically, and its implementation on the IBM RP3 parallel processor. Our algorithm exploits the random access capability of main-memory databases to minimize the overhead of adaptation. Our experiments on the 64 processor RP3 have shown almost linear speedup for as many as 60 processors and beyond. In contrast, almost no speedup could be obtained beyond 10 processors without load balancing--even for only moderate skew. The work presented in this dissertation is based on two paradigms. (1) When the operands of a relational operation do not fit in main memory, accurate estimates of skew and correlation, obtained using random samples, are made available for adaptation through preprocessing. We show how memory-resident random samples can be used in obtaining such estimates without accessing disks. (2) When the operands fit in main memory, the dynamic adaptation techniques we develop ensure equitable distribution of workload without any preprocessing. The sum of the above methods constitutes a set of adaptive techniques that solve critical problems in the effective use of parallel computers in all phases of the processing of large database queries.

Read the paper · More papers on PaperTik