I built a distributed key‑value store in Go using the Redis source code to guide me
I spent four weeks building a distributed key‑value store from scratch in Go, using Redis as a reference.
I spent the last four weeks building a distributed key-value store from scratch in Go, with the Redis source code open in another window as a reference. The project broke into four phases: starting a cluster, initializing it, propagating state between nodes, and resharding when nodes get added or removed. Each phase had a few subproblems of its own.
If you'd rather watch than read, here's the video version.
Phase 1: Starting Server Processes
This was the easiest part of the project. All I needed was a CLI command that accepts cluster start and hands it to the main program, which then launches as many server processes as the CLI asked for.
Redis requires at least three nodes and allows up to a thousand, and it configures each one through a config file. I eventually moved to Docker Compose to spin up the containers instead.
At this point, each node knows about itself and nothing about the other nodes or their state, so you end up with a handful of unconfigured servers sitting idle. Fixing that is the job of the next command.
Phase 2: Configuring Clusters
After reading into how Redis handles cluster create, I found it solves two subproblems. It has to divide the data among the nodes, and it has to get the nodes talking to each other so they can keep their state up to date.
Splitting the data with hash slots
Redis splits its data into 16,384 hash slots and divides them almost evenly across every node in the cluster. Building this took three steps:
- Have the CLI accept
cluster create. - Check that every node listed in the command is valid, exists, and is alive, then connect to each one.
- Keep some state on the CLI side that assigns each node its share of the slots, then push that state out to every node.
Making the nodes meet
With slots configured everywhere, the nodes still had to find each other. I pick a bootstrap node (usually the first argument passed to cluster create) and have every other node open a TCP connection to it.
This runs right after we divide the slots for all nodes on the CLI side.
If you look closely at the animation, you'll notice the bootstrap node now knows about every node and every node knows about the bootstrap node, but the rest of the nodes still don't know about each other.
There's also a deeper problem. Each node's state keeps changing over time, and with no way to propagate those changes across the cluster (yet), it would be hard to call the data distributed or sharded at all.
Phase 3: State Propagation
This is implemented using a cron job that runs roughly 10 times every second.
There are two data structures worth walking through for this phase.
ClusterState
ClusterState is what actually gets propagated between server processes. Every server keeps its own copy, maintains it, and sends it to the others. Among other fields, it holds:
- a pointer to the ClusterNode this server process represents
- a map from node ID to node, for fast lookups by name
- a global array of slot ownership that says which node owns which slot, which is the most important field in the struct
ClusterLink
ClusterLink does most of the work in this phase. A link lives as long as its cluster node lives, so I don't have to keep dialing a new TCP connection per node, and connection handling stays in its own structure away from the node state.
The structure itself is pretty simple. It wraps a Go channel. When a node wants to send something, it writes to the channel, and a writer goroutine dedicated to that link copies whatever arrives into the TCP buffer of the destination node. On the receiving end, a reader goroutine pulls it out of the TCP buffer and updates that server node's ClusterState accordingly, and sends a message back to the source node with an OK. That's the whole propagation mechanism.
Phase 4: Rebalancing and Resharding
Once state propagation worked, I could start rebalancing the cluster whenever a node joins or leaves, which really means moving slots and the keys inside them. It happens in three steps:
- Recalculate how many slots each node should have.
- Sort the nodes by how many slots they need to give away or take on.
- Move the keys in the affected slots from each source node to its destination.
Migrating keys
Fetching and moving keys was my favorite part of the whole project. The code breaks into four chunks: fetch a batch of keys, decode them, migrate them to the target node, and delete the migrated keys from the source.
The migration step is the interesting one. On the source node, slots own keys, and you're lifting a batch of keys out of that slot-to-key structure and dropping them into the same structure on the target node.
That covers one iteration. Both Redis and my version keep querying the requested slots for keys until every one of them is empty, and at that point the resharding is done (for that node).
Wrapping up
I ended up solving every problem I set out to solve. Along the way I got real practice with concurrency in Go and learned a lot about Redis internals and how distributed data systems are built, though I'm only scratching the surface of all three.
I'm planning more projects like this over the next year, and I'll post each one here and on YouTube. If you've made it this far, thanks for reading, and consider leaving a star :)
GitHub: github.com/omavashia2005/emberdb