Skip to main content

Storage Format

This document describes the storage layout and file formats used by OSO Kafka Backup.

Directory Structure​

Backups are organized in a hierarchical directory structure:

<storage-root>/
├── <backup-id>/
│ ├── manifest.json # Backup metadata
│ ├── topics/
│ │ ├── <topic-name>/
│ │ │ ├── partition=<id>/
│ │ │ │ ├── segment-00000000000000000000.bin.zst # Segment files, named by first offset
│ │ │ │ ├── segment-00000000000000104857.bin.zst
│ │ │ │ ├── segment-00000000000000209714.bin.zst
│ │ │ │ └── ...
│ │ │ └── ...
│ │ └── ...
│ ├── offsets.db # SQLite offset store; only with continuous: true or offset_storage:
│ └── consumer-groups-snapshot.json # Only with backup.consumer_group_snapshot: true
├── <another-backup-id>/
│ └── ...
├── evidence-reports/<report-id>/<YYYY>/<MM>/<report-id>.json | .pdf | .sig # validation evidence
└── offset-snapshots/<snapshot-id>/snapshot.json | metadata.json # offset-rollback snapshots

<storage-root> is the filesystem path, or <bucket>/<prefix> for S3, Azure and GCS (prefix is joined as {prefix}/{key}; filesystem and memory backends have no prefix). CLI --path URLs: s3://bucket/prefix, az://account.blob.core.windows.net/container, gs://bucket, file:///path. Only s3:// keeps a path prefix — az:// / gs:// discard it, so use the YAML prefix: key for those (#162).

There is no state/, checkpoints/ or checkpoint.json; backup progress lives in offsets.db (see Progress state).

Manifest File​

manifest.json is the index of the backup: which topics and partitions were captured and where every segment lives. It is written by the backup engine and merged on every save — segments are de-duplicated by (key, start_offset) with the existing entry winning, so re-runs and resumes never drop segments.

{
"backup_id": "production-backup-001",
"created_at": 1733220000000,
"source_cluster_id": "prod-cluster-east",
"source_brokers": [],
"compression": "zstd",
"topics": [
{
"name": "orders",
"original_partition_count": 3,
"partitions": [
{
"partition_id": 0,
"segments": [
{
"key": "production-backup-001/topics/orders/partition=0/segment-00000000000000000000.bin.zst",
"start_offset": 0,
"end_offset": 124999,
"start_timestamp": 1733216400123,
"end_timestamp": 1733219999871,
"record_count": 125000,
"uncompressed_size": 268435456,
"compressed_size": 41943040
}
]
},
{
"partition_id": 1,
"segments": [ "..." ],
"gaps": [
{
"start_offset": 90000,
"end_offset": 90512,
"reason": "offset_out_of_range",
"detected_at": 1733218800500
}
]
}
]
}
]
}

Manifest Fields​

FieldTypeDescription
backup_idstringBackup identifier (also the key prefix)
created_atintCreation time, epoch milliseconds
source_cluster_idstring | nullbackup.source_cluster_id, if set
source_brokersarrayReserved; currently always []
compressionstringSegment codec: zstd, lz4 or none (level is not recorded)
topics[].namestringTopic name
topics[].original_partition_countintPartition count on the source at backup time; drives create_topics on restore
topics[].partitions[].partition_idintSource partition
topics[].partitions[].segments[].keystringFull storage key of the segment — this is how restore and validate locate it
…segments[].start_offset / end_offsetintFirst / last source offset in the segment
…segments[].start_timestamp / end_timestampintEarliest / latest record timestamp (epoch ms); used for PITR segment selection
…segments[].record_countintRecords in the segment
…segments[].uncompressed_size / compressed_sizeintBytes before / after compression
topics[].partitions[].gaps[]arrayPresent only when retention deleted offsets before they were fetched: start_offset, end_offset, reason (offset_out_of_range), detected_at. Reported by validate and describe.

Totals (records, segments) are computed from the manifest on read; there is no statistics, time_range or version object.

Segment Files​

Segment files hold the record data of one partition, in offset order. They are named segment-<first offset, 20 digits>.bin.<codec> — .zst for zstd, .lz4 for lz4, .bin for none — and live under topics/<topic>/partition=<id>/.

