📨 Section 15 · Question #27

What is Retention Policy

The Retention Policy in Apache Kafka is a set of configuration rules governing how long and in what volume brokers retain messages inside topic partitions, and under what condit...


🟢 Junior Level

What is Retention Policy in Plain English

The Retention Policy in Apache Kafka is a set of configuration rules governing how long and in what volume brokers retain messages inside topic partitions, and under what conditions obsolete data is cleaned up or deleted.

Unlike traditional message queues (RabbitMQ, IBM MQ, ActiveMQ) where a message is permanently purged immediately after a consumer acknowledges it (ack), Kafka retains all messages, regardless of how many consumers have read them or whether any consumer has connected at all.

Partition Commit Log on Disk:
[Segment 1: 8 days ago] [Segment 2: 4 days ago] [Segment 3: 1 day ago] [Segment 4: Active Segment]
         │
         ▼ (retention.ms = 7 days)
[CANDIDATE FOR DELETION] ──────> Purged completely by background broker thread

Core Log Cleanup Policies (cleanup.policy)

The topic configuration parameter cleanup.policy dictates the fundamental mechanism for freeing disk space:

Cleanup Policy Description Typical Use Cases
delete (Default) Purges records once they exceed time thresholds (retention.ms) or partition size ceilings (retention.bytes). Event streams, audit trails, server metrics, user telemetry.
compact Preserves only the latest value for each unique message key (key). Older records matching the same key are deleted. Entity state changelogs, user profiles, current account balances, cache stores.
compact,delete Hybrid mode: performs key deduplication while also discarding records that have not received updates within retention.ms. Temporary user sessions, transient cache invalidation changelogs.

Essential Configuration Parameters

# 1. Message time-to-live (milliseconds)
retention.ms=604800000          # 7 days (7 * 24 * 60 * 60 * 1000)

# 2. Maximum disk storage limit per partition (bytes)
retention.bytes=10737418240     # 10 GB per partition (NOT across the entire topic!)

# 3. Cleanup policy
cleanup.policy=delete

Kafka CLI Management

# Modify data retention on an existing topic to 3 days (259,200,000 ms)
kafka-configs.sh --bootstrap-server localhost:9092 \
  --entity-type topics --entity-name orders \
  --alter --add-config retention.ms=259200000

# Inspect active retention settings for a topic
kafka-configs.sh --bootstrap-server localhost:9092 \
  --entity-type topics --entity-name orders --describe

🟡 Middle Level

Log Segments and Deletion Mechanics

Kafka never deletes individual messages from the middle of a continuous file. A partition log on disk is partitioned into immutable fixed-size files called log segments (.log), accompanied by offset index files (.index) and timestamp index files (.timeindex).

Partition Directory: /var/lib/kafka/data/orders-0/
├── 00000000000000000000.log       (Closed segment, 1 GB, created 10 days ago) -> PURGED
├── 00000000000000000000.index
├── 00000000000000000000.timeindex
├── 00000000000001948520.log       (Closed segment, 1 GB, created 3 days ago)   -> RETAINED
├── 00000000000001948520.index
├── 00000000000003891440.log       (ACTIVE segment, 400 MB, open for writes)    -> NEVER PURGED

Rules Governing Segment Deletion

  1. Deletion Granularity is the Entire Segment: The background broker thread kafka-log-retention wakes up every log.retention.check.interval.ms (default: 5 minutes / 300,000 ms) to inspect closed segments.
  2. A Segment is Deleted Only When ALL Messages Inside It Expire: The broker checks the maximum timestamp (maxTimestamp) recorded in the .timeindex of a closed segment. If: \(\text{CurrentTime} - T_{\text{segment\_max}} > \text{retention.ms}\) the segment is marked for deletion (renamed with .deleted extension) and physically unlinked from disk after log.segment.delete.delay.ms (default: 1 minute).
  3. The Active Segment is NEVER Deleted: The active segment currently receiving writes from producers is protected against deletion, even if its oldest messages or physical size exceed retention.ms or retention.bytes.

