How storage works

Each partition is a directory of segments. A segment is a log file of length-prefixed, CRC-checked records and an index file that maps offsets to positions in it. This page covers the format, when data is durable, how a crash is recovered, and how retention and compaction rewrite segments safely.

Files on disk

data directory
data/
├── topics/orders/
│   ├── topic.json                       # partition count and config
│   ├── 0/
│   │   ├── 00000000000000000000.log
│   │   ├── 00000000000000000000.index
│   │   ├── 00000000000000000147.log     # rolled at offset 147
│   │   └── 00000000000000000147.index
│   └── 1/ …
├── groups/analytics.json                # committed offsets and owner
├── producers.json, producers.journal    # idempotent producer state
├── schemas.json
└── api_keys.json                        # key hashes, never secrets

A segment is named after the offset of its first record, padded to 20 digits, so sorting the names sorts the segments. Every file and directory the broker creates is readable only by its own user.

Records

Integers are big-endian. record_len counts the bytes after itself; the CRC covers everything from offset to the end of value. A record is rejected unless its length is plausible and its CRC matches, before anything is allocated or copied for it.

recordbig-endian
record_len     u32      bytes that follow
offset         u64
timestamp_ms   i64
key_len        i32      -1 = no key
key            [u8]
value_len      u32
value          [u8]
crc32          u32      over offset … value

Index and reads

The index holds 16-byte entries, (offset within segment, file position): one for the first record, then one every 4 KiB of log. A billion retained records need about 4 MB of index in memory. A read finds the segment, then the nearest entry at or before the offset, and scans forward at most 4 KiB.

Each segment keeps its own file handle open and reads with positional reads in large windows: a read of a hundred records is one or two system calls. A segment that compaction or retention removes keeps serving the readers that were already using it, from the same bytes.

Writes and durability

A produce request is one write per partition: its records are encoded into one buffer and appended with a single write, together with their index entries. With the default --flush-every-records 1, the response is sent only after an fsync that covers the write.

That fsync is shared. The first request to need one syncs the newest data it can see; requests waiting behind it find their write already covered and return without another system call. This is group commit, and it is why durability costs little under concurrent load.

--flush-every-records 1
Default. Everything acknowledged is on disk.
--flush-every-records N
Fsync once N records are waiting, and at least every --flush-interval (1 s). A crash can lose what arrived since the last fsync.

If a write fails part-way (a full disk, an I/O error), the partial bytes are cut off the files before the error is returned. If that is impossible, or an fsync fails, the partition stops accepting writes until a restart runs recovery: after a failed fsync nothing written since the last good one can be trusted.

Metadata files (topic config, offsets, keys, schemas) are replaced atomically: written to a unique temporary file, fsynced, renamed over the old one, and the directory fsynced.

Recovery after a crash

At startup each sealed segment is checked against its index: the first record and everything after the last index entry must decode, pass their CRC and end exactly at the end of the file. A segment that passes is trusted without reading it all; one that does not is scanned in full and its index rebuilt. The segment being written to is always scanned in full.

  • A torn tail in the active segment, from a crash mid-write, is cut off. Those records were never acknowledged.
  • A corrupt record in a sealed segment is not a torn write. The rest of that segment is skipped and reported, the file is left in place, and no other segment is touched.
  • The next offset continues after the highest valid record.

Retention and compaction

Two background tasks visit every partition. Each topic's cleanup policy decides which runs. The segment being written to is never deleted or compacted.

Retention

  • By age: a sealed segment goes once its newest record is older than retention_ms.
  • By size: the oldest sealed segments go until the partition is under retention_bytes.
  • With --cold-storage-dir, a segment is copied there and fsynced first. If the copy fails, the segment stays and is retried.

Compaction

  • Keeps the newest record of each key across the sealed segments, and every record without a key, at its original offset.
  • An empty value is a tombstone: it deletes its key and is itself dropped after tombstone_retention_ms.
  • Two streaming passes: the first remembers only each key's newest offset, the second copies the survivors. Memory grows with the number of keys, not the size of the data.
  • A pass that would drop nothing is skipped. An unreadable record aborts the pass rather than losing what follows it.

Replacing segments safely

  1. Write the result to temporary files and fsync them.

  2. Check nothing moved. Under the partition's write lock, the sealed segments must be exactly the ones that were read; a roll or a retention pass in the meantime cancels the swap.

  3. Record the plan in compaction.marker and fsync it. From here a crash is finished at the next start, so the new segment and the ones it replaces never coexist.

  4. Swap. Rename the result into place, delete the replaced segments, fsync the directory, and publish the new segment list. Readers already in progress keep their old files.

Throughput

Measured with crates/es-broker/tests/bench.rs: 256-byte records over loopback, release build with full link-time optimisation, on a 32-core server with NVMe storage. Every case runs with fsync before acknowledgement. The earlier figures are version 0.2.0 running the same bench.

Casev0.2.0Now
Binary produce, 8 connections, batches of 500, 1 partition15k/s4.6M/s
Binary produce, 8 connections, batches of 500, 4 partitions34k/s4.0M/s
HTTP produce, 8 clients, batches of 500, 4 partitions18k/s3.7M/s
Binary consume, 1 connection, 4 MiB fetches1.6M/s3.0M/s
64 pipelined requests of 10 records, 1 connection15k/s118k/s
Raft partition, single node, 4 connections, batches of 1005k/s1.0M/s

The runs are short; compare ratios rather than digits, and run the bench on your own hardware:

shell
cargo test --release -p es-broker --test bench -- --ignored --nocapture

What made the difference, largest first:

  1. One write and at most one fsync per request, instead of two fsyncs per record under the partition lock.
  2. One Raft entry per request, on an append-only Raft log, instead of rewriting the whole log as JSON per record.
  3. Group commit and concurrent request handling, so small requests share fsyncs.
  4. Windowed positional reads instead of opening the file and making four system calls per record.

Batch size dominates: small requests each pay for an fsync. More partitions stop helping once one partition is no longer the bottleneck.