Segment File Format​

┌─────────────────────────────────────────┐
│ Segment Header (32 bytes) │
├─────────────────────────────────────────┤
│ Compressed record block │
│ (Record 1, Record 2, ... Record N) │
├─────────────────────────────────────────┤
│ Segment Footer (8 bytes) │
└─────────────────────────────────────────┘

All integers are little-endian.

Segment Header​

OffsetSizeFieldDescription
04MagicKBAK
41VersionFormat version (currently 1)
51Compression0 = none, 1 = zstd, 2 = lz4
62ReservedZero
88Record CountNumber of records in the segment (u64)
168Start OffsetOffset of the first record (i64)
248End OffsetOffset of the last record (i64)

The record block between header and footer is compressed as one unit with the codec named in the header. There is no per-batch framing: once decompressed, it is a plain sequence of records.

OffsetSizeFieldDescription
04CRC32CRC32 of every preceding byte (header + compressed block)
44End MagicBKAE

Record Format​

Each record in the decompressed block is length-prefixed. Lengths of -1 mean null, which is different from 0 (empty) — Kafka distinguishes the two for keys, values and header values, and so does the segment.

┌─────────────────────────────────────────┐
│ Total Length (u32) — bytes that follow │
├─────────────────────────────────────────┤
│ Timestamp (i64, epoch ms) │
├─────────────────────────────────────────┤
│ Offset (i64, source offset) │
├─────────────────────────────────────────┤
│ Key Length (i32, -1 = null) │
│ Key bytes │
├─────────────────────────────────────────┤
│ Value Length (i32, -1 = null) │
│ Value bytes │
├─────────────────────────────────────────┤
│ Header Count (u16) │
│ Headers, in order: │
│ Key Length (u16), Key bytes (UTF-8) │
│ Value Length (i32, -1 = null), bytes │
└─────────────────────────────────────────┘

Because the source offset is stored in every record, restore can build the source→target offset mapping without relying on any header.

Offset-tracking headers​

With backup.include_offset_headers: true (the default), the backup appends these headers after the record's own headers:

Header KeyValueAdded when
x-original-offsetsource offset, little-endian i64 (8 bytes)always
x-original-timestampsource timestamp (epoch ms), little-endian i64always
x-source-clusterbackup.source_cluster_id, UTF-8source_cluster_id is set

A restore with include_original_offset_header: true (or consumer_group_strategy: header-based) adds x-original-offset, x-original-timestamp and x-source-partition (little-endian i32) to the records it produces; restore.strip_offset_headers: true (v0.19.0+) removes all of these from archived records before producing. See Offset-tracking headers.

Fidelity across versions

