Reliable Stateful Systems at Netflix
There is quite an interesting long-read article explaining how Netflix implements the reliability of its stateful services.
To understand what reliability actually means the author proposes answering three fundamental questions:
- How often does the system fail?
- When it fails, how large is the blast radius?
- How long does it take to recover from an outage?
Ideally, systems don’t fail, have minimal impact on failure and recover very quickly😀. That is the simple part.
So how is that achieved? That is actually the hard part:
📍Single tenancy. There are no multi-tenant data stores. It minimizes blast impact.
📍Capacity management. Netflix programs special workload capacity models that generate specifications for a cluster according to the system requirements and SLOs.
📍Data replication. Data is replicated in 12 availability zones across 4 regions.
📍Overprovisioning. In case of one region degradation, the traffic is spread across other regions so each region must keep an extra 33% capacity reserved for failover.
📍Snapshot restoration. To replace one instance with another - load snapshot from s3 and then apply delta.
📍Performance monitoring. Continuous monitoring is essential to detect failures quickly, remediate them and recover.
📍Cache in front of services. The idea is to use caching for complex business logic rather than underlying data. Service with business logic is quite expensive to operate, so approach shifts load to the cache.
📍Reliable clients. Quite complex approach to inform clients what timeouts and level of service that they can rely on. Better to read the original.
📍Load Balancing. Netflix uses improved choice-of-2 algorithms with weight requests. It takes into account availability zones, replicas, health state.
📍Stateful APIs. All requests are idempotent with built-in work pagination. That approach requires additional idempotent tokens or global transaction ids.
#architecture #reliability #usecase
Post #8
261