📨 Section 15 · Question #28

How are Old Messages Deleted from a Topic

In Apache Kafka, messages are never deleted individually.


🟢 Junior Level

Physical Deletion Mechanics in Kafka

In Apache Kafka, messages are never deleted individually.

Because partition commit logs are designed as high-throughput, sequential, append-only logs, deleting an individual record from the middle of a continuous file would require rewriting gigabytes of binary data on disk and recalculating physical offsets—completely destroying broker I/O performance.

Instead, deletion occurs in entire file blocks called closed Log Segments (LogSegment):

Operating System Filesystem (/var/lib/kafka/data/events-0/):
[Segment 0: 00000.log]   [Segment 1: 00500.log]   [Segment 2: 01000.log (ACTIVE)]
       (1 GB)                   (1 GB)                     (400 MB)
    Offsets 0 - 499          Offsets 500 - 999          Offsets 1000 - ...
          │
          ▼ (retention.ms expired)
   [PURGED ENTIRELY]

Segment Structure on Disk

Every partition segment consists of a coordinated set of files sharing the identical base offset as their filename:

  • 00000000000000000000.log — Physical binary stream of serialized message batches.
  • 00000000000000000000.index — Memory-mapped sparse offset index (offset -> physical byte position).
  • 00000000000000000000.timeindex — Timestamp index (timestamp -> offset).
  • 00000000000000000000.txnindex — Transaction index (tracking aborted transactions under read_committed).

When the broker determines that a segment must be deleted, all associated index and log files are deleted from the operating system atomically as a single unit.


🟡 Middle Level

Lifecycle of Segment Deletion

Deletion of obsolete segments is orchestrated by the internal broker scheduler through the following sequential pipeline:

1. Background thread `kafka-log-retention` awakens every `log.retention.check.interval.ms` (default 5 min).
2. Every closed segment is evaluated against retention criteria:
   - Time: maxTimestamp from .timeindex < (CurrentTime - retention.ms)
   - Size: total partition disk size > retention.bytes (oldest closed segments targeted first).
3. Candidate segment files are renamed:
   00000.log       -> 00000.log.deleted
   00000.index     -> 00000.index.deleted
   00000.timeindex -> 00000.timeindex.deleted
4. Delay timer `log.segment.delete.delay.ms` (default 60 seconds) is initiated.
5. Upon timer expiration, the broker executes an OS unlink() syscall to remove the files.
6. The partition's `LogStartOffset` shifts forward to the base offset of the oldest surviving segment.

Purpose of .deleted Suffix and the 60-Second Delay

The deletion delay (log.segment.delete.delay.ms = 60000) is a vital safety mechanism for distributed concurrent reads:

  1. Protection of In-flight Consumer Reads: When the retention cleaner targets a segment, a consumer or follower replica might already have an open file descriptor reading byte slices via the sendfile(2) system call.
  2. Renaming files with the *.deleted suffix immediately unregisters the segment from metadata lookups for new FetchRequest calls, preventing new consumers from targeting the dying segment.
  3. Existing consumer threads already reading the file complete their active I/O buffer transfers within the 60-second window without suffering sudden I/O crashes or IOException.

Programmatic Deletion via AdminClient.deleteRecords()

Starting with Kafka 0.11, applications can programmatically trigger old message deletion directly from Java code without waiting for automated retention policies:

Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

try (AdminClient adminClient = AdminClient.create(props)) {
    TopicPartition partition = new TopicPartition("orders", 0);
    
    // Define the offset BEFORE which all records must be purged (up to offset 150000)
    Map<TopicPartition, RecordsToDelete> recordsToDelete = Map.of(
        partition, RecordsToDelete.beforeOffset(150000L)
    );

    // Dispatch async request to the partition leader
    DeleteRecordsResult result = adminClient.deleteRecords(recordsToDelete);
    result.all().get();
    
    System.out.println("LogStartOffset for orders-0 successfully advanced to 150000");
}

