Implementing Log Compaction

1. Understanding Compaction Concept

Concept Description Detail
Compaction Keep latest value per key Not time-based
Use Case Changelogs, state snapshots Materialized view
Guarantee Latest record per key retained Older removed

Example: Compacted topic concept

k1=v1, k2=v2, k1=v3, k2=v4  →  compacted →  k1=v3, k2=v4

2. Enabling Log Compaction

Config Value Detail
cleanup.policy compact Per topic
Requires Keys Null key not allowed Else error
log.cleaner.enable true (broker) Default on

Example: Enable compaction

kafka-configs.sh --alter --entity-type topics --entity-name user-state \
  --add-config cleanup.policy=compact --bootstrap-server localhost:9092

3. Understanding Compaction Process

Phase Action Detail
Dirty Ratio Cleaner selects logs to clean Threshold-based
Build Map Offset map of latest keys In memory
Rewrite Copy survivors to new segment Drop old

Example: Cleaner status in logs

grep "Cleaner" /var/lib/kafka/logs/log-cleaner.log

4. Understanding Tombstone Messages

Concept Description Detail
Tombstone Record with null value Delete marker
Effect Removes key on compaction After retention
delete.retention.ms How long tombstone kept Default 1d

Example: Send a tombstone

// null value deletes the key from the compacted topic
producer.send(new ProducerRecord<>("user-state", userId, null));

5. Configuring Compaction Lag

Config Description Default
min.compaction.lag.ms Min age before record compacted 0
max.compaction.lag.ms Max delay before compaction Long.MAX
Use Case Keep recent history readable Consumer catch-up

Example: Set compaction lag

kafka-configs.sh --alter --entity-type topics --entity-name user-state \
  --add-config min.compaction.lag.ms=3600000 \
  --bootstrap-server localhost:9092

6. Setting Delete Retention

Config Description Default
delete.retention.ms Tombstone visibility window 86400000 (1d)
Consumer Window Must consume deletes within Else missed
Tradeoff Longer = more disk Safety vs cost

Example: Extend delete retention

kafka-configs.sh --alter --entity-type topics --entity-name user-state \
  --add-config delete.retention.ms=172800000 \
  --bootstrap-server localhost:9092

7. Configuring Segment Size

Config Description Detail
segment.bytes Smaller = compact sooner Active never cleaned
segment.ms Force roll by time Frees active
Tradeoff Smaller = more files/overhead Balance

Example: Smaller segments for compaction

kafka-configs.sh --alter --entity-type topics --entity-name user-state \
  --add-config segment.bytes=52428800,segment.ms=600000 \
  --bootstrap-server localhost:9092

8. Setting Min Cleanable Ratio

Config Description Default
min.cleanable.
dirty.ratio
Dirty fraction to trigger clean 0.5
Lower Compact more aggressively More CPU/IO
Higher Less frequent, more disk Less overhead

Example: Aggressive cleaning

kafka-configs.sh --alter --entity-type topics --entity-name user-state \
  --add-config min.cleanable.dirty.ratio=0.1 \
  --bootstrap-server localhost:9092

9. Configuring Compaction Thread

Config Description Default
log.cleaner.threads Number of cleaner threads 1
log.cleaner.
io.max.bytes.per.second
Cleaner IO throttle Double.MAX
Scaling More threads for many topics Broker-level

Example: More cleaner threads

log.cleaner.threads=4
log.cleaner.io.max.bytes.per.second=10485760

10. Setting Compaction Buffer

Config Description Default
log.cleaner.
dedupe.buffer.size
Memory for offset map 134217728
log.cleaner.
io.buffer.size
Read/write IO buffer 524288
Sizing Larger maps more keys/pass Fewer passes

Example: Larger dedupe buffer

log.cleaner.dedupe.buffer.size=536870912  # 512MB

11. Monitoring Compaction Progress

Metric Description Source
max-clean-time-secs Longest cleaning time LogCleaner JMX
cleaner-recopy-percent Fraction recopied Efficiency
max-buffer-utilization Dedupe buffer use Sizing signal

Example: Cleaner JMX bean

# kafka.log:type=LogCleanerManager,name=max-dirty-percent