skip to content

Tiered Storage Internals

The internals of KIP-405 tiered storage: the plugin interfaces, local versus remote segments, and how a fetch is served from an object store. Interviewers ask it as a modern-Kafka question about decoupling retention from broker disk.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

What is Kafka Tiered Storage (KIP-405), and what problem does it solve at a high level?

level: juniorimportance: must knowfreq 70%

answer

  1. KIP-405, Kafka 3.6+
  2. local hot tier vs remote object store
  3. closed segments offloaded then deleted locally
  4. decouple storage from compute
  5. non-compacted topics only

basics

~10 s

Tiered Storage (KIP-405) moves old, closed log segments from local broker disks to cheaper remote object storage (like S3). Recent data stays local for fast reads; older data is fetched from remote when needed.

solid answer

~40 s

Kafka Tiered Storage, introduced by KIP-405 and production-ready in Kafka 3.6+, splits a topic's log into two tiers. Brokers keep recent ('local' or 'hot') segments on local disk for low-latency reads and writes, while older closed segments are copied ('offloaded') to a remote tier — typically an object store like S3, GCS, or HDFS — and then deleted locally. This decouples storage from compute: you can retain weeks or months of data cheaply without sizing broker disks to hold it all, and you store far less per broker so cluster rebalancing and recovery are faster. Consumers reading old offsets transparently fetch from remote storage. It only applies to non-compacted topics and is enabled per-topic with remote.storage.enable=true on a cluster where remote storage is configured.

go deeper

for a junior

Know the one-line purpose: old segments go to cheap remote storage, recent data stays local, total retention gets cheaper.

for a middle

Explain the two tiers, the per-topic remote.storage.enable flag, and that only closed non-compacted segments are offloaded.

for a senior

Articulate the decoupling of storage and compute, the RSM/RLMM plugin model, and operational benefits like faster rebalancing/recovery.

for a principal

Reason about when tiered storage is the right architecture vs. alternatives, cost/latency tradeoffs, and the failure/consistency implications of an external object store in the durability path.

## The problem A Kafka topic is stored as an append-only **log**, physically a sequence of **segment files** on each broker's local disk. Traditionally, how much history you can keep is bounded by local disk: to retain 30 days of data you must provision enough disk on every broker to hold 30 days. This couples two unrelated needs — **compute/throughput** (how many brokers you need to serve traffic) and **storage** (how much history you want to keep). It also makes brokers 'heavy': moving a large partition during rebalancing or recovering a failed broker means copying huge amounts of data. ## What KIP-405 does **Tiered Storage** introduces a second storage tier. Each partition's log is conceptually split: - **Local tier** — recent segments on the broker's local disk, serving writes and low-latency reads of fresh data. - **Remote tier** — older, **closed** (rolled, no longer being written) segments copied to an external **object store** (Amazon S3, Google Cloud Storage, Azure Blob, HDFS, etc.). Once a closed segment is safely uploaded to the remote tier, the broker can delete the local copy (subject to a separate local-retention setting). The **total retention** of the topic is now governed by remote retention, which can be far larger and far cheaper than local disk. ## Why it matters - **Cost**: object storage is ~10x cheaper per GB than provisioned block storage, so long retention becomes affordable. - **Elasticity / faster ops**: brokers hold much less data, so adding/removing brokers and recovering failed ones is faster because less data has to be replicated. - **Decoupling**: you scale storage (retention) and compute (throughput) independently. ## Key facts to know - Available as production-ready from **Apache Kafka 3.6** (early access in 3.6, GA later). Enabled per-topic via **remote.storage.enable=true**, with the cluster-level **remote.log.storage.system.enable=true**. - Works only on **non-compacted** topics (cleanup.policy=delete), not log-compacted topics. - It is implemented through pluggable interfaces — **RemoteStorageManager** (moves segment bytes) and **RemoteLogMetadataManager** (tracks which segments live where) — so vendors/clouds can supply their own object-store backends. - Reading old data is **transparent** to clients: the consumer asks for an offset, and the broker decides whether to serve it locally or fetch it from the remote tier. ## Edge cases / caveats - Only **closed** segments are offloaded; the active segment being written is never in the remote tier. - Reads from the remote tier have higher latency than local reads, so a workload that frequently re-reads very old data pays an object-store round-trip. - Sizing local retention vs. capacity is an **operational** concern owned by a different topic; here we focus on the internals.

  • Does Tiered Storage work with log-compacted topics?
    No. KIP-405 supports only delete-policy (non-compacted) topics. Compaction rewrites segments in place by key, which is incompatible with immutable offloaded remote segments.
  • Which segment is never offloaded to remote storage?
    The active (currently-being-written) segment. Only closed/rolled segments that are immutable can be uploaded to the remote tier.

saying these in an interview costs you the question

  • Saying tiered storage replaces local disk entirely — recent data and the active segment always stay local.
  • Claiming it works for compacted topics.
  • Confusing it with simply increasing retention.ms — that still needs local disk.
  • Thinking consumers must know whether data is local or remote — fetches are transparent.

context

open as a page

What are the RemoteStorageManager (RSM) and RemoteLogMetadataManager (RLMM) plugins, and how do their responsibilities differ?

level: middleimportance: must knowfreq 60%

basics

~10 s

RemoteStorageManager moves the actual segment bytes to/from remote storage (the data plane). RemoteLogMetadataManager tracks metadata — which segments exist remotely, their offset/epoch ranges, and their state (the control plane).

open as a page

Trace the remote-fetch read path: what happens inside a broker when a consumer requests an offset that lives only in remote storage?

level: seniorimportance: must knowfreq 55%

basics

~20 s

The broker sees the requested offset is below the local log's start, asks the RemoteLogMetadataManager which remote segment holds it, uses the remote offset index to find the byte position, streams the bytes back via the RemoteStorageManager, and returns them in the fetch response — transparently to the consumer.

open as a page

Walk through how a closed log segment gets offloaded to remote storage. What does the broker upload, and when can it delete the local copy?

level: middleimportance: should knowfreq 50%

basics

~20 s

When a segment rolls (closes), the leader broker's RemoteLogManager uploads it plus its indexes to remote storage and records metadata. Only after the upload is confirmed (COPY_SEGMENT_FINISHED) can the local copy be deleted, subject to local retention.

open as a page

Tiered Storage relies on an external object store and a metadata store. As a principal engineer, what consistency and durability invariants must hold, and where can things go wrong?

level: principalimportance: should knowfreq 35%

basics

~20 s

A segment must be durably in the object store before its metadata is marked finished/readable; metadata must be the single source of truth for what's readable. Risks include orphaned objects, eventual-consistency reads, metadata/data divergence, and double-offload across leader changes.

open as a page