A design framework and a scalable storage platform to simplify internet service construction

Steven D. Gribble, Eric Brewer · 2000

The Internet infrastructure has evolved from a collection of loosely organized data repositories and web pages to a rich landscape populated with industrial strength applications and services. Although this evolution has been rapid, the task of building and maintaining these services nonetheless remains challenging, primarily because services must be exceptionally robust, remaining available and performing well in the face of voluminous and growing traffic demands. Coupled with a lack of suitable reusable building blocks and design methodologies for service construction, this challenge unfortunately implies that only organizations with very capable engineering and operations staff can currently successfully build and maintain new Internet services. This dissertation represents a step towards ameliorating this situation; in it, we address two sets of challenges: the design and implementation of a programming model, concurrency model, and I/O substrate specifically geared towards Internet service construction, and the design and implementation of a storage platform that shields service authors from the complexities of robust, scalable persistent data management. The first half of this dissertation focuses on the development of a programming model and design framework that is well suited to the needs of scalable, highly concurrent services. The framework consists of a set of design patterns that can be applied to code in order to “condition” it against load, concurrency, failure, and performance bottlenecks. The framework also describes a way of structuring programs using queues to separate and decouple the program's stages. The second half of this dissertation describes a scalable storage platform that we built using our design framework. This platform, called a distributed data structure (DDS), greatly simplifies the task of implementing a new Internet service by completely shielding authors from the complexities of scalable, available, consistent storage management. We describe the design, implementation, and performance of a distributed hash table, as well as a number of services implemented using it. The hash table design makes several assumptions and optimizations based on the properties of clusters of workstations. We believe that this dissertation makes several contributions that greatly reduce the complexity of implementing new Internet services. As such, we hope that it will accelerate the evolution of the Internet by encouraging more people to implement and deploy new, innovative, creative services.

Read the paper · More papers on PaperTik