How do you do capacity and cost planning when offloading Kafka segments to remote object storage, and what hidden costs should you account for?
answer
- local = ingest × local.retention × RF
- remote = ingest × (total−local) × compression, single copy (no RF)
- hidden: PUT/GET/LIST requests + egress + metadata topic
- segment.bytes big ⇒ fewer PUTs; chunk small ⇒ more GETs
- replay/lag converts cheap storage into expensive reads
basics
~20 sModel 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 sLocal 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
Knows local disk shrinks and remote storage holds the bulk of retained data more cheaply.
Can compute local vs remote bytes from retention and ingest and knows remote keeps one copy.
Accounts for per-request and egress costs and tunes segment/chunk sizes against them.
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.