skip to content

How do you do capacity and cost planning when offloading Kafka segments to remote object storage, and what hidden costs should you account for?

level: principalimportance: should knowfreq 30%

answer

  1. local = ingest × local.retention × RF
  2. remote = ingest × (total−local) × compression, single copy (no RF)
  3. hidden: PUT/GET/LIST requests + egress + metadata topic
  4. segment.bytes big ⇒ fewer PUTs; chunk small ⇒ more GETs
  5. replay/lag converts cheap storage into expensive reads

basics

~20 s

Model local disk from local.retention and ingest rate; model remote cost from total retention × ingest × replication-of-1 (remote stores one copy). Add hidden costs: API request charges (PUT/GET/LIST), egress on remote fetches, and the metadata topic.

solid answer

~50 s

Local disk is sized by local.retention × ingest rate × replication factor across partitions — keep it small so most data lives remotely. Remote storage cost is total_retention × ingest_rate × compression_ratio, and crucially remote keeps a single copy (no Kafka RF multiplier), so offloading slashes both disk and the cross-AZ replication footprint of cold data. But object stores bill more than GB-months: per-request charges (PUT per segment upload + multipart parts, GET/LIST on fetch), data-transfer/egress when consumers replay cold data across AZ/region/internet, and the __remote_log_metadata topic's own storage. Tune segment.bytes (bigger segments ⇒ fewer PUTs but coarser deletes) and chunk size (smaller ⇒ more GETs but less over-fetch). Hot-path reads stay local; only lagging/replay consumers incur remote GET + egress, so model your replay patterns. Net: tiered storage trades cheap, deep retention against request/egress costs and remote-fetch latency.

go deeper

for a junior

Knows local disk shrinks and remote storage holds the bulk of retained data more cheaply.

for a middle

Can compute local vs remote bytes from retention and ingest and knows remote keeps one copy.

for a senior

Accounts for per-request and egress costs and tunes segment/chunk sizes against them.

for a principal

Builds a full TCO model (RF-free remote, request/egress, metadata overhead, replay patterns) and drives placement and config decisions across clusters.

**Goal.** Decide local disk size, remote storage budget, and configs (retention, segment size, chunk size) so the cluster meets retention and replay SLAs at minimum cost. **1. Local disk model.** Per partition: `local_bytes ≈ ingest_rate_per_partition × local.retention.ms` (or directly `local.retention.bytes`). Multiply by partitions per broker and by **replication factor** (every replica stores the local log). This is your expensive SSD/disk tier — minimize `local.retention` to just cover consumer lag, replay window, and replica catch-up. **2. Remote storage model.** `remote_bytes ≈ ingest_rate × (retention.ms − local.retention.ms) × compression_ratio`. Key insight: the remote store keeps **one copy** of each segment — Kafka's replication factor does **not** multiply remote storage (the RSM uploads once; the object store provides its own durability/redundancy internally). So moving cold data remote removes both the local-disk RF multiplier *and* the cross-AZ replication traffic for that cold data. This is the dominant saving. **3. Hidden / non-obvious costs:** - **Per-request (API) charges.** Object stores bill per PUT/GET/LIST. Uploads generate PUTs (and one per multipart part); reads on the cold path generate GETs (and LIST for discovery). High partition counts × small segments ⇒ many PUTs. Tune **segment.bytes** larger to cut upload request count (trade-off: coarser remote deletion granularity, larger minimum fetch). - **Data transfer / egress.** Reading cold data may cross AZ, region, or the internet — egress is often the biggest surprise. Same-AZ/same-region access and gateway/VPC endpoints reduce this. Replay-heavy workloads (backfills, reprocessing) can dominate. - **Metadata topic.** `__remote_log_metadata` consumes broker disk and replication like any topic; size its partitions/retention for the volume of remote segments. - **Early-delete waste.** Deleting before the typical object-store minimum-storage-duration (if the backend/storage class has one) can incur minimum-duration charges; pick the right storage class. - **Fetch cache / chunk size.** Smaller chunk size reduces over-fetch per ranged GET but increases GET count; larger chunks do the opposite. Some plugins keep a local fetch cache (extra local disk) to dampen repeated GETs. **4. Latency / SLA trade-off.** Hot reads inside `local.retention` are local (low latency). Cold reads pay RSM fetch latency + object-store round trips + egress. Capacity planning must include *who reads cold data and how often*: a consumer that constantly lags past local.retention will continuously hit remote, inflating GET/egress cost and tail latency. **5. Putting it together (worked sketch).** For a topic at 10 MB/s ingest, RF=3, 30-day total retention, 6h local retention, ~0.5 compression: - Local: 10 MB/s × 6h × 3 ≈ ~648 GB across replicas (vs. ~78 TB if all 30 days were local × RF3 — the whole point). - Remote: 10 MB/s × ~29.75d × 0.5 ≈ ~12.8 TB, single copy. - Plus: PUTs ≈ (bytes/segment.bytes) per upload; GETs/egress driven by replay volume; metadata topic overhead. **Principal-level framing:** present it as a TCO comparison (local-only RF3 disk vs. tiered local + remote single-copy + request/egress), choose segment/chunk sizes to balance request count vs. fetch granularity, place buckets in-region with VPC endpoints to kill egress, and govern replay patterns since they convert 'cheap storage' into 'expensive reads'.

  • Why does tiered storage remove the replication-factor multiplier from cold-data cost?
    Kafka's RF replicates the local log across brokers, but the RSM uploads each segment to the object store once; the object store provides its own internal redundancy. So remote bytes are stored a single time regardless of Kafka RF, eliminating both the RF disk multiplier and cross-AZ replication of cold data.
  • A team complains tiered storage made their cluster cheaper on disk but the cloud bill barely dropped. What would you investigate?
    Per-request and egress charges: many small segments/high partition count inflating PUT/GET counts, cross-AZ or cross-region bucket placement causing egress on fetches, and replay-heavy or constantly-lagging consumers repeatedly reading cold data. Fix by enlarging segment.bytes, co-locating the bucket in-region with VPC endpoints, tuning chunk size, and reducing unnecessary backfills.
  • How does segment.bytes interact with remote request cost?
    Each uploaded segment is one (multipart) PUT and remote deletion is per-segment. Larger segment.bytes means fewer PUTs and less metadata, but coarser deletion granularity and a larger minimum unit to fetch; smaller segments mean more requests but finer control.

saying these in an interview costs you the question

  • Assuming remote storage cost is multiplied by Kafka's replication factor (it's a single copy).
  • Modeling only GB-month storage and ignoring per-request (PUT/GET/LIST) and egress charges.
  • Forgetting that lagging/replay consumers turn cheap cold storage into repeated expensive remote reads.
  • Overlooking the __remote_log_metadata topic's own footprint.

context