Important Detail: Calling deleteRecords(beforeOffset) does not split segment files in half. It advances the partition’s LogStartOffset pointer to 150,000. Any closed segments whose highest offset is strictly less than 150,000 are immediately marked as .deleted and unlinked.


🔴 Senior Level

Low-Level Interaction with Linux VFS and Inode Semantics

At the Linux kernel level, Kafka’s segment deletion depends on POSIX file descriptor and inode reference counting:

[LogRetention Thread] ──unlink("00000.log.deleted")──> Removes dentry from directory structure
                                                                   │
                                                                   ▼
                                                     Inode hard link count drops to 0
                                                                   │
        ┌──────────────────────────────────────────────────────────┴───────────────────────────────────────┐
        ▼                                                                                                  ▼
If no open file descriptors (FD count = 0):                                     If consumer is actively reading (FD count > 0):
- OS frees disk blocks immediately.                                              - OS retains disk blocks allocated on disk!
- Space is instantly reclaimed and visible in df -h.                             - Consumer continues streaming bytes via sendfile(2).
                                                                                 - Only when process calls close(fd) does OS free the blocks!

The “Ghost Disk Leak” Gotcha: If a process (such as a stuck backup agent, file watcher, or hung consumer connection) holds an open file descriptor on a deleted segment, rm or .deleted unlinks the directory entry, but the physical disk blocks are not freed. ls shows an empty directory, while df -h continues to report 100% disk usage. This is diagnosed on production brokers using:

lsof +L1 | grep deleted

LogStartOffset Progression and OffsetOutOfRangeException

When segments are physically purged, the lowest addressable offset (LogStartOffset) advances monotonically:

Initial State:
LogStartOffset = 0 ─────────────────────────────> LogEndOffset = 250000

After Purging 3 Segments:
                     LogStartOffset = 150000 ───> LogEndOffset = 250000
                     (Offsets 0..149999 are physically wiped from disk)

If a slow consumer attempts to fetch offset 120,000, the broker checks: \(\text{Requested Offset (120,000)} < \text{LogStartOffset (150,000)}\)

  1. The broker responds with an OFFSET_OUT_OF_RANGE error code.
  2. The consumer client’s reaction is governed by its auto.offset.reset configuration:
    • earliest: Automatically resets its fetch pointer to the new LogStartOffset (150,000) and resumes reading. Messages 120,000–149,999 are permanently skipped.
    • latest: Skips immediately to LogEndOffset (250,000).
    • none: Throws a fatal OffsetOutOfRangeException directly into consumer application code, crashing the consumer thread.

Segment Rolling Triggers and Jitter

For retention to purge old records, active segments must be closed into immutable files (Segment Rollover). An active segment is rolled when any of the following conditions are met:

  1. Physical file size reaches segment.bytes (default: 1 GB).
  2. Elapsed time since the first record was written exceeds segment.ms (or log.roll.ms, default: 7 days).
  3. The offset index reaches segment.index.bytes (default: 10 MB).
Segment Configuration Parameter Default Value Description
segment.bytes 1,073,741,824 (1 GB) Maximum size of a single log segment file before rolling.
segment.ms 604,800,000 (7 days) Maximum time an active segment can remain open before rolling.
segment.index.bytes 10,485,760 (10 MB) Maximum size of .index before rolling the segment.
log.roll.jitter.ms Random (up to log.roll.jitter.ms) Introduces random variance into rollover times to prevent I/O stampedes.

Thundering Herd Rollover Protection: If a topic with 64 partitions is created on Monday at 10:00 AM, all 64 active segments would reach their 7-day segment.ms limit at the exact same second the following Monday. Rolling 64 segments simultaneously triggers concurrent file creation and fsync spikes, saturating disk I/O. Kafka introduces log.roll.jitter.ms to stagger rollovers across randomized intervals.


4 Tricky Questions

1. What happens if a partition contains only a single segment (which is active), and all records in it were written 30 days ago under retention.ms=7 days?