The Critical Trap: retention.bytes Applies PER PARTITION

A common operational disaster is assuming that retention.bytes limits the aggregate storage consumed by an entire topic across the cluster:

\[\text{Max Topic Disk Usage} = \text{retention.bytes} \times \text{Partition Count} \times \text{Replication Factor}\]

Production Outage Example: An engineer sets retention.bytes = 50 GB, expecting the topic to consume 50 GB of disk. The topic has 32 partitions and a replication factor of 3. Total cluster disk footprint becomes: \(50\text{ GB} \times 32 \times 3 = 4,800\text{ GB (4.8 TB)!}\) Broker disks fill to 100%, causing brokers to halt writes and enter read-only panic state.


🔴 Senior Level

Log Compaction Internals: Clean vs Dirty Log

For topics configured with cleanup.policy=compact, background threads in the LogCleaner pool (kafka.log.LogCleaner) prune redundant historical states:

                      LOG COMPACTION SEGMENT STRUCTURE
┌───────────────────────────────────────────────┬───────────────────────────────┐
│                 CLEAN SECTION                 │         DIRTY SECTION         │
│ (Already compacted closed segments)           │ (Uncompacted closed segments) │
│ Each key appears exactly once                 │ May contain duplicate keys    │
└───────────────────────────────────────────────┴───────────────────────────────┘
                                                ▲                               ▲
                                                │                               │
                                          firstDirtyOffset                activeSegment
                                                                    (Never compacted)
  1. SkimpyOffsetMap Memory Footprint: The LogCleaner thread scans the dirty log section and builds an off-heap hash table named SkimpyOffsetMap (configured via log.cleaner.dedupe.buffer.size, default: 128 MB). It maps a 128-bit MD5 hash of each message key to its latest observed offset.
  2. Compaction Pass: The cleaner reads through clean and dirty sections sequentially. If a record’s key has a higher offset stored in SkimpyOffsetMap, the older record is dropped. Surviving records are written to a replacement .clean segment, which is then atomically swapped with the old segment files.
  3. Tombstone Protocol (delete.retention.ms): To delete a key completely, a producer publishes a Tombstone record (key != null, value == null):
    • The tombstone is retained in the log to ensure downstream consumers read the deletion event.
    • Once delete.retention.ms (default: 24 hours / 86,400,000 ms) has elapsed, subsequent compaction runs permanently purge the tombstone marker itself.

KIP-405: Tiered Storage Architecture

Starting with Kafka 3.x, Tiered Storage (KIP-405) decouples storage capacity from local broker compute:

[Kafka Broker]
     │
     ├── Fast Local Storage (NVMe / SSD):
     │   local.retention.ms = 86400000 (24 hours) / local.retention.bytes
     │   (Serves real-time consumers with sub-millisecond zero-copy latency)
     │
     └── Background Offload via RemoteStorageManager (RSM):
         retention.ms = 31536000000 (1 year)
         [Cloud Object Storage: AWS S3 / GCS / Azure Blob / MinIO]
         (Retains cold closed segments at commodity object-storage pricing)
  • Local disks enforce strict, aggressive local.retention.ms (e.g., 24 hours), keeping expensive broker SSDs small and lean.
  • Historical consumers (backfill, audit, analytics, ML training) seamlessly fetch older segments directly from cloud object storage via standard Kafka consumer protocols without impacting broker page-cache memory.

4 Tricky Questions

1. What happens if a producer sends a batch with an erroneous future timestamp (e.g., year 2099) when message.timestamp.type=CreateTime?

Answer: This causes a permanent retention deletion lock (disk leak):

  1. Under CreateTime, the broker records the producer-supplied timestamp directly into the segment .timeindex.
  2. The closed segment records its maxTimestamp as the year 2099.
  3. When evaluating segment deletion eligibility, the broker calculates: \(\text{CurrentTime} - \text{segment.maxTimestamp} > \text{retention.ms}\) Because the result is negative, this condition will never evaluate to true during normal operational lifespans.
  4. This segment and all preceding segments in the partition are blocked from deletion based on time, accumulating indefinitely until disk space is exhausted (unless bounded by retention.bytes).
  5. Mitigation: Enforce broker-side timestamp validation using log.message.timestamp.difference.max.ms to reject records whose timestamps diverge excessively from broker wall-clock time.