Versions before v0.18.0 wrote a null header value with length 0 instead of -1 (#155); such archives cannot be repaired. Duplicate header keys on one record are collapsed to the last one at backup time (#156).

Legacy JSON segments​

Segments written before the binary format are a compressed JSON array of records with base64-encoded key, value and header value fields (null for a null value). Restore recognises them by the missing KBAK magic and still reads them.

Progress state​

There is no checkpoint.json. Two separate mechanisms track progress:

Backup progress: the SQLite offset store​

The backup engine keeps per-partition progress in a local SQLite database and uploads that file byte-for-byte to <backup-id>/offsets.db:

CREATE TABLE IF NOT EXISTS offsets (
backup_id TEXT NOT NULL,
topic TEXT NOT NULL,
partition INTEGER NOT NULL,
last_offset INTEGER NOT NULL,
checkpoint_ts INTEGER NOT NULL DEFAULT (strftime('%s', 'now') * 1000),
PRIMARY KEY (backup_id, topic, partition)
);

CREATE TABLE IF NOT EXISTS backup_jobs (
backup_id TEXT PRIMARY KEY,
source_cluster_id TEXT,
status TEXT NOT NULL DEFAULT 'running',
created_at INTEGER NOT NULL DEFAULT (strftime('%s', 'now') * 1000),
last_heartbeat INTEGER NOT NULL DEFAULT (strftime('%s', 'now') * 1000),
last_checkpoint INTEGER
);
  • The store exists only when backup.continuous: true or an offset_storage: section is configured. A one-shot backup without either has no offsets.db and cannot resume.
  • Local path: offset_storage.db_path, else $TMPDIR/<backup-id>-offsets.db.
  • Progress is written after every fetched batch and synced to storage at most every backup.sync_interval_secs (default 30 s), at the end of every cycle, and on completion (backup_jobs.status = 'completed').
  • Resume: at startup the engine downloads <backup-id>/offsets.db if the local database has no rows, then starts each partition at last_offset + 1.
Options that are currently ignored

backup.checkpoint_interval_secs, offset_storage.s3_key and offset_storage.sync_interval_secs are accepted but have no effect — the key is always <backup-id>/offsets.db and the cadence is backup.sync_interval_secs (#161).

Restore checkpoint​

restore.checkpoint_state names a local JSON file (it is not written to object storage):

{
"backup_id": "production-backup-001",
"start_time": 1733220000000,
"last_checkpoint_time": 1733220315000,
"segments_completed": [
"production-backup-001/topics/orders/partition=0/segment-00000000000000000000.bin.zst"
],
"segments_in_progress": [["production-backup-001/topics/orders/partition=1/segment-00000000000000000000.bin.zst", 12582912]],
"records_restored": 125000,
"bytes_restored": 268435456,
"config_hash": "a3f1c9…"
}
caution

The file is only loaded and updated if it already exists — a fresh restore does not create it, so resumable restores currently require a pre-existing checkpoint file (#160).

Object Storage Layout​

On S3, Azure Blob and GCS the same tree maps to object keys under the configured prefix:

s3://my-bucket/kafka/
├── backup-001/manifest.json
├── backup-001/offsets.db
├── backup-001/consumer-groups-snapshot.json
├── backup-001/topics/orders/partition=0/segment-00000000000000000000.bin.zst
├── backup-001/topics/orders/partition=0/segment-00000000000000125000.bin.zst
├── backup-001/topics/orders/partition=1/segment-00000000000000000000.bin.zst
└── backup-002/...

kafka-backup list --path s3://my-bucket/kafka discovers backups by listing */manifest.json.

kafka-backup writes every object with the bucket's default storage class; there is no storage_class option (#163). Use a bucket lifecycle policy to transition older backups to infrequent-access or archive tiers, remembering that objects in GLACIER / DEEP_ARCHIVE must be restored before a kafka-backup restore can read them.

Compression​

backup.compressionExtensionNotes
zstd (default).zstLevel from backup.compression_level (default 3; only zstd uses it)
lz4.lz4lz4_flex block format with a size prefix — not the standard LZ4 frame format
none(none)Raw record block

Compression covers only the record block between the 32-byte header and the 8-byte footer, so a .zst segment is not a standalone zstd file — zstd -d will not decompress it; use kafka-backup validate --deep or a restore to read it. There is no gzip or snappy segment codec (broker-side gzip/snappy batches are decoded on fetch).

Integrity Verification​

Every segment carries one checksum: a CRC32 of the header plus the compressed record block, stored in the footer before the BKAE end magic. It is verified on every read (restore included), not only during validation. There is no per-record or per-batch checksum.

# Manifest loads; every segment key exists; sizes match compressed_size; gaps reported
kafka-backup validate --path /data --backup-id backup-001

# Additionally opens every segment: end magic + CRC32, decompresses, and checks
# record_count / start_offset / end_offset against the manifest
kafka-backup validate --path /data --backup-id backup-001 --deep

The report is printed to stdout (JSON with --format json); nothing is written back to storage. A missing or unreadable manifest, missing segment, or corrupt segment makes the command exit non-zero.

Compatibility​

  • Segment format version is 1 (KBAK header byte 4). Readers reject any other version outright; there is no forward-compatibility mode.
  • Manifest has no version field. All fields added since the first release (source_cluster_id, source_brokers, compression, original_partition_count, gaps, uncompressed_size, compressed_size) default when absent, so old manifests load in new readers, and new manifests load in old readers as long as the base fields are present.
  • Legacy JSON segments (pre-binary-format) are still restorable — see Legacy JSON segments. Compressed legacy segments could not be decompressed before v0.19.2. validate --deep does not understand them and reports them as corrupt.
  • Null header values are written correctly from v0.18.0; archives from earlier versions store them as empty and cannot be repaired (#155).