Reading the Amazon Aurora Limitless Paper
In this post, we’ll walk through how Amazon Aurora Limitless works under the hood, based on the paper Aurora PostgreSQL Limitless Database: Building a Highly Scalable OLTP Database. Rather than covering every detail of the paper, I’ll focus on what I found most interesting: how Limitless implements MVCC across a distributed database. Amazon Aurora Limitless TL;DR How can we horizontally scale Aurora? Limitless answers this question by sharding data across multiple database instances. A Limitless database consists of routers and shards. Routers receive client queries, build query plans, and distribute work among shards. Each shard stores and processes a portion of the data. Conventional XID-based MVCC does not scale naturally to a distributed database because it requires global transaction information. Limitless instead implements time-based MVCC using Amazon Time Sync Service and Hybrid Logical Clocks (HLCs). Limitless can scale vertically by increasing the capacity of individual routers and shards, and horizontally by adding more of them. Routers and Shards 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 Client | v +-------------------+ | DNS Endpoint | +-------------------+ / \ v v +----------+ +----------+ | Router 1 | | Router 2 | +----------+ +----------+ | | +-------+-------+ | +------------+------------+ | | | v v v +---------+ +---------+ +---------+ | Shard 1 | | Shard 2 | | Shard 3 | +---------+ +---------+ +---------+ | | | v v v +---------+ +---------+ +---------+ | Storage | | Storage | | Storage | | Volume | | Volume | | Volume | +---------+ +---------+ +---------+ In Limitless, customers can specify a sharding key for a table. The table is divided into partitions using hash-based partitioning, and those partitions are distributed across shards. Typically, a sharded table is divided into 256 partitions. ...