skip to content

What is Kafka Tiered Storage (KIP-405), and how does it change the economics and architecture of using Kafka as a long-retention backbone for analytics?

level: seniorimportance: should knowfreq 40%

answer

  1. local tier (recent) + remote tier (old) in object store
  2. RemoteStorageManager / RemoteLogManager
  3. decouple retention from broker disk
  4. transparent remote fetch for old offsets
  5. GA in Kafka 3.9; compaction caveat

basics

~20 s

Tiered Storage (KIP-405) lets a Kafka broker offload older log segments to cheap remote object storage (like S3) while keeping recent data on local disk. This makes long retention affordable and lets brokers store far more history without huge local disks.

solid answer

~50 s

KIP-405 introduces **Tiered Storage**: each topic partition's log is split into a **local tier** (recent segments on broker disk) and a **remote tier** (older, closed segments offloaded to object storage like S3/GCS via a pluggable `RemoteStorageManager`). A `RemoteLogManager` tracks remote segment metadata. Reads of old offsets transparently fetch from remote storage; producers and recent consumers hit local disk as before. This **decouples retention from local disk capacity**: you can keep weeks/months of history cheaply without scaling broker storage, and brokers recover/rebalance faster because less data lives locally. For analytics it means Kafka itself becomes a durable, replayable long-term log — you can reprocess months of events into the lakehouse, backfill new consumers, or rebuild derived tables without a separate archival system. It went GA in Apache Kafka 3.9. Caveats: compacted topics weren't initially supported, remote reads add latency, and you still pay object-storage and retrieval costs.

go deeper

for a junior

Know KIP-405 offloads old Kafka data to cheap object storage so you can keep more history.

for a middle

Explain local vs remote tier, transparent old-offset fetch, and decoupling retention from disk.

for a senior

Discuss RemoteStorageManager, replay/backfill use cases, latency/cost caveats, and the compaction limitation.

for a principal

Reason about Kafka-as-system-of-record trade-offs, object-store dependency, and when tiered storage vs. lakehouse archival is right.

## The problem before KIP-405 A Kafka topic's data lives as **log segments** on the **local disks** of brokers, replicated across the cluster. Retention was bounded by how much local disk you were willing to buy and replicate. Keeping months of history meant huge, expensive, replicated SSDs — and big disks make broker recovery and rebalancing slow because all that data must be re-replicated when a broker is added or replaced. ## What KIP-405 does **Tiered Storage** splits each partition's log into two tiers: - **Local tier**: the most recent segments, on broker disk (low-latency for producers and tailing consumers). - **Remote tier**: older, already-closed (rolled) segments **offloaded to remote object storage** — S3, GCS, Azure Blob, HDFS — through a pluggable **`RemoteStorageManager` (RSM)** interface. A **`RemoteLogManager`** on the broker handles copying eligible segments to remote storage and tracks **remote-log-segment metadata** (which offsets live where), persisted via a `RemoteLogMetadataManager` (default implementation uses an internal Kafka topic). `local.retention.ms`/`local.retention.bytes` control how much stays local; overall `retention.ms`/`retention.bytes` can now be far larger, with the bulk sitting cheaply in object storage. ## How reads work - Producers and **recent** consumers read from the local tier — no change in hot-path performance. - A consumer reading **old** offsets (a backfill, a new derived job, replay after an incident) triggers a **transparent remote fetch**: the broker streams the needed segment from object storage. This adds latency and retrieval cost but requires no client changes. ## Why it changes analytics architecture 1. **Cheap long retention**: months of event history at object-storage prices instead of replicated SSD prices. Kafka becomes a viable **system of record / replayable log**, not just a transient buffer. 2. **Replay and backfill**: you can re-run ingestion to rebuild lakehouse tables, onboard a new consumer from the beginning of history, or reprocess after a bug — without a separate archive. 3. **Faster operations**: less local data means quicker broker startup, recovery, and rebalancing, since remote segments don't need re-replication. 4. **Decoupling storage from compute**: scale retention independently of broker disk. ## Caveats and edge cases - **GA in Apache Kafka 3.9** (early access in 3.6); production hardening matured over the 3.x line. - **Compacted topics** were **not supported** in the initial design (only delete-retention topics offload) — check your version. - **Latency/cost on cold reads**: remote fetches are slower and incur per-request and egress charges; a workload that constantly reads old data may cost more, not less. - **Operational dependency**: your Kafka durability now also depends on the object store's availability and the correctness of the RSM plugin. - It is **not** a replacement for a lakehouse table format — Kafka remains an append-only log; you still sink into Iceberg/Parquet for query-optimized analytics. Tiered storage just makes Kafka a durable, replayable source for that pipeline.

  • Does Tiered Storage replace sinking data into an Iceberg/Parquet lakehouse?
    No. Kafka remains an append-only log optimized for streaming, not columnar analytical queries. Tiered storage makes Kafka a cheap, durable, replayable source; you still sink into a query-optimized table format (Iceberg/Parquet) for efficient analytics. The two are complementary.
  • What new failure/cost considerations does Tiered Storage introduce?
    Durability now depends on the remote object store's availability and the RemoteStorageManager plugin's correctness; cold reads of old offsets add latency and incur per-request/egress costs, so read-heavy historical workloads can get slower and more expensive, not cheaper.

saying these in an interview costs you the question

  • Saying it lets Kafka replace the lakehouse/Iceberg — it complements, doesn't replace, query-optimized storage.
  • Claiming remote reads are as fast as local — cold fetches add latency and cost.
  • Assuming all topic types (including compacted) are supported in every version.
  • Thinking recent/produce-path performance degrades — hot data stays local.

context