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 |
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 |
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 |
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 |
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 |