Design of a Scalable Distributed Database System: SD-SQL Server

Soror Sahri · 2006

Databases are now often huge and growing at a high rate. Large tables are then typically hash or range partitioned into segments stored at different storage sites. Current Data Base Management Systems (DBSs), e.g., SQL Server, Oracle or DB2, provide static partitioning only. The database administrator (DBA) in need to spread these tables over new nodes has to manually rede redistribute the database (DB). A better solution has become urgent. This situation is similar to that of file users forty years ago in the centralized environment. The indexed Sequential Access Method (ISAM) was in use for ordered (range partitioned) files. Likewise, only static hash access methods were known. Both approaches required file reorganization whenever the file grew too large. B-trees and extensible (linear, dynamic) hash methods were invented to avoid this need for file reorganization. Instead of reorganizing a complete file, these methods deal with file growth by incremental splitting one or a few buckets (pages, leaves, segments...) at certain inserts. These dynamic methods were successful enough to make ISAM and centralized static hashing obsolete. Efficient management of distributed data adds specific needs. Scalable Distributed Data Structures (SDDSs) address these needs for files. SDDS can use hashing, range-partitioning or k-d trees to distribute its data in buckets spread over the nodes of a multicomputer. These nodes can form a peer-to-peer (P2P) or grid network. An SDDS grows to more buckets by splitting the overflowing ones. The splits are triggered by the (overflowing) inserts.

Read the paper · More papers on PaperTik