skip to content

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