0:00
Today we're looking at key value stores,
0:03
a type of database that keeps track of
0:05
everything from your shopping cart to
0:06
your chat messages. And what starts as a
0:09
simple idea quickly becomes a hot
0:12
problem in software engineering. So what
0:14
is a key value store? Think of it as a
0:17
giant dictionary. You have a key, maybe
0:20
user 1 2 3 4 5 cart and a value. All the
0:24
items in that user's shopping cart. You
0:27
can put data in, you can get data out.
0:29
Here's the thing, though. Amazon stores
0:32
everything. We're talking about
0:33
terabytes of data per major region with
0:36
billions of key value pairs that need to
0:38
be accessed millions of times per
0:40
second. No single computer can handle
0:43
that load. So, you need to spread it
0:45
across thousands of servers. But now you
0:48
have a problem. When someone asks for
0:50
user 1 2 3 4 5 cart, how do you know
0:53
which server has it? Your first instinct
0:55
might be just hash the key and use
0:57
modulo. Hash the key, divide by the
1:00
number of servers, use the remainders to
1:03
pick a server. Simple math. But here's
1:05
where it gets tricky. What happens when
1:07
you need to add a new server? Suddenly,
1:10
you're dividing by a different number
1:12
and almost every key maps to a different
1:14
server. You would have to move nearly
1:16
all your data just to add one machine.
1:19
This is where consistent hashing comes
1:20
in, and it's pretty clever. Instead of
1:23
mapping keys directly to servers,
1:25
imagine both keys and servers living on
1:27
a giant circle. Think of a clock face,
1:30
but instead of 12 hours, you have
1:32
millions of positions. You place your
1:34
servers at random spots around this
1:36
circle. Maybe server A is at position
1:39
100. Server B is at position 1,000.
1:42
Server C is at position 5,000. Now, when
1:45
you want to store user 1 2 3 4 5 cart,
1:48
you hash that key to get a position on
1:50
that circle. Let's say position 750. You
1:53
site at that spot and walk clockwise
1:56
until you hit the first server. In this
1:58
case, that would be server B at position
2:01
a th00and. The magic happens when you
2:03
add a new server. Say you place server D
2:06
at position 500. Now keys that hash to
2:09
position 101 through 500 go to server D
2:12
instead of server B. But everything else
2:15
stay exactly where it was. You only move
2:17
a fraction of your data. But wait, what
2:20
happens when server B crashes? Now all
2:22
those keys has nowhere to go. This is
2:25
where you need copies. Instead of
2:27
storing each piece of data on just one
2:29
server, you store it on multiple
2:31
servers. One way to do this is to keep
2:33
using that circle. When user 1 2 3 4 5
2:37
card hashes to position 750, you don't
2:40
just store it on server B. You also
2:42
store copies on the next two servers
2:44
clockwise. Maybe server C and server A.
2:48
Now, if server B goes down, you still
2:50
have the data. Great. Now, your data is
2:52
safe, but you've created a much bigger
2:55
headache. Let's say two people are
2:57
shopping for the same family account.
2:59
They both add items to the cart at the
3:01
exact same time, but their requests hit
3:03
different servers. Now, you have two
3:05
different versions of the same shopping
3:07
cart. Which one is correct? This brings
3:09
us to one of the fundamental concepts of
3:12
distributed systems. You cannot have
3:14
perfect consistency, perfect
3:16
availability, and perfect network
3:18
reliability all at the same time. You
3:20
have to pick two. If you choose
3:22
consistency, making sure everyone always
3:25
sees the latest data. You might have to
3:27
refuse requests when server can't
3:29
communicate. Banks do this because
3:31
showing the wrong account balance is
3:33
unacceptable. If you choose
3:35
availability, keeping the system running
3:37
no matter what, you might occasionally
3:39
serve stale data. Most web apps do this
3:42
because it's better to show an old
3:44
shopping cart than no shopping cart at
3:46
all. The solution most big systems use
3:49
is called eventual consistency. The idea
3:52
is simple. Given enough time, all the
3:55
copies will match up. But right now,
3:57
this instance, they might be a little
3:59
different. This creates a new challenge.
4:02
How do you handle conflicting versions?
4:04
There are different approaches to this
4:06
problem. One clever solution is called
4:08
vector clocks. Think of it like a
4:10
version number but smarter. Every time a
4:12
piece of data gets modified, it gets
4:15
tagged with information about which
4:17
server did the modification and when.
4:19
When you detected two conflicting
4:21
versions, you have a few options.
4:23
Sometimes you can automatically merge
4:25
them like combining two shopping carts
4:27
to include all the items. Sometimes you
4:30
have to ask the users to choose.
4:32
Sometimes you just pick the most recent
4:34
one and hope for the best. Next problem.
4:36
In a system with thousands of servers,
4:39
machines are failing constantly. How do
4:41
you even know when a server is down? The
4:44
naive approach is to have every server
4:46
ping every other server. But with
4:48
thousand servers, that's nearly a
4:50
million connections. It doesn't scale.
4:53
One solution that works well is called a
4:55
gossip protocol. Each server keeps a
4:57
list of all the other servers and
4:59
occasionally shares that list with a few
5:01
random neighbors. If server X stops
5:04
responding, the gossip spreads
5:06
throughout the entire cluster. It's
5:08
exactly like how rumor spreads in high
5:10
school. Surprisingly effective, and you
5:12
don't need everyone talking to everyone.
5:14
What started as remember this shopping
5:16
cart turns into a master class in
5:18
distributed systems engineering. And
5:21
this is just the beginning. We haven't
5:23
even touched storage engines, data
5:25
center failures, or performance
5:27
optimization across thousands of
5:29
machines. The next time you add
5:31
something to your cart and it just
5:33
works, you'll know there's an intricate
5:35
dance happening across data center
5:37
around the world. All to make sure your
5:39
data is exactly where you expect it to
5:41
be.
5:43
Ready to ace your next technical
5:44
interview? Join our community where we
5:47
offer comprehensive courses on system
5:49
design, coding, behavioral questions,
5:52
machine learning, and object-oriented
5:54
design. Learn more at byitebico.com.