sonirico.dev

Marcos Sánchez
SRE @ Chess.com
Madrid

The Art of Not Buying Disk, Part II

Part I ended with rpkv: a sidecar that reads compacted Redpanda topics by key while storing only an index, key -> (partition, offset). The values never leave the log. It works, and one of its benchmark numbers is where this story starts. I’m telling it as a timeline because that’s the order I’ll want it back in.

The floor

rpkv answers a read in 11.1 ms p50 over loopback. I went looking for my bug and found the design instead: every lookup is a real Kafka fetch, so it pays TCP, protocol framing, and a fetch request the broker schedules like any other, all to retrieve one record the broker could read from its own disk in microseconds. In other words, the sidecar sits outside the process that owns the bytes, and every read pays admission.

There’s a second problem too. A sidecar can’t know where a key lives without scanning, so a key whose segment was evicted to tiered storage costs a scan through the remote read path, billed against the object store. You can’t fix either of these from outside the broker. So I went inside.

Reading someone else’s broker

I cloned redpanda and started reading. First discovery: Seastar. One shard per core, no locks, no shared memory between shards; every partition lives on exactly one shard and if you want something from it, you send your code to the data with submit_to instead of reaching for a mutex. Everything, meanwhile, is a future. I came from Go, where concurrency is cheap and crossing cores is nobody’s business, whereas here the core you’re on is the whole architecture. It took me a day of reading before I trusted myself to write ten lines.

Second discovery, the one that made the project feasible: almost every tool I needed was already in the tree.

So the job shrank on contact: not building a key-value engine, just threading a wire through machinery somebody else had already installed.

The wire

Fourteen commits, 1,708 lines added, on a PoC branch with a staging PR so the series can be read as a review. The order matters: feature flag first, then a kv_index_enabled cluster property, then the redpanda.kv.index.enabled topic property, and only then the index itself. Everything defaults to off, and the flags went in before the feature.

The index is storage::kv_index: one lsm::database per partition, under <ntp dir>/kv_index/, no WAL. The stored value is the record’s Kafka offset, eight bytes, translated at index time:

if (r.is_tombstone()) {
    wb.remove(sv, lsm::sequence_number(o()));
} else {
    iobuf v;
    auto be = ss::cpu_to_be(
      static_cast<uint64_t>(ot.from_log_offset(o)()));
    v.append(reinterpret_cast<const char*>(&be), 8);
    wb.put(sv, std::move(v), lsm::sequence_number(o()));
}

The translation happens at index time for a reason: tiered storage prefix-truncates the local offset translator state when it evicts segments, so a raw raft offset stored today might be untranslatable tomorrow. A Kafka offset, on the other hand, stays valid for as long as the record exists anywhere.

Feeding the index is one line in the append path, in disk_log_appender, right between the segment write and the offset translator:

auto stop = co_await append_batch_to_segment(batch);
co_await _log.kv_index_batch(batch);
_log.offset_translator().process(batch);

Write path: one producer batch lands in the segment, and eight bytes of pointer land in the kv_index

Every replica runs this, so every replica has its own index and can serve its own reads. And since there’s no WAL, durability is borrowed: on open, if last_applied() is ahead of the log’s dirty offset, the index is thrown away and rebuilt from the log; a suffix truncation at or before the last indexed offset triggers the same rebuild. The index never argues with the log. It re-reads it.

The read path

The endpoint is GET /kv/{topic}/{key} on Pandaproxy, in kv_handlers.cc. Without an explicit ?partition=N it hashes the key with the default murmur2 partitioner, the same one a producer would use. It finds the owning shard through the shard table, and everything after that happens on that shard, in-process: index lookup, gate checks, record read through partition_proxy. No network hop, which was the whole point of moving in.

Read path: partition resolution, then the index lookup, then three gates any of which can turn the answer into a 404, then the record read and the key recheck

The gates, in the order the request meets them:

if (
  koff >= model::offset_cast(cluster::kafka_high_watermark(*p))) {
    co_return kv_lookup_result{
      kv_lookup_result::status::not_found};
}
auto local_start = model::offset_cast(
  p->log()->from_log_offset(p->raft_start_offset()));
if (
  koff < local_start
  && !config::shard_local_cfg().kv_index_remote_read_enabled()) {
    co_return kv_lookup_result{
      kv_lookup_result::status::not_local};
}

The high watermark gate exists because the index is fed at append time, before raft commit. An offset the index knows about may still be truncated away by a leadership change, so anything at or beyond the high watermark is answered 404 rather than shown to a client and later unwritten. The other gate, local_start, is the tiered storage decision: a lookup below the local start would be a cold random read against the object store, one remote segment chunk per miss. That bill belongs to whoever runs the cluster, so it sits behind kv_index_remote_read_enabled and stays off until they sign it.

Even after the gates, the record that comes back is not trusted. The read path rechecks that the key of the record it fetched equals the key that was asked for:

if (!found->has_key() || !(found->key() == key)) {
    co_return kv_lookup_result{
      kv_lookup_result::status::not_found};
}

That recheck is the same discipline rpkv already had, and it covers the stale pointer case: compaction can remove or supersede the offset the index still points at, and a mismatch is just a miss. Success returns the raw value with X-Kv-Partition and X-Kv-Offset headers, while a broker that doesn’t host the partition answers 421 rather than forwarding. And a replica that arrived by partition move, recovered from a snapshot, only indexes what was in its local log plus what came after, so two replicas can disagree about old keys.

Left open, on purpose

Records from aborted transactions are indexed and served; the fetch path filters them and this path doesn’t yet. There are no quotas on the endpoint, no metrics on the index, and the per-replica rebuild could arguably be replication instead. All of it is listed in the design note as unresolved, because a proof of concept that hides its holes is asking for a yes instead of asking for a review.

The secondary index of the work

I posted the whole thing as discussion #31597 on August 17. Zero comments so far, which I’m assured is how upstream RFCs work: you write your key into someone else’s compacted topic and wait to find out whether it survives the next pass. Maintainer attention has an eviction policy too. Nobody publishes the settings.

So this post is my hedge. If they answer, it’s the story behind the PR. If they never do, it’s the index entry: the key, and where the work lives, one fetch away from a future me who’ll have forgotten every line of it. Either way, I’m not paying for my memory twice.