> For the complete documentation index, see [llms.txt](https://jaywin.gitbook.io/leetcode/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://jaywin.gitbook.io/leetcode/system-design/database-and-file-system/dynamodb.md).

# DynamoDB

Paper: <https://www.allthingsdistributed.com/files/amazon-dynamo-sosp2007.pdf>

Dynamo is distributed key-value store, i.e. Distributed Hash Table (DHT) , DynamoDB is AWS database product based on Dynamo.

## Why DynamoDB

* highly available (AP model, whereas Bigtable is CP focus on consistency)
* highly scalable
* schema-less
* why not
  * if requires strong consistency

## Architecture

* Data distribution: consisten hashing
* Replication: optimistic, eventual consistency
* Handling failure: sloppy quorum(hinted handoff)
* Inter-node communication: goosip protocol
* Conflict resolution: vector lock

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-98120df87752fe50e840c57ff9a5227e1324220d%2Fimage.png?alt=media)

## Data model

* Dynamo is key-value store. Value is a set of attributes.
* Partition key: input for hash function(MD5) -> partition(physical storage)
* Partition key + sort key
  * composite primary key
  * items with the same partition key value are stored together, in sorted order by sort key value.
* Secondary index
  * global
  * local
* Query pattern
  * Range query is not well supported, unlike Bigtable scan. E.g. `select * from A where age > 20`
    * Since data is distributed in many nodes, this is basically a full scan or the cluster.
  * Solutions (either way, we may still end up doing full scan)
    * use range key(aka. sort key)
    * create global secondary index ([reference](https://aws.amazon.com/blogs/big-data/scaling-writes-on-amazon-dynamodb-tables-with-global-secondary-indexes/))
* Local persistent storage
  * These (key, value) pairs are stored within that node using various storage systems depending on application needs. A few examples of such storage systems are:
    * BerkeleyDB Transactional Data Store
    * MySQL (for large objects)
    * An in-memory buffer (for best performance) backed by persistent storage

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-298554e32f8b1d783d6a559b936e339217b6fa6f%2Fimage.png?alt=media)

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-055c8ace259341373fd0e338d3bf377b9a7d72e1%2Fimage.png?alt=media)

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-afbcee408a7071bb862f6dce0f98c9806fdc42b1%2Fimage.png?alt=media)

## Implementation details

### Optimistic replication

Coordinator node stores data first, then asynchronously replicates to next N nodes in the background. Replicas are not guaranteed to be identical at all times.

If a client cannot contact the coordinator node, it sends the request to a node holding a replica.

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-e816bcb125a75087d3959b0af4b272242378c017%2Fimage.png?alt=media)

### Add & remove nodes

* Add nodes, aka. bootstrapping.
  * Define how many VNodes the new machine is responsible
  * Allocation algorithm in the new node pick random VNodes(reprented as tokens).
  * New node requests current replicas of those tokens to stream data.
* Remove nodes.
  * Send command to remove a node.
  * The old node reassigns tokens to other nodes. Also replicates data to those nodes.
* Reference: <https://cassandra.apache.org/doc/latest/cassandra/operating/topo_changes.html>

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-8679a24467f844c36c1df20cf4d0b2002fe18d47%2Fimage.png?alt=media)

### Sloppy quorum

* Preference list
  * The list of nodes responsible for storing a particular key. E.g. for key "K", preference list is \[server 1, 2, 3, 4] where 4 is to store hinted replica.
* Sloppy quorum ensures "always writable"
  * When a node is unreachable, another node can accept writes on its behalf. The write is then kept in a local buffer and sent out once the destination node is reachable again.
* Drawbacks: data conflict

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-c99a243b42bf095d1fefe26e45acf4d7e10ea773%2Fimage.png?alt=media)

### Vector clock

* Clock skew
  * In distributed system, we cannot assume that wall clock time t on node a happened before time t + 1 on node b.
  * Hardware solution: GPS unit, atomic clock.
* Dynamo returns conflicting versions to **client who are responsible to explicitly reconcile** the conflict.
* Alternative way to handle conflict
  * Model data as Conflict-free replicated data types (CRDTs). E.g. adding item A, B to shopping cart can be handled anywhere anytime, because the end result is 2 items in the cart anyway.
    * AKA strong eventual consistency
    * Downsides: not easy to model every data as CRDT
  * Last-write-wins (LWW).
    * Downsides
      * using wall clock doesn't guarantee a real last write.
      * can easily lost data, e.g. for 2 concurrent data, 1 of them is thrown away.

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-0544d604141753aead36f796a05a01f47e3a041d%2Fimage.png?alt=media)

### Put & Get operation

* Choosing coordinator node
  * Via load balancer
    * Pros: system loosely coupled, helps scalability
    * Cons: extra hop increase latency, waste of resources
  * Dynamo use client library that maintains server address
    * Pros: zero-hop DHT, low latency
    * Cons: less control of load distribution
* Quorum
  * Read + Write > N(replicas)
  * A Common (N, R, WN,R,W) configuration used by Dynamo is (3, 2, 2).
    * (3, 3, 1): fast WW, slow RR, not very durable
    * (3, 1, 3): fast RR, slow WW, durable
* Write(put) process
  * coordinate node generate vector clock
  * saves data in local
  * sends write request to N-1 nodes from preference list
  * returns successful after receiving W-1W−1 confirmation
* Read(get) process
  * coordinate node request data from N-1 nodes from preference list
  * wait until R-1 replies
  * handle data versions (then update it back to nodes, aka. read repair)
  * returns all relevant data versions(may have conflicts) to the caller
* State machine
  * Each client request results in creating a state machine on the node. (like request context?)
  * Coordinate node selection
    * Any top N nodes from preference list
    * A usual request patter is read-then-write, so a node replied fastest in previous read operation is chosen to handle write. The info stored in request context. This increases the chances of getting “read-your-writes” consistency.

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-3742d1f030ccee23040805e42556b3e4efc5f73d%2Fimage.png?alt=media)

### Anti-entropy Through Merkle Trees

* Read repair: use vector clock to resolve conflicts. But it's too slow if a replica falls too behind.
* Merkle trees is used to resolve conflicts in the background.
* Workflow
  * Data is split into smaller parts. Merkle trees is a binary tree of hashes.
  * Compare 2 tree from root to leaves. If encounter difference, resolve conflicts in this small subtree. Thus data transmission and work load is minimized.
* Drawbacks: have to recalculate the tree when data changes.

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-65ae065507929aea77f181830f9f27b1e052c8b5%2Fimage.png?alt=media)

### Gossip protocol

* Peer-to-peer communication
* Each node periodically exchange info(basically the copy of hash ring) with random nodes.
* Use seed nodes to bootstrap. More read on [Uber's Ringpop](https://eng.uber.com/ringpop-open-source-nodejs-library/).

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-fab3b0519f6a31dee354cb495fc7922315bc5f21%2Fimage.png?alt=media)

### Inspired by Dynamo: Cassandra and Riak

![](https://3398971849-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2F-M9HSM8Jn9nrPTh34ajG%2Fuploads%2Fgit-blob-299dd606ffde86a39bb37fd21a437a6c17eac465%2Fimage.png?alt=media)

## Criticism on Dynamo

* Each Dynamo node contains the entire Dynamo routing table. This is likely to affect the scalability.
* Security concerns, e.g. no ACL.
* Dynamo’s design can be described as a “leaky abstraction,” where client applications are often asked to manage inconsistency.

## Reference

* <https://www.educative.io/courses/grokking-adv-system-design-intvw/xoEXr9614RB>
