Exploiting trade-offs in the design of fault-tolerant distributed databases
Amitabh Shah · 1990
In a fault-tolerant distributed database, data replication can improve the availability of data as well as the response times of transactions, but only to a point. More replication can limit concurrency of transaction execution in the event of communication failures that partition the database. It can also increase the execution time of the replica management protocols. In fact, there are trade-offs involved between the amount of replication, the availability of data, and the transaction response times; the trade-offs are a function of the failure characteristics and the transaction access patterns. This work focuses on analyzing the nature of these trade-offs and develops strategies for exploiting them to enhance system performance. First, the trade-offs between replication and availability are statically analyzed--where it is assumed that the underlying communication network has partitioned into several groups--and upper bounds on availability are derived for different cases of data replication and data access patterns. It is shown that partial replication can improve availability despite partition failures, if the partitioned database satisfies a particular property of data access. To ensure that a database indeed satisfies this property with a high probability, a network design strategy is developed. This is done by first identifying the most desirable partition of the set of sites--one characterized by high transaction access within every group of the partition, and low access across the groups--for a given cost of connecting the database sites, and then defining the connectivity of the database in a hierarchical manner. While the problem of finding the most desirable partition is shown to be NP-hard under most natural definitions of desirability, excellent approximation algorithms exist for finding such a partition. The resulting parameterized network, called a Harary Network, can support a wide spectrum of database assumptions and behaviors. Finally, to analyze the dynamic performance of such a database, in particular the trade-offs between availability and response times, a stochastic tool is developed. This tool enables a database designer to perform an analysis of the expected system performance and fine-tune the system parameters for improved performance. As an example, we extend Thomas's majority consensus protocol to tolerate both site and link failures and analyze its performance using our stochastic tool.