2. A consumer was offline for 48 hours for maintenance. The topic uses cleanup.policy=compact with default delete.retention.ms=86400000 (24 hours). What bug occurs when the consumer restarts?

Answer: Silent State Desynchronization (State Inconsistency):

  1. During the consumer’s 48-hour outage, a tombstone record (key="user_123", value=null) was published to delete user 123.
  2. Because the downtime exceeded delete.retention.ms (48 hours > 24 hours), the LogCleaner completed a compaction pass and physically erased the tombstone marker from disk.
  3. When the consumer comes back online, there is no record of the deletion ever occurring in the log segment.
  4. The consumer reads remaining offsets, and in its local relational database or cache, user 123 remains active permanently.
  5. Architectural Rule: delete.retention.ms must always be configured to exceed the maximum guaranteed consumer recovery SLA window (e.g., 7 days).

3. Can Kafka delete an active segment if partition size exceeds retention.bytes by 10x?

Answer: No. An active segment (activeSegment) is strictly immune to deletion under all circumstances. Even if retention.bytes = 100 MB and the active segment has accumulated 1 GB of data, the broker will never delete it while it remains open for writes. The segment must first be rolled (closed) into an immutable segment via segment.bytes or segment.ms. Only after a segment is closed does the retention coordinator evaluate it for deletion.

4. What is the operational danger of having log.roll.ms set much higher than retention.ms on low-throughput topics?

Answer:

  • retention.ms specifies how long a closed segment is retained before deletion.
  • log.roll.ms (or segment.ms) dictates the maximum duration an active segment stays open before being forced to roll.
  • The Low-Throughput Trap: If a quiet topic receives only a few messages a day, and segment.bytes is 1 GB while segment.ms defaults to 7 days, an active segment stays open for 7 days. Even if retention.ms is configured for 24 hours, data will not be purged until: \(\text{Effective Retention} = \text{segment.ms} + \text{retention.ms} = 7\text{ days} + 1\text{ day} = 8\text{ days}\) The application expects data to disappear within 24 hours for compliance/GDPR reasons, but data remains readable for over a week.

🎯 Interview Cheat Sheet

30-Second Elevator Pitch

“Retention Policy governs Kafka log lifecycle management via cleanup.policy. Under delete (default), closed segments are pruned once they exceed retention.ms or retention.bytes. Under compact, Kafka retains only the latest record for each key using background LogCleaner threads, removing keys permanently via tombstones (value=null) after delete.retention.ms. Crucially, retention operates on entire closed log segments, never on individual records, and active segments are never deleted. Modern Kafka 3.x architectures leverage Tiered Storage (KIP-405) to maintain short retention on local NVMe drives (local.retention.ms) while archiving cold segments in object storage (S3) for long-term replay.”

Key Production Monitoring Metrics

  • kafka.log:type=Log,name=Size,topic={topic},partition={id}: Partition disk consumption in bytes.
  • kafka.log:type=LogCleanerManager,name=max-dirty-percent: Percentage of uncompacted dirty data. Values persistently $> 0.5$ indicate the cleaner is falling behind.
  • Host-level disk_used_percent: Broker disk alerts (triggering alarms before reaching 85% capacity).

Red Flags (DO NOT Say)

  • ❌ “retention.bytes limits the total storage used by the entire topic across all brokers.” (retention.bytes is enforced per individual partition).
  • ❌ “Kafka inspects each individual message and deletes it the exact millisecond its retention.ms expires.” (Pruning occurs at the segment file level during retention check sweeps).
  • ❌ “The active segment is purged if the broker hard drive reaches 100% capacity.” (Active segments are never purged; full disks crash broker write threads).
  • ❌ “Log compaction erases deleted keys immediately upon receiving a tombstone.” (Tombstones remain in the log until delete.retention.ms expires to allow consumers to observe deletions).