Welcome the fourth post in my Builder’s Library notes series containing my takeaways on the remaining pieces released prior to ReInvent 2020. These fall less neatly into a theme but contain interesting design architectures and considerations suggested by Amazon. One piece goes into caching strategies, another discusses tips of leader election and how it’s useful, a third discusses a design architecture that puts a smaller service in control of a larger one to protect against overload. The fourth piece talks about shuffle sharding, a way of placing workloads on resources in a way that makes each resilient. The last piece discusses fallback systems and why Amazon avoids them, in a first piece discusses a strategy they don’t recommend.
Caching challenges and strategies by Matt Brinkley and Jas Chhabra
Caches are appealing for better performance without expensive database changes, but a surge in uncached queries can cause servers to overload and cause an outage. This happens when a service is addicted to its cache and modal behavior based on if an object is cached can cause an issue. Amazon caches items when there are uneven request patterns leading to hot keys. Partition throttling can be good if a good cache hit ratio occurs across multiple requests and operations. Systems should account for how tolerant the system can be to eventual consistency as rate of change for the source data and cache refresh policies determine how inconsistent data tends to be. These are related as data that changes less frequently can have a longer cache time.
Two categories of caches are in memory and external caches. Caching on box can be quick and provide significant improvements with minimal work. This is commonly implemented with inprocess memory and typically the first approach tried. It requires no additional operational over head and is low risk, often using an in memory hash table managed with application logic or embedded into a service client. The downsides to in memory caches are that there can be inconsistency between servers allowing for different results based on what server is reached. Cache hits and misses should be monitored as downstream load can still overwhelm dependencies. Servers will be susceptible to cold start issues as initial requests will hit dependencies. Request coalescing can help mitigate these issues.
External caches live on a separate fleet of services and examples include Memcached and Redis. Except with update failures, these servers are likely to stay in sync and and can reduce the load on downstream service. The number of cache servers here doesn’t have to be proportional to fleet size either. As the cache is populated throughout the deployment, cold starts aren’t present and there are more options for storage. The drawbacks to external fleets are increased complexity and servers with different availability metrics and more metrics needed to track failures. Load tests with realistic traffic should be done and it’s important to know how scaling works as some services require down time when adding nodes. Other variations in different libraries include consistent caching and node discovery. These behaviors should be tested as well as ensuring that cache data serialization doesn’t cause issues when rolling forward and back. Falling back to a dependency has a risk to overwhelm the system and it’s better to fail to an in-memory cache. Load shedding and request caps can help avoid overwhelming the dependency.
Two types of caches include inline and side caches. Inline are read through, write through caches that embed a cache into the main data access API such as DAX for Dynamo or NGINX http caching implementations. Side caches are generic object stores such as Elasticache where the app code manipulates the cache itself and checks for objects before calling dependencies. Inline cache requires less work but also gives the client less control and a cache failure could appear potentially as a service failure.
Correctly choosing a cache size, expiration time and eviction policy can be additional challenge when configuring a cache. An expiration policy determines how long to keep an item in the cache and commonly associates time to live with with an object based on client’s requirements and tolerance to stale data. The eviction policy tells the cache how to remove objects once its limit is hit, with “least recently used” or LRU being a popular policy. Metrics and testing should be used to determine these numbers. Two “time to lives” can be used in case the downstream is unreachable, with a soft TTL used to refresh the cache and a hard TTL in case of refresh delay. This pattern is used by IAM and the downstream service can use backpressure to tell the cache to use old data until the hard TTL to help recover from overload. A negative cache can also be used to help cache errors as well. Security concerns with caching can be encrypting data and avoiding cache poisoning attacks. Response time can tell a user what data has been accessed before. Request coalescing is the practice of limiting a server or cache to a single request to downstream to avoid a sudden burst from a “thundering herd.”
Best practices for using caches are to ensure its needed for availability improvements and that the data is used by multiple requests. Caching and scaling behavior should be monitored and tested with alarms set on CPU, memory and hit rate. Ensure that cache works and scaling doesn’t cause down time. A load test with the cache disabled can help ensure that the dependencies can stay up on cache failure. Design the storage format to evolve and be able to change versions and consider security, encryption and the possibility of poisoning and side channel attacks.
For the full article these notes are taken on, check out https://aws.amazon.com/builders-library/caching-challenges-and-strategies/
Leader election in distributed systems by Marc Brooker
Leader election is giving one piece of a distributed system special powers, such as modifying data, assigning work or handling all requests. It can improve efficiency and reduce required coordination, simplifying architecture but also adding failure modes and scaling bottlenecks. Generally another option should be considered first such as workflow services such as step functions or using idempotent apis and other patterns to allow retries. Some advantages of leader election are that a single leader with all concurrency in one place is easy to think about. Logs can be found in one place and partial failure is reduced. Other servers can be informed of changes instead of waiting for consensus and only a single cache is required. Also there’s no risk of other servers doing the same work.
Disadvantages are that there’s a single point of trust (with a large blast radius), failure and scaling and a rearchitecture will be needed to grow beyond one leader. Partial deployments such as One-box or A/B testing can be difficult as well. Many of the drawbacks can be mitigated by choosing a leader’s scope such as sharding with each shard having its own leader. Dynamo, EFS, EBS and more AWS services use this approach. Sharding adds design complexity of own and needs to be carefully considered.
When it comes to electing a leader, there are many ways including Zookeeper, the Paxos algorithm and custom hardware. Amazon typically uses leases where a database stores the current leader and listens for heartbeats. If no heartbeat in period of time, another node can try and take over. Leases typically depend on elapsed time duration and Dynamo’s lock client is a good example of an implementation used at Amazon. It’s hard to make sure the leader only does work while holding a lock and doesn’t believe its holding a lock longer than it is. Garbage collection pauses can create an issue and hardening against lock issues can be a major challenge. DynamoDB and Zookeeper both have lease based locking clients with fault tolerant leader election and Amazon typically reuses these unless a specific need for something else.
Examples of systems that use leader elections are systems with relational database management systems, which handles all writes. Election is often done by a human operator. EBS distributes read and write across servers but has a primary in each area of the volume to order the actions. If it fails, followers can use the same election mechanism. This improves consistency and avoids coordinating with the data plane and similar approaches are used by DynamoDB, QLDB and Kinesis. The Kinesis client library uses leases to ensure each shard is processed by its owner making scaleout processing easier.
When a leader fails, its important that tasks are idempotent so that the next leader can redrive them. Leadership can pass from server to server and its not possible to guarantee that there’s a single leader, which behavior depending on whether a system can handle two leaders. Potential for two leaders can help with high availability and weaker election processes with it being harder to build a system that can’t tolerate that. In systems that allow a maximum of only one leader, the leader election system must be correct and make sure that the previous leader is removed first.
Best practices are to check remaining lease time frequently and avoid heartbeating in background threads. Metrics should show how much work a leader does compared to its potential and scaling should be used before needed. An audit trail can help tell who the leader is at any moment and algorithms should be modeled and verified for correctness with tools such as TLA+ to catch hard to find bugs.
For the full article these notes are taken on, check out https://aws.amazon.com/builders-library/leader-election-in-distributed-systems/
Avoiding overload in distributed systems by putting the smaller service in control by Joe Magerramov
This article talks more about control plane and data plane interaction as well as how to manage different services with limited responsibilities within distributed systems. Services can interact over apis and scale or change independently such as with data and control planes. Control and data planes’ interactions should be set up to avoid overload. Often times the larger data plane will call the smaller control plane, but there are benefits to putting the smaller service in control.
The two planes need to stay in sync as the data plane will need configuration updates while the control plane needs to know the state of the data plane. Data plane often has a fleet of 100 times as many instances so choosing the direction of calls can be important.
One benefit of the data plane calling the control plane’s load balancer to request updates is that the control plane doesn’t have to keep track of all data plane instances. However, a big drawback with the larger service calling the smaller is that too many calls at once can lead to overload. This can be caused by a code change, outage recovery or retries causing all servers to check for updates at once. When requests are uncorrelated, things tend to be fine, but when clients act in correlated matter, issues arise. Careful tuning is required to avoid this breaking point, such as load shedding on the smaller service and the data plane can use backoff and jitter to help avoid too many requests at once.
Scale mismatch is the biggest architecture challenge here and a popular architecture at Amazon uses S3 to help with it. Data plane apis aren’t directly exposed. The control plane writes config updates to a bucket, the data plane polls the bucket to get the update and caches it locally. The data plane then can write its state periodically to a bucket for the control plane to poll. S3 can help reduce load in both directions to remove the issue of size mismatch and also stops the outage of one plane from affecting the other as both read form S3, providing static stability with data in buckets. Hyperplane, an internal network function virtualization system behind Network Load Balancer, NAT Gateway and Private Link, uses a similar architecture. Hyperplane scans dynamo tables with customer configurations and writes them to S3 files, which the data plane downloads to update its internal routing configuration. Using S3 as an in between service doesn’t work for all scenarios. If the configuration is too dynamic, constantly writing may not be practical (EC2 needs constantly changing configuration for IAM, VPCs, EBS, etc.). It’s also a problem if a change has to be reflected in ten seconds or less such as when starting up a new container.
When using S3 doesn’t work, putting the smaller fleet in control can help control the pace of requests. The control plane can push the config change to the data plane fleet which doesn’t have the same worry of overload. In this case, the control plane will need to have up to date inventory of the data plane servers and be able to handle issues if they become unreachable, making sure all servers are accounted for by coordinating with other control plane servers. One way of handling is to have a data store each control plane server reads from to know its assigned data plane servers it’s responsible for based on partitions as well as which servers sent recent heartbeats. A potential compromise is to have the data plane initiate a connection via the control plane’s load balancer but let the control plane server make calls across it or reject if busy. This requires a more sophisticated communications protocol but lets control plane control ace and know which data plane servers want an update. Enough connection requests could still potentially overload if there’s a large enough mismatch (1000 to 1), but overall, this can be a good compromise and a variation is used for EC2’s nitro system to get updates.
For the full article these notes are taken on, check out https://aws.amazon.com/builders-library/avoiding-overload-in-distributed-systems-by-putting-the-smaller-service-in-control/
Workload isolation using shuffle-sharding by Colm MacCarthaigh
Probably the first piece to catch my interest, this describes a system for mapping customers’ workloads to nodes in a way that avoids an attack on or failing of one customer from being able to take down another. Having a requirement to host a customer’s domain while working with Route53 to return changing ip addresses based on specific AWS services at the route of a domain, they’d also have to take on defending the domain when it’s a lot cheaper to scale fake clients than it is to add more service and the risk of one client being taken down by an attack on another. The goal for this approach was to target resources only at the domain being attacked.
A mentioned downside to scaling is that if any worker can handle any request for any client/domain, then it’s easy for a bug or excess requests to shut down the entire system. Sharding helps make the system more resilient by limiting requests for a client to specific nodes, meaning that only the nodes in that shard would be taken down. With lots of clients, the goal of shuffle sharding is to assign each customer to its own shard of 2 or more workers and use an efficient algorithm to assign customers to shards so that no customers share all their shards. More isolation is possible when with more shards/workers. A given example is that Route53 uses 2048 virtual name servers and assigns each customer’s domain to a shard with 4 to ensure that no customers share more than 2 servers with each other.
Shuffle-sharding is described as a great tool for both lowering overall risk and helping to counter the risk still present. In the event of a DDOS attack, the worst case scenario would be one customer taken down, but this is less likely as Amazon can focus all their effort on defending specifically the nodes being attacked since the attack surface is much smaller. Shuffle-sharding is referenced in many other pieces in the Builder’s Library as well and Colm’s paper goes into a lot more detail and contains diagrams for visualizing the process.
For the full article these notes are taken on, check out https://aws.amazon.com/builders-library/workload-isolation-using-shuffle-sharding/
Avoiding fallback in distributed systems by Jacob Gabrielson
This is a paper about the idea of having fallback systems and why it’s not the best. Amazon’s services have to handle a majority of critical failures to be reliable. Four categories of strategies to handle failure, retry, proactive retry (run multiple times, use first success), failover (perform action on different endpoints) and fallback, which is using a different mechanism to get the same result. Fallback is rarely used at Amazon even if backup plans are common in real life.
In the case of a single machine, a system would try one function to get a result and a second function if the first fails. This can be hard to test as tests would have to be written for both functions and failure would have to be simulated to make sure failover works. There’s also the problem that if the second function was better than the first, it would be the primary function so you’re always committing to using a worse function. In that case why not spend the resources improving the generally used primary. Calling a backup function can also hurt if the error was caused by a lack of memory. It can also change what kind of resources are used, leading to unpredictable load (ex. CPU to I/O) and if not used can cause a latent bug that doesn’t get found for a while which can be even worse in distributed systems. The paper walks through an example of a memory allocation function with a single machine fallback.
While on a single machine, a fallback isn’t ideal, it’s much worse in a distributed system. The customer and service don’t have a shared fate and it can be much harder to test as many more cases exist and overload scenarios have to be mitigated against. Fallback can fail even if the odds of success are improved and can end up making an outage worse if all traffic is hitting worse fallback function on failure. Why use a worse fallback function if something is already wrong? Latent bugs occur with unlikely coincidences. An example of this was implementing a cache that fell back to a database query. When all caches failed at once, all traffic went to the database which locked up and the site went down. The database was also used by fulfilment centers causing a worldwide halt until a fix was made. In the example, having no fallback would have been better. All problems in single machine have worse consequences in distributed systems.
Amazon’s techniques to avoid using fallback are to improve the primary service, let the caller handle the error via retry, push data proactively, or potentially converting to failover with sufficient testing if needed. It’s important to make sure that timeouts and retries don’t fall into the same traps and are thoroughly tested.
For the full article these notes are taken on, check out https://aws.amazon.com/builders-library/avoiding-fallback-in-distributed-systems/





Leave a Reply