Answer: The data will NOT be deleted.

  1. The active segment (activeSegment) is the sole segment currently receiving new writes for that partition.
  2. In Kafka’s storage architecture, the active segment is strictly immune to deletion under all circumstances, regardless of record age or size limits.
  3. The segment can only be deleted after a rollover occurs—either because a new record triggers segment.bytes, or segment.ms forces a roll. Only once closed into an immutable segment does the retention cleaner evaluate and purge it.

2. How does the broker determine whether a segment is obsolete: by filesystem file creation time, first record timestamp, or last record timestamp?

Answer: The broker relies exclusively on the maximum record timestamp (maxTimestamp) across all records in the segment, recorded in the .timeindex file header:

  1. The broker does not use OS filesystem metadata (ctime/mtime), which can be corrupted or reset during partition reassignment or directory restores.
  2. For every record appended, the broker updates the segment’s maxTimestamp.
  3. A segment is eligible for deletion only if its newest record is older than CurrentTime - retention.ms. If a closed segment contains even a single message within the retention window, the entire segment file is preserved on disk.

3. What happens if AdminClient.deleteRecords() is called with an offset located in the middle of the active segment?

Answer:

  1. The broker accepts the request and immediately updates partition metadata, setting LogStartOffset = 500.
  2. Because the active segment cannot be split or deleted, no physical file deletion occurs.
  3. However, offsets 0 to 499 become logically inaccessible to all consumers. Any fetch request targeting an offset below 500 receives OFFSET_OUT_OF_RANGE.
  4. Physical disk reclamation will occur only after the active segment rolls, becomes closed, and all records inside it fall below the current LogStartOffset.

4. Can segment deletion on the partition leader break lagging follower replicas into Under-Replicated Partitions (URP)?

Answer: Yes, if a follower replica lags beyond the retention window.

  1. Each broker executes segment deletion on its own local kafka-log-retention schedule.
  2. If a follower experiences a network partition for 8 days on a cluster with retention.ms = 7 days, the leader will delete the segments corresponding to the follower’s missing offsets.
  3. When connectivity restores, the follower requests data starting from its last known offset. The leader responds with OFFSET_OUT_OF_RANGE because those segments have been unlinked.
  4. The follower cannot catch up via standard fetch replication, is dropped from the In-Sync Replicas (ISR) list, and transitions into a persistent Under-Replicated state requiring manual administrator intervention or full partition reassignment.

🎯 Interview Cheat Sheet

30-Second Elevator Pitch

“Kafka deletes old messages by pruning entire closed log segments (LogSegment), never by rewriting individual records. The background thread kafka-log-retention runs every 5 minutes checking closed segments. If a segment’s maxTimestamp in .timeindex exceeds retention.ms (or partition size exceeds retention.bytes), files are renamed to *.deleted. After a safety delay of log.segment.delete.delay.ms (60 seconds) to allow in-flight consumer reads via sendfile(2) to complete, the OS unlinks the files and advances LogStartOffset. The active segment is strictly protected and never deleted.”

Key Production Monitoring Metrics

  • kafka.log:type=Log,name=LogStartOffset,topic={topic},partition={id}: Lowest addressable offset in the partition.
  • kafka.log:type=Log,name=NumLogSegments,topic={topic},partition={id}: Total segment files on disk for the partition.
  • Host-level lsof +L1: Checks for unlinked files with open file descriptors holding disk space.

Red Flags (DO NOT Say)

  • ❌ “Kafka locates expired messages in the middle of a file and deletes their bytes.” (Kafka logs are append-only; deletion happens only at whole-segment granularity).
  • ❌ “When Kafka deletes a segment, active consumers crash immediately with FileNotFoundException.” (Files are held open by Linux inode semantics until open file descriptors close).
  • ❌ “Calling deleteRecords() splits the segment file on disk.” (It only updates LogStartOffset; segment files are unlinked in whole blocks).
  • ❌ “The active segment is purged if no new messages arrive for a long time.” (Active segments are never purged; only rolled segments are eligible for deletion).