0:00
How does Amazon maintain high
0:02
availability during Black Friday traffic
0:04
spikes? Why do financial systems rarely
0:07
lose transactions despite network
0:09
failures? Today, we'll explore the
0:11
reliability and consistency concepts
0:14
that power high-scale distributed
0:16
system. Let's dive right in. System
0:19
reliability isn't about preventing all
0:21
failures. It's about continuing to work
0:23
correctly when failures inevitably
0:26
happen. Network partitions will occur.
0:29
Servers will crash, hard drives will
0:31
fail. The question is, how do we build
0:33
systems that handle these realities
0:35
gracefully? There are seven key concepts
0:38
that make this possible. Let's examine
0:40
each one. The CAP theorems establishes a
0:43
fundamental constraint. Distributed
0:45
system can provide at most two of these
0:47
three guarantees simultaneously.
0:50
Consistency. All notes see the same data
0:53
at the same time. Availability. Every
0:56
request to a non-failing note receives a
0:58
response. Partition tolerance. The
1:01
system continues operating despite
1:03
network failures. Since network
1:05
partitions are unavoidable in
1:07
distributed systems, you're really
1:09
choosing between consistency and
1:11
availability during partition events.
1:13
Let's look at how real systems make this
1:16
choice. Google Spanner chooses
1:18
consistency. It uses atomic clocks and
1:20
synchronized time across global data
1:23
centers to maintain linearizable
1:25
transactions. During network partitions,
1:28
the majority partition remains available
1:30
for both reads and rights while minority
1:32
partitions become read only. The
1:35
trade-off is that notes in the minority
1:37
partition lose right availability until
1:39
the partition heals. Amazon Dynamo DB
1:42
chooses availability. During network
1:45
partitions, it continues accepting
1:46
rights using eventual consistency and
1:49
resolve conflicts to last right wins
1:51
based on timestamps. User always get
1:54
responses, but they might occasionally
1:56
see stale data. Neither approach is
1:58
right or wrong. It depends on your
2:00
requirements. Banking systems typically
2:03
need consistency. Social media feeds can
2:06
tolerate eventual consistency.
2:09
Eventual consistency makes a simple
2:10
promise. If you stop making updates, all
2:13
replicas will eventually converge to the
2:15
same state. This might sound sloppy, but
2:18
it's incredibly powerful for large scale
2:20
systems. Here's why. In eventually
2:23
consistent system, rice can complete
2:25
immediately without waiting for
2:27
confirmation from all replicas. This
2:29
improves performance and availability.
2:32
Amazon's shopping cart works this way.
2:34
You can add items even if some servers
2:36
are temporarily unreachable. But how the
2:38
system handle conflicts when different
2:40
replicas receive different updates? When
2:43
conflicts occur, systems use different
2:45
resolution strategies. Last right winds
2:48
use timestamp to pick the most recent
2:50
update. Simple but can lose data.
2:52
Conflict free replicated data types use
2:54
mathematical properties to guarantee
2:57
that all replicas converge to the same
2:59
state regardless of update order.
3:01
Application defined merge functions let
3:04
developers write custom logic for
3:06
resolving conflicts based on business
3:08
rules. Modern systems like Dynamob
3:11
typically achieve consistency within
3:13
milliseconds under normal conditions
3:15
making an eventual park barely
3:17
noticeable to users.
3:20
Low balances distribute incoming
3:22
requests across multiple servers. Sounds
3:24
simple, but there's sophisticated
3:26
engineering behind effective load
3:28
balancing. There are two main types.
3:30
Layer 4 load balancers operate at the
3:32
transport layer. They make routing
3:34
decisions based on IP addresses and TCP
3:37
UDP ports. They are fast because they
3:40
don't need to inspect application data,
3:42
but they are also limited in their
3:44
routing intelligence. Layer 7 load
3:46
balancers operate at the application
3:48
layer. They can examine HTTP headers,
3:51
URLs, and even request content to make
3:54
smarter routing decisions. More
3:56
powerful, but also more computationally
3:59
expensive.
4:00
Modern load balancers go beyond simple
4:02
roundrobin distribution. Lease
4:05
connections routing sends new requests
4:07
to the server handling the fewest active
4:09
connections. Lease time factors in how
4:12
quickly each server responds avoiding
4:14
slower servers. For high availability,
4:17
low balances themselves are deployed in
4:19
primary secondary configurations with
4:22
heartbeat protocols. If the primary
4:24
fails, the secondary takes over within
4:26
milliseconds.
4:29
Consistent hashing ensures the same
4:31
client consistently hits the same
4:33
server, which is critical for
4:35
maintaining session state. In the last
4:37
video, we talked about horizontal
4:39
scaling. Here's a problem that large
4:41
horizontally scale systems face. You
4:44
have data spread across multiple nodes
4:46
and you need to add or remove nodes
4:48
without moving massive amounts of data.
4:51
Traditional hashing doesn't work well
4:52
here. If you use simple modular hashing,
4:55
then adding one node requires remapping
4:58
almost all of your data. That's
5:00
expensive and disruptive. Consistent
5:02
hashing solves this elegantly. Instead
5:04
of mapping keys directly to notes, both
5:07
keys and notes are placed on a circular
5:09
hash ring. Here's how it works step by
5:12
step. First, hash each note to determine
5:15
its position on the ring. Second, to
5:17
find where data belongs, hash the key
5:19
and walk clockwise around the ring until
5:21
you hit the first note. Third, replicate
5:24
the data to the next n minus one notes
5:26
clockwise on the ring. When you add a
5:29
new note, it only takes ownership of
5:31
keys from its immediate neighbors. When
5:33
you remove a note, it keys get
5:35
redistributed to the next note on the
5:37
ring. The result adding or moving notes
5:40
only requires moving k over n keys
5:43
instead of nearly all keys where k is
5:46
total keys and n is the number of nodes.
5:49
Amazon Dynamo DB and Apache Cassandra
5:52
both use consistent hashing for exactly
5:54
this reason. It lets them scale
5:56
horizontally without massive data
5:58
reorganization.
6:01
When one service in a distributed system
6:03
fails, it can trigger a cascade of
6:06
failures across dependent services.
6:08
Circuit breakers prevent this. A circuit
6:10
breaker monitors the failure rate array
6:12
of calls to the dependent service. It
6:15
has three states. Closed requests flow
6:18
normally. Open requests are block
6:21
immediately returning fast failures.
6:23
Half open. A few test requests allow
6:26
through to check if the service has
6:28
recovered. Here's how it works step by
6:30
step. The circuit breakers tracks
6:32
successful and failed requests to a
6:35
service. When the failure rate exceeds a
6:37
threshold, the circuit breaker trips to
6:40
open state. In the open state, requests
6:42
fail immediately without even attempting
6:45
to call the failing service. This
6:46
prevents resource exhaustion and gives
6:49
the failing service time to recover.
6:51
After a timeout period, the circuit
6:53
breaker moves to half open and allows a
6:56
few test requests through. If the test
6:58
requests succeed, the circuit breaker
7:00
closes and normal operation resumes. If
7:03
they fail, it opens again. Netflix
7:06
pioneered this pattern with their
7:07
history library and is now standard in
7:10
microservices architectures. The key
7:13
insight is that fast failures are better
7:15
than slow failures that consume
7:17
resources.
7:19
Rate limiting controls how many requests
7:21
client can make within a given time
7:22
window. It protects systems from
7:25
overload and abuse. There are several
7:27
algorithms each with different
7:29
characteristics. Token bucket
7:31
accumulates tokens at a fixed rate up to
7:34
a maximum capacity. Each request
7:36
consumes a token. This allows control
7:39
burst while maintaining an average rate.
7:41
Leaky bucket processes requests at a
7:43
constant rate smoothing out traffic
7:46
spikes. Excess requests are cute or
7:49
dropped. Fixed window can request in
7:51
discrete time intervals. simple but
7:54
prone to boundary effects where clients
7:56
can double the rate by timing requests
7:59
around window boundaries. Sliding
8:01
windows uses a rolling average that
8:04
smooths out the boundary problems of
8:06
fixed windows. In distributed systems,
8:09
ray limiting gets more complex. You
8:11
can't just count requests on a single
8:13
server. You need coordination across
8:15
multiple instances. One solution is rad
8:18
space ray limiters that use l script to
8:21
atomically increment counters and set
8:23
ttls. This ensures consistent ray
8:26
limiting across your entire service
8:28
cluster. Many systems implement tier
8:31
array limiting with different limits for
8:33
authenticators versus anonymous users
8:36
and progressively stricter limits as
8:38
suspicious patterns are detected.
8:41
You can manage what you can measure.
8:43
Monitoring provides visibility into
8:45
system behavior and performance. Modern
8:48
observability focuses on four key signal
8:51
types. Metrics, time series numerical
8:53
data like CPU usage, request rate, and
8:57
error counts. Logs structured records of
9:00
discrete events with contextual
9:02
information. Traces end to end request
9:05
flows showing how a single request moves
9:08
through your distributed system. events,
9:11
significant occurrences like deployments
9:13
or configuration changes. The challenge
9:16
is processing massive amounts of
9:18
telemetry data. Large scale systems
9:20
generate terabytes of monitoring data
9:22
daily. Effective alerting balances two
9:25
concerns. You want to catch real
9:27
problems quickly, but you don't want to
9:29
be overwhelmed by false alarms. Static
9:32
thresholds work for stable metrics, but
9:34
fail when traffic patterns change.
9:36
Modern systems use statistical anomaly
9:39
detection that learns normal patterns
9:41
and alerts on deviations. Composite
9:44
alerts combine multiple signals to
9:46
reduce noise. Instead of alerting on
9:48
high CPU alone, you may alert when CPU
9:51
is high and error rates are increasing
9:54
and response times are slow. The goal is
9:56
service level objectives that measures
9:59
user experience.
10:01
Not every systems needs all these
10:03
patterns. A simple web application
10:05
serving a few thousand users doesn't
10:07
need consistent hashing or circuit
10:09
breakers. The concepts we covered give
10:11
you the tools to make this trade-off
10:13
consciously rather than discovering your
10:15
systems limitation during an outage.
10:18
Start simple, measure everything and add
10:21
complexity only when you have clear
10:23
evidence you need it. That's how you
10:25
build systems that scale reliably.
10:28
If you like our videos, you might like
10:30
our system design newsletter as well. It
10:32
covers topics and trends in large scale
10:34
system design. Trusted by 1 million
10:37
readers. Subscribed at blog.by go.com.