Managing Scalable Persistent Data

Francesc D. Mu, J. R. Garc · 2011

Distributed applications should be able to manage dynamic workloads; i.e., the amount of client requests per time unit may vary frequently and servers should rapidly adapt their computing efforts to those workloads. This implies that these applications should be able to scale out without problems in order to handle workload peaks and to reduce their number of replicas when the workload diminishes, at least when a pay-per-use utility model is assumed. Cloud systems provide a solid basis for this kind of applications. This paper surveys different techniques being used in different modern systems in order to increase the scalability and adaptability in the management of persistent data. Those techniques follow two basic principles: (a) to minimise distributed coordination, and (b) to eliminate any sources of delay in local operation service. These principles are implemented following six complementary mechanisms: (1) to replicate data in order to improve read access parallelisation, (2) to partition the database in order to increase update access concurrency, (3) to relax the resulting replica consistency in order to ensure network-partition tolerance, (4) usage of simple operations in order to reduce concurrency control efforts, (5) declaration of simple schemas, thus reducing the dependency on elaborate indexing techniques and eliminating the need of join operations, and (6) to bound coordination for directly achieving the first

Read the paper · More papers on PaperTik