After watching the re:Invent 2024 DynamoDB (DAT406) and S3 (STG302) deep dives, I learned that DynamoDB uses an in-memory data store (MemDS) that serves their request routers.

DynamoDB is a low-latency and highly available system. AWS reports single-digit-millisecond average service latency for many single-item operations, excluding client and network overhead, and I learned it uses Request hedging and Constant work as design patterns to help achieve this.

S3 also uses a pattern called Shuffle sharding (I’m not sure if DynamoDB also uses that, but it seems likely to me).

Request hedging

Request hedging is a design pattern to deal with tail latencies in distributed systems. It’s both used by DynamoDB and S3.

The main idea is to make a second request, and use the fastest response.

This works when the requests take sufficiently independent paths. If the first request is in the tail, the second request may still complete sooner.

Request hedging is a bet. It will not always pay off. But in practice, it’s very effective. Dean and Barroso describe this pattern in The Tail at Scale.

Note

The second request can start immediately or after a delay. Hedging adds work and load, so only use it for operations that are safe to duplicate, and cancel the slower request when possible.

Shuffle sharding

Shuffle sharding is a design pattern used by S3 to allocate workloads across subsets of resources with limited overlap. This isolates contention and “hot” workloads, by keeping their blast radius bounded.

Shuffle sharding can make request hedging more effective when it helps send the hedged requests through sufficiently independent paths. It does not guarantee independence.

But shuffle sharding also helps achieve high availability:

Power of two random choices

Assigning requests completely at random can still leave some resources much busier than others. A surprisingly effective improvement is to sample two resources and pick the less busy one.

Constant work

DynamoDB uses a design pattern called constant work to make sure a (part of a) system always runs in the same “steady state”. This prevents potentially overloading the system, and helps make it highly available.

This is especially important for systems that use (in-memory) caches.

DynamoDB’s request routers have their own in-memory cache, but even for cache hits, they send request(s) to their storage (MemDS), so that when they lose the in-memory cache in the router for whatever reason, the storage “doesn’t notice”.