AtomicKV, part 1 of 2
Building a Key-Value Database from Scratch in C++: epoll, LRU Cache and B-Tree
How one AtomicKV node works: a single-threaded Linux epoll event loop, an O(1) LRU cache behind a Bloom filter, and an on-disk B-Tree in C++17, with benchmarks.
- C++
- Databases
- Systems
- Linux
TL;DR. AtomicKV is a key-value database I wrote from scratch in C++17. One node runs a single-threaded Linux epoll event loop over non-blocking sockets, keeps hot keys in an O(1) LRU cache, puts a Bloom filter in front of the disk so that missing keys are rejected without any I/O, and persists everything in an on-disk B-Tree whose values live in a separate append-only file. On my laptop, a single node served 10,000+ requests per second at ~16 ms average latency (p99 ~40 ms) with 200 concurrent clients. This post covers one node; part 2 covers how nodes form a cluster.
The code is on GitHub: stym01/AtomicKV.
Why build a database from scratch?
I wanted to understand what actually happens inside stores like Redis and Cassandra: how one thread serves thousands of sockets, how a cache decides what to forget, how a B-Tree lays itself out on disk, and what "eventually consistent" means in code rather than on a slide. Reading about these ideas was not enough, so AtomicKV implements each of them itself, with nothing beyond the C++ standard library and POSIX.
What a single node looks like
clients (TCP, one command per line)
│
▼
┌──────────────────────────────────────────────┐
│ epoll event loop (one thread) │
│ accept → read → parse → execute → reply │
└──────────────────────┬───────────────────────┘
▼
┌──────────────────────────────────────────────┐
│ KVStore │
│ LRU cache unordered_map + list │
│ Bloom filter 1,000,000 bits, 3 hashes │
│ B-Tree index btree.idx (keys, offsets)│
│ value log database.dat (append-only)│
└──────────────────────────────────────────────┘
A write goes to disk first, then to the cache. A read tries the cache, then asks the Bloom filter whether the key can exist at all, and only then walks the B-Tree.
The wire protocol
AtomicKV speaks a newline-terminated text protocol over raw TCP, so any terminal is a client:
$ nc localhost 8081
SET user_1 alice
OK
GET user_1
alice
SET session abc123 10 ← expires after 10 seconds
OK
DEL user_1
DELETED
A text protocol costs a little parsing, but it made debugging much easier: I could watch a request with nc and reproduce bugs by pasting commands. The cluster adds internal commands (INTERNAL_SET, INTERNAL_GET, MIGRATE, GOSSIP), which part 2 covers.
One thread, thousands of sockets: the epoll loop
The obvious server design is a thread per connection. It is simple, but every thread costs a stack and every switch between threads costs CPU time, which adds up quickly with thousands of mostly idle connections. AtomicKV instead uses one thread and Linux's epoll: the server registers every socket with the kernel, sleeps in epoll_wait, and wakes up only for sockets that actually have data.
int epoll_fd = epoll_create1(0);
event.events = EPOLLIN;
event.data.fd = server_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, server_fd, &event);
while (true) {
int num_ready = epoll_wait(epoll_fd, events, MAX_EVENTS, -1);
for (int i = 0; i < num_ready; i++) {
int fd = events[i].data.fd;
if (fd == server_fd) {
// New client: make it non-blocking and watch it for input
int client_socket = accept(server_fd, (struct sockaddr*)&client_addr, &client_len);
if (client_socket != -1) {
set_nonblocking(client_socket);
event.events = EPOLLIN;
event.data.fd = client_socket;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, client_socket, &event);
}
} else {
// Existing client sent data (or hung up)
char buffer[1024] = {0};
int valread = read(fd, buffer, sizeof(buffer));
if (valread <= 0) {
epoll_ctl(epoll_fd, EPOLL_CTL_DEL, fd, NULL);
close(fd);
} else {
process_client_request(fd, epoll_fd, std::string(buffer, valread));
}
}
}
}
A few details matter here:
- Non-blocking sockets.
set_nonblockingsetsO_NONBLOCKwithfcntl, so a slow client can never park the only thread insideread. - Level-triggered mode. I left epoll in its default level-triggered mode: if data is still unread after one
read, the nextepoll_waitreports the socket again. That makes the simple one-read-per-event loop safe. Edge-triggered mode is faster in some workloads, but it requires draining every socket untilEAGAIN, which is easy to get wrong. SIGPIPEis ignored, so a client that disconnects while the server is replying cannot kill the process.
The event loop owns all client I/O. Most cluster work (replication, gossip, read repair, anti-entropy and key migration) runs on five background threads that receive work through mutex-protected queues, so it does not block the loop. The one exception, forwarding a request to the node that owns the key, is covered in part 2.
The RAM tier: an O(1) LRU cache
Hot keys live in a least-recently-used cache built from two standard containers: an std::unordered_map for O(1) lookup, and an std::list that keeps keys in order of last use. Each map entry stores the value and an iterator to the key's list node, so moving a key to the front is a constant-time splice, with no search:
// map: key -> { entry, position of the key in the recency list }
std::unordered_map<std::string, std::pair<Entry, std::list<std::string>::iterator>> store;
std::list<std::string> lru_list; // front = most recently used
if (store.find(key) != store.end()) {
store[key].first = {value, expiry, version};
lru_list.splice(lru_list.begin(), lru_list, store[key].second); // mark as most recent
} else {
if (store.size() >= capacity) {
std::string lru_key = lru_list.back(); // least recently used
lru_list.pop_back();
store.erase(lru_key); // still on disk: this is tiered storage
}
lru_list.push_front(key);
store[key] = {{value, expiry, version}, lru_list.begin()};
}
Eviction only removes a key from RAM. The key stays in the B-Tree, which is why the cache can be small without losing data.
TTLs are lazy. SET key value 3600 stores an absolute expiry time. Nothing scans for expired keys; instead, a GET that finds an expired entry deletes it (from the cache and, as a soft delete, from disk) and answers NULL. That avoids a background sweeper thread. The trade-off is that expired keys nobody reads keep taking space until someone touches them, which is why Redis combines lazy expiry with periodic random sampling.
Rejecting missing keys before touching the disk: the Bloom filter
On a cache miss, the slow path is a walk down the B-Tree on disk. For keys that do not exist at all, that walk is pure waste, so a Bloom filter sits in front of the disk. It is a 1,000,000-bit array (about 122 KB) and three hash functions: std::hash, FNV-1a and DJB2.
bool BloomFilter::possibly_contains(const std::string& key) const {
if (bits.empty()) return false;
if (!bits[hash1(key) % bits.size()]) return false;
if (num_hashes > 1 && !bits[hash2(key) % bits.size()]) return false;
if (num_hashes > 2 && !bits[hash3(key) % bits.size()]) return false;
return true;
}
If any of the three bits is zero, the key was never written, and the server answers NULL without any disk read. If all three are set, the key might exist, and only then does the B-Tree get involved. On startup, the node rebuilds the filter from the keys already on disk.
A Bloom filter never gives false negatives, but its false-positive rate grows with the number of keys n. For m bits and k hashes, the rate is about (1 − e^(−kn/m))^k. With m = 1,000,000 and k = 3:
| Keys stored | False-positive rate |
|---|---|
| 50,000 | ~0.3% |
| 100,000 | ~1.7% |
| 200,000 | ~9% |
| 500,000 | ~47% |
So this filter is sized for roughly 100,000 keys. At 10 bits per key, the optimal number of hashes is (m/n)·ln 2 ≈ 7, which would bring the rate down to about 0.8%. Sizing the filter from the expected key count is on my list. One more property worth knowing: a standard Bloom filter cannot forget. A deleted key still passes the filter and falls through to the B-Tree, which knows the key is deleted.
The disk tier: a B-Tree index with values in a separate file
Everything is persisted in a B-Tree that I wrote from scratch. The tree lives in btree.idx: the first 8 bytes hold the file offset of the root, followed by fixed-size nodes. Values do not live in the tree at all; they are appended to database.dat, and each tree entry stores where its value starts and how long it is:
#define MAX_KEYS 3
const int MAX_KEY_LEN = 64;
// One fixed-size block in btree.idx
struct BTreeNode {
bool is_leaf;
int num_keys;
char keys[MAX_KEYS][MAX_KEY_LEN];
size_t value_offsets[MAX_KEYS]; // where the value starts in database.dat
size_t value_lengths[MAX_KEYS]; // and how long it is
bool is_deleted[MAX_KEYS]; // soft deletes (tombstones)
uint64_t versions[MAX_KEYS]; // Lamport clock version of each key
size_t child_pointers[MAX_KEYS + 1]; // file offsets of child nodes
};
Separating keys from values keeps every node the same size, so a node can be rewritten in place at a known offset no matter how long its values are. Values are only ever appended, so writes to database.dat are sequential.
MAX_KEYS is 3 on purpose. With nodes this small, splits happen constantly, which kept the split logic under test the whole time I was building it. A production engine sizes each node to a disk page (4 KB or more), so that a node holds dozens to hundreds of keys and the tree stays three or four levels deep even with millions of keys.
Inserting with proactive splits
Insertion uses the single-pass algorithm from CLRS (Introduction to Algorithms): on the way down, any full child is split before the insert descends into it, so an insert never has to walk back up the tree. When the root itself is full, a new root is allocated above it and the tree grows one level taller:
if (root.num_keys == MAX_KEYS) {
size_t new_root_offset = _allocate_node();
BTreeNode new_root;
new_root.is_leaf = false;
new_root.num_keys = 0;
new_root.child_pointers[0] = root_offset;
_split_child(new_root_offset, new_root, 0, root_offset, root);
// ... insert into the correct half, then persist the new root pointer
root_offset = new_root_offset;
idx_file.seekp(ROOT_POINTER_OFFSET);
idx_file.write(reinterpret_cast<const char*>(&root_offset), sizeof(size_t));
idx_file.flush();
}
When a key already exists, its entry is overwritten only if the incoming version is at least the stored one. That single comparison is what lets the cluster in part 2 resolve conflicting writes with last-write-wins.
A lost-update bug, and why proactive splitting caused it
While writing this post I wrote a test that builds a tree, updates every key and reads each one back. 209 of 6,000 updates were lost: a read returned the old value.
The cause is a subtle interaction with proactive splitting. Suppose a leaf holds [c, d, e] and you update d. On the way down, the leaf is full, so it is split before the descent, and its median, d, moves up into the parent. The insert then compares the key with the separator it just created, finds that d is not greater than d, descends into the left half, and inserts a second d there. Searches stop at the first match, the old d in the parent, so the new value is unreachable.
The fix is to check, right after each split, whether the key being written is the median that just moved up, and if so update it in place:
_split_child(node_offset, node, i, node.child_pointers[i], child);
// The split moved the child's median up into this node. If that median is the
// key being written, update it here instead of inserting a duplicate below it.
if (key == node.keys[i]) {
if (overwrite_if_newer(node, i, val_offset, val_len, version)) {
_write_node(node_offset, node);
}
return;
}
The same check is needed when the root splits. With both in place, the same test loses 0 of 6,000 updates. The lesson I took away: test updates, not just inserts.
Deletes are tombstones
DEL does not restructure the tree; it sets is_deleted on the entry. Real B-Tree deletion (borrowing from siblings and merging nodes) is the hardest part of the data structure, and tombstones also turn out to be useful in a replicated system, because a deletion has to be remembered in order to be replicated. The cost is space: neither tombstones nor overwritten values in database.dat are reclaimed yet. That needs compaction.
The write path and the read path
Putting the pieces together, this is what one SET does:
- Increment the node's Lamport clock to get the write's version.
- Append the value to
database.dat, then insert or update the key in the B-Tree, flushing both files. - Add the key to the Bloom filter.
- Insert or refresh the key in the LRU cache, evicting the least recently used key if the cache is full.
- Reply
OK, then queue replication to two other nodes (part 2).
And one GET:
- Cache hit: return the value and move the key to the front of the recency list.
- Cache miss: if the Bloom filter says the key cannot exist, return
NULLimmediately. - Otherwise walk the B-Tree, read the value from
database.dat, put it in the cache, and return it.
What "durable" means here. The value and the index update are both written before the client gets OK, but only as far as the operating system: nothing calls fsync, so a power failure can lose recent writes. The index pages also go through a buffered std::fstream, so a crash of the process right after OK can lose the newest index update; an explicit flush closes that gap. An fsync on every write would cost a lot of throughput, which is why databases use a write-ahead log with group commit instead. That is a natural next step.
Locking. The event loop is single-threaded, but the background threads also read and write the store, so KVStore is guarded by a std::shared_mutex. Even GET takes the exclusive lock, because a read changes the LRU order and may delete an expired key. That serialises reads, which is fine for one event-loop thread, but sharding the store by key hash would be the first step towards using more cores.
Benchmark: 10,000+ requests per second on one node
benchmark.py opens 200 client threads. Each keeps one persistent connection and sends 400 requests, half SET and half GET over its own set of keys, waiting for each reply before sending the next request. Latency is the round trip measured by the client.
| Metric | Result |
|---|---|
| Requests | 80,000 (200 clients × 400) |
| Throughput | 10,365 – 11,167 requests/s across runs |
| Average latency | 15 – 16 ms |
| p99 latency | 38 – 41 ms |
Setup: a single node on WSL2 (Ubuntu 24.04) on a Ryzen 5 5500H laptop with 8 vCPUs, running on AC power, with client and server on the same machine. On battery in power-saver mode, the same test drops to about 7,500 requests/s at ~22 ms, which is a good reminder to always report the machine state with a number.
Two sanity checks on these numbers:
- Little's law says the number of requests in flight equals throughput × latency. With 200 clients and about 10,500 requests/s, each request should spend about 200 / 10,500 ≈ 19 ms in the system. The measured average of 15–16 ms is a little lower, which makes sense: the client also spends time in Python between receiving one reply and sending the next request.
- The disk path is actually exercised. The server runs with a deliberately tiny cache, and the 200 clients touch tens of thousands of distinct keys (each picks from its own 1,000), so almost every
GETmisses RAM. It is answered either by the Bloom filter, for keys that were never written, or by a walk down the B-Tree. The throughput is not just an in-memory hash map.
The client is Python, and the GIL probably limits it before the server does. A load generator in C++ or Go would tell me how much headroom the server really has.
What I would change next
- Framing and buffering. The loop handles one
readof up to 1 KB per event and treats it as a complete command. A command split across TCP segments, or two commands in one segment, needs a per-connection input buffer and proper framing. Large replies also need a per-connection output buffer andEPOLLOUT. - Page-sized B-Tree nodes, plus compaction of
database.datto reclaim overwritten values and tombstones. - A write-ahead log with group commit for real durability without an
fsyncper request. - Persisted TTLs. Expiry times live only in the cache today, so a key that is evicted and later reloaded from disk loses its TTL.
- A properly sized Bloom filter (about 7 hashes for 10 bits per key), or a counting Bloom filter so that deletes can clear bits.
- Sharded locks so that cache hits can run on more than one core.
Part 2 covers how AtomicKV becomes a cluster: consistent hashing, gossip, Lamport clocks, read repair and anti-entropy.
Frequently asked questions
What is AtomicKV?
AtomicKV is an open-source distributed key-value database that Satyam Kesharwani wrote from scratch in C++17. Each node runs a single-threaded Linux epoll event loop, keeps hot keys in an LRU cache, uses a Bloom filter to skip disk reads for missing keys, and stores data in an on-disk B-Tree. Nodes form a masterless cluster with consistent hashing, gossip, Lamport clocks, read repair and anti-entropy.
How fast is AtomicKV?
On a laptop (WSL2, Ryzen 5 5500H, AC power), a single node handled 10,000+ requests per second (10,365–11,167 across runs) at about 16 ms average latency and about 40 ms p99, with 200 concurrent clients sending 80,000 requests, half SET and half GET.
Why use epoll instead of one thread per connection?
A thread per connection costs a stack per client and CPU time for context switches, which adds up with thousands of mostly idle connections. With epoll, one thread registers every socket with the kernel and wakes up only for sockets that have data.
Why does AtomicKV keep values in a separate file from the B-Tree?
The B-Tree file holds fixed-size nodes with each key and the offset and length of its value, while the values themselves are appended to a separate data file. Fixed-size nodes can be rewritten in place regardless of value size, and appending values is a sequential write.
How does a Bloom filter speed up a key-value store?
A Bloom filter answers 'definitely not present' or 'maybe present' from a small bit array. AtomicKV checks it on every cache miss, so a GET for a key that was never written returns without reading the B-Tree from disk. Its filter uses 1,000,000 bits and three hash functions.