AtomicKV, part 2 of 2
Distributing a C++ Key-Value Store: Consistent Hashing, Gossip and Read Repair
AtomicKV as a masterless cluster: consistent hashing with virtual nodes, gossip failure detection, Lamport clocks, replication, read repair and anti-entropy, in C++.
- Distributed Systems
- C++
- Databases
TL;DR. This is part 2 of the AtomicKV series. Part 1 built a single node: an epoll server, an LRU cache, a Bloom filter and an on-disk B-Tree. This post turns it into a masterless cluster. A consistent hash ring with 100 virtual nodes per server decides which node owns a key. Any node accepts any request and forwards it to the owner. The owner writes locally, replies, and replicates asynchronously to the next two nodes on the ring. Lamport-clock versions settle conflicting writes, gossip heartbeats detect failures, and key migration, read repair and anti-entropy pull the cluster back into agreement. While writing this post I also found and fixed two membership bugs; the before-and-after numbers are near the end.
The cluster at a glance
client
│ SET user_42 alice
▼
node 8081: hash("user_42") lands in an arc owned by 8082
│ forward the request
▼
node 8082 (owner): write locally with version v ──▶ reply OK
│
└─ background: INTERNAL_SET user_42 alice v ──▶ the next two nodes on the ring
There is no leader and no coordinator: every node runs the same binary and the same five background threads.
| Thread | Wakes up | Job |
|---|---|---|
| Gossip | every 2 s | Send its heartbeat and its view of the cluster to one random live peer; mark peers that have been silent for 10 s as dead |
| Migration | every 1 s, if the ring changed | Push keys this node no longer owns to their new owner |
| Replication | when work is queued | Send INTERNAL_SET, INTERNAL_DEL and MIGRATE commands to other nodes |
| Read repair | after every successful GET | Compare versions with the replicas and fix whichever side is stale |
| Anti-entropy | every 30 s | Compare every key this node owns with its replicas |
In CAP terms, AtomicKV is an AP system: when nodes fail or cannot reach each other, it keeps accepting reads and writes and reconciles afterwards, instead of refusing requests to stay strictly consistent.
Partitioning: consistent hashing with virtual nodes
The naive way to spread keys over N servers is hash(key) % N. It breaks badly when N changes: going from 3 to 4 servers moves about three quarters of all keys. Consistent hashing places both servers and keys on a ring of hash values, and each key belongs to the first server clockwise from it. Adding or removing a server then only moves the keys in the arcs next to it.
AtomicKV's ring is an ordered map from hash value to node. Each physical node is inserted 100 times, as ip:port#0 through ip:port#99:
void ConsistentHashRing::add_node(const std::string& ip, int port) {
std::unique_lock<std::shared_mutex> lock(ring_mutex);
std::string node_id = ip + ":" + std::to_string(port);
// ... return if the node is already on the ring
ClusterNode node = {ip, port, node_id};
for (int i = 0; i < vnodes_per_node; ++i) {
std::string vnode_key = node_id + "#" + std::to_string(i);
ring[hash_func(vnode_key)] = node;
}
ring_version++; // wakes up the migration worker
}
ClusterNode ConsistentHashRing::get_node_for_key(const std::string& key) {
std::shared_lock<std::shared_mutex> lock(ring_mutex);
if (ring.empty()) return {"", 0, ""};
auto it = ring.lower_bound(hash_func(key));
if (it == ring.end()) it = ring.begin(); // wrap around the ring
return it->second;
}
Why 100 virtual nodes instead of one point per server? With one point each, arcs have random sizes and one server can easily own twice its fair share. A hundred small arcs per server average out. They also spread the load of a failure: a dead server's arcs are taken over by many different successors instead of one unlucky neighbour.
Lookups vastly outnumber membership changes, so the ring sits behind a std::shared_mutex: lookups take a shared lock and joins or leaves take an exclusive one. Every change also bumps an atomic ring_version, which is how the migration thread notices that it has work to do.
Replicas come from the same ring: starting at the owner, walk clockwise and collect the next two distinct physical nodes, skipping further virtual nodes of servers already chosen.
std::set<std::string> seen_nodes;
seen_nodes.insert(it->second.id); // the owner
auto curr_it = it;
while ((int)replicas.size() < count) {
if (++curr_it == ring.end()) curr_it = ring.begin();
if (curr_it == it) break; // fewer nodes than replicas wanted
if (seen_nodes.insert(curr_it->second.id).second) {
replicas.push_back(curr_it->second);
}
}
One subtle caveat: the hash function is std::hash<std::string>, whose output is implementation-defined. All nodes must therefore run builds of the same standard library, or they will disagree about who owns what. A cluster running mixed builds would need a fixed hash such as MurmurHash or xxHash.
Routing: any node accepts any key
Clients do not need to know the topology. Whichever node receives a request hashes the key; if another node owns it, the receiver opens a connection to the owner, forwards the raw request and relays the reply. Internal commands (INTERNAL_SET, INTERNAL_GET, MIGRATE) skip this step, because they are already addressed to the right node.
if (!key.empty() && !ring.is_empty()) {
owner = ring.get_node_for_key(key);
if (owner.id != my_node_id
&& command != "INTERNAL_SET" && command != "INTERNAL_DEL"
&& command != "MIGRATE" && command != "INTERNAL_GET") {
needs_redirect = true; // forward to the owner and relay its reply
}
}
The benefit is that a plain nc session against any node works. The cost is an extra network hop, plus a limitation I would fix first: the forward is a blocking call (with 2-second timeouts) made inside the event loop, so while one request waits on a slow owner, that node serves nobody else. There are two standard fixes. One is to make the forward asynchronous, registering the owner's socket with epoll and completing the client's request when the reply arrives. The other is to stop forwarding and answer with a redirect, as Redis Cluster does with MOVED, letting smart clients cache the topology.
Replication: acknowledge first, copy in the background
The owner applies a write locally, replies OK, and only then queues copies for the two replicas. Each copy carries the version the owner assigned:
uint64_t version = db->set(key, actual_value, ttl);
result = "OK\n";
std::vector<ClusterNode> replicas = ring.get_replica_nodes(key, 2);
std::string internal_cmd = "INTERNAL_SET " + key + " " + actual_value
+ " " + std::to_string(version) + "\n";
for (const auto& rep : replicas) {
std::lock_guard<std::mutex> lock(rep_mutex);
replication_queue.push({internal_cmd, rep.ip, rep.port});
rep_cv.notify_one();
}
A dedicated thread waits on the queue with a condition variable and sends each command over a short-lived blocking socket with 2-second timeouts. Blocking is fine there, because it never touches the event loop.
Acknowledging before replicating keeps writes fast, but it opens a window: if the owner dies after replying OK and before the copies go out, that write survives only on the dead node. Dynamo-style systems make this trade-off tunable with quorums: a write succeeds only after W of the N replicas confirm it, and choosing R + W > N makes every read overlap the latest write. AtomicKV is effectively W = 1 today, and quorum writes are on my list.
Versions: Lamport clocks and last-write-wins
Replicas receive writes in different orders and at different times, so every copy needs a version to decide which value is newer. Wall-clock timestamps are risky because server clocks drift. AtomicKV uses a Lamport clock instead: an std::atomic<uint64_t> counter per node.
- Every local write increments the counter and uses the result as the write's version.
- Every version received from another node moves the counter forward to at least that version, so this node's next write is guaranteed to look newer than anything it has seen.
// Advance the Lamport clock when a newer version arrives from another node
uint64_t current = version_counter.load();
while (version > current && !version_counter.compare_exchange_weak(current, version)) {}
The compare-and-swap loop is there because background threads and the event loop can update the counter at the same time. Storage then keeps whichever copy has the higher version (the B-Tree check from part 1), which makes the cluster last-write-wins.
Last-write-wins is simple and convergent, but it has known limits. If two clients write the same key on different nodes at truly the same time, one write is silently discarded. Two nodes can also produce the same version number, and AtomicKV does not yet break that tie (comparing (version, node id) pairs would). Systems that must not lose concurrent writes keep both versions and let the client merge them, which is the vector-clock approach described in Amazon's Dynamo paper.
Failure detection: gossip
There is no master to declare nodes dead, so the nodes have to agree among themselves. Every 2 seconds, each node:
- increments its own heartbeat counter,
- marks every peer it has not heard a higher heartbeat from in 10 seconds as dead, and removes it from the ring,
- picks one random live peer and sends it a
GOSSIPmessage.
A gossip message carries the sender's own heartbeat, followed by the sender's view of every live member as ip:port:heartbeat. The receiver merges that view, adding nodes it did not know about, and replies with its own view in a GOSSIP_ACK, which the sender merges in turn:
std::string build_gossip_message(const std::string& type, int heartbeat) {
std::string msg = type + " " + my_ip + " " + std::to_string(my_port) + " " + std::to_string(heartbeat);
for (const auto& node : gossip_mgr.get_alive_nodes()) {
msg += " " + node.ip + ":" + std::to_string(node.port) + ":" + std::to_string(node.heartbeat);
}
return msg + "\n";
}
Heartbeats only ever increase, so it does not matter which peer relays the news: a node is alive as long as its heartbeat keeps rising somewhere in the cluster, and information spreads like an epidemic in O(log N) rounds. Removing a dead node from the ring bumps ring_version, which wakes the migration thread.
Two improvements I would make here: SWIM-style indirect probes (the SWIM paper) to detect failures faster with fewer false alarms, and generation numbers so that a restarted node, whose heartbeat starts again from zero, is not ignored by peers that remember a higher count.
Healing: migration, read repair and anti-entropy
Replicas drift apart: a copy is lost when a node is down, a node rejoins with old data, or the ring changes. Three mechanisms pull them back together.
Key migration runs whenever the ring changes. The node scans every key on disk, and any key whose owner is now someone else is pushed to the new owner with a MIGRATE command that keeps its version, so it cannot overwrite newer data.
Read repair runs after every successful GET. The client gets its answer immediately; a background task then asks each replica for its copy and compares versions:
if (replica_version > task.local_version) {
// The replica has newer data: adopt it locally
db->set_versioned(task.key, replica_value, replica_version);
} else if (replica_version < task.local_version) {
// The replica is stale: queue our version for it
std::string fix_cmd = "INTERNAL_SET " + task.key + " " + task.local_value
+ " " + std::to_string(task.local_version) + "\n";
// ... push fix_cmd to the replication queue
}
Anti-entropy covers keys nobody reads. Every 30 seconds, a node walks every key it owns and runs the same comparison against both replicas. That costs a full scan and one round trip per key per replica, so it gets expensive as data grows. Dynamo and Cassandra solve this with Merkle trees: replicas compare hashes of key ranges and only exchange the ranges that differ. That is the next step for AtomicKV.
Two bugs I found by following my own README
Before writing this post, I started a three-node cluster exactly as the README describes (one node, then two more that join through it as a seed) and pushed traffic through every node. It did not go well.
Bug 1: joining nodes never put themselves on their own ring. A node started with a seed added the seed to its ring, but added itself only when it had no seeds at all. So node 8082 believed 8081 owned every key, while 8081 believed 8082 owned some of them. A request for such a key bounced between the two nodes until a 2-second timeout fired. The fix is one line: every node adds itself to its ring at startup.
Bug 2: membership did not spread, and healthy nodes were declared dead. Gossip used to carry only the sender's own heartbeat, and the reply was ignored. Two nodes that joined through the same seed therefore never learned about each other, and their rings disagreed. Worse, a node refreshed a peer's heartbeat only when that peer happened to pick it at random. With three nodes, each direction between two nodes gets a message only about every 4 seconds on average, so a 10-second gap happens by chance quite often, and healthy nodes kept being declared dead and re-added. The fix is the merge-the-views gossip described above.
I ran the same script against the old and the fixed code: write 18 keys round-robin through the three nodes, read every key through every node, watch the logs for failure detections, then kill one node and read everything again through the two survivors.
| Three-node cluster, README setup | Before | After |
|---|---|---|
| Writes spread over all three nodes | 0 of 18 | 18 of 18 |
| Every key read through every node | 10 of 54 | 54 of 54 |
| Healthy nodes declared dead | 16 | 0 |
| Reads after killing one node | 0 of 36 | 36 of 36 |
The old run took several minutes, because requests hung until their timeouts, which also gave its gossip more time to misfire; the fixed run took about a minute. After the kill, both survivors removed the dead node from their rings within 15 seconds, migrated its keys, and answered every read from the replicas. The lesson for me: a distributed system has to be tested as a cluster, with failures, by following the same instructions a stranger would follow. Every single-node test I had passed.
What I would change next
- Asynchronous forwarding, or redirects, so a slow owner cannot stall another node's event loop.
- Quorum writes and reads (W and R out of N) plus hinted handoff, so acknowledged writes survive a single failure.
- Merkle-tree anti-entropy instead of comparing every key.
- SWIM-style failure detection and generation numbers for restarts.
- Tie-breaking on
(version, node id), or vector clocks where concurrent writes must not be lost. - A fixed hash function (MurmurHash or xxHash) so that differently built nodes agree on ownership.
The code is on GitHub at stym01/AtomicKV. A live instance runs on AWS EC2 if you want to try it: nc 100.53.37.203 8081.
Frequently asked questions
How does AtomicKV decide which node stores a key?
It uses consistent hashing. Every node is placed on a hash ring 100 times as virtual nodes, and a key belongs to the first node clockwise from the key's hash. The key's two replicas are the next two distinct physical nodes on the ring, and adding or removing a node only moves the keys next to it.
How does AtomicKV detect failed nodes without a master?
With gossip. Every 2 seconds each node increments its heartbeat and sends its view of the cluster to a random live peer, which merges it and replies with its own view. A node whose heartbeat has not increased for 10 seconds is marked dead and removed from the ring, which triggers key migration.
How does AtomicKV resolve conflicting writes?
Every write carries a Lamport-clock version. A node increments its counter for each local write and advances it past any version it receives, and storage keeps the copy with the higher version, so the cluster is last-write-wins.
Is AtomicKV a CP or an AP system?
AP. It keeps serving reads and writes when nodes fail, acknowledges writes before replicating them, and converges afterwards through read repair, anti-entropy and key migration.
What did testing the AtomicKV cluster reveal?
Following the README's three-node setup showed that nodes joining through a seed never added themselves to their own hash ring, and that gossip carried only the sender's own heartbeat, so requests looped between nodes and healthy nodes were declared dead. After both fixes, 54 of 54 reads through every node succeeded, and all 36 reads succeeded after one node was killed.