Practical Guide to Kafka Log Compaction for Event Sourcing Architectures
Intended Reader
This guide is for software engineers, architects, and platform operators who design or maintain event sourcing systems using Apache Kafka as the event backbone. It assumes familiarity with Kafka fundamentals and Java client usage.
Concrete Outcome
By the end of this guide, you will understand when and how to use Kafka log compaction effectively to support event sourcing architectures, including code examples for topic setup, producing, and consuming compacted topics, along with verification, troubleshooting, and operational best practices.
Prerequisites
- Kafka version 2.0 or newer (log compaction is stable and well-supported)
- Java 11+ for client examples
- Basic Kafka producer and consumer knowledge
- Access to a Kafka cluster with Admin API privileges
When and Why Use Kafka Log Compaction for Event Sourcing
Event sourcing stores application state changes as immutable events. Kafka’s append-only logs naturally capture these events in order, but storage and retention policies impact data durability and recovery.
Traditional Kafka time- or size-based retention policies can:
- Delete events prematurely, risking data loss
- Cause uncontrolled storage growth if retention is unlimited
Log compaction solves this by retaining only the latest record per key indefinitely, enabling:
- Efficient storage of current entity state alongside full event histories for critical keys
- Fast state reconstructions without scanning entire event logs
- Combining immutable append logs with a "latest state" snapshot model
Use log compaction when:
- You need to maintain current state for entities without losing event histories
- Your event sourcing relies on replaying events keyed by entity IDs
- Storage optimization is critical for long-lived topics
When not to use log compaction:
- If your use case requires full event history retention without deletion
- For logs where keys are mostly null or events do not represent entity updates
- When precise event ordering of all duplicate keys is mandatory (compaction removes intermediate duplicates)
Alternatives:
- Pure delete-based retention with periodic backups/snapshots
- External state stores for snapshots (Cassandra, RocksDB in conjunction with Kafka Streams)
Trade-offs:
- Log compaction reduces storage but introduces background CPU overhead
- Compacted topics require careful key design
- Consumers must handle potentially missing intermediate events
Understanding Kafka Log Compaction Mechanism
Kafka’s log compaction runs asynchronously, scanning the log segments of configured topics to retain only the latest record for each unique key.
- Key importance: Each record’s key defines identity. Only records with non-null keys participate in compaction.
- Value: The latest value per key is retained; older duplicates are removed over time.
- Null value: Writing a record with a null value for a key acts as a delete marker (tombstone) after compaction.
This differs from time-based retention that deletes records older than a certain age regardless of keys.
Topic configuration to enable compaction:
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic user-events \
--partitions 4 --replication-factor 3 --config cleanup.policy=compact
Setting cleanup.policy=compact marks the topic for compaction. Optionally, you can combine with delete (cleanup.policy=compact,delete) to allow time-based deletion alongside compaction.
Practical Implementation with Java Example
Creating Compacted Topics via Admin API
Using Kafka's AdminClient, create a compacted topic programmatically:
import org.apache.kafka.clients.admin.*;
import java.util.*;
public void createCompactedTopic(AdminClient adminClient) throws Exception {
NewTopic topic = new NewTopic("user-events", 4, (short)3)
.configs(Collections.singletonMap("cleanup.policy", "compact"));
adminClient.createTopics(Collections.singleton(topic)).all().get();
System.out.println("Created compacted topic: user-events");
}
This ensures Kafka periodically compacts the topic.
Producing Events with Meaningful Keys
In event sourcing, keys typically represent entity or aggregate IDs to group related events. See the example below:
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
Properties producerProps = new Properties();
// Configure bootstrap servers, serializers, acks, idempotence, etc.
producerProps.put("bootstrap.servers", "localhost:9092");
producerProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
producerProps.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
producerProps.put("acks", "all");
producerProps.put("enable.idempotence", true);
KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);
String userId = "user-123";
String eventPayload = "{ \"type\": \"UserUpdated\", \"details\": { ... } }";
ProducerRecord<String, String> record = new ProducerRecord<>("user-events", userId, eventPayload);
producer.send(record, (metadata, exception) -> {
if (exception != null) {
System.err.println("Failed to send record: " + exception.getMessage());
} else {
System.out.println("Sent to partition " + metadata.partition() + ", offset " + metadata.offset());
}
});
producer.flush();
producer.close();
How It Works Together
- The topic
user-eventsis compacted, so Kafka keeps only the latest event peruserIdkey. - Multiple updates with the same
user-123key will eventually be compacted to the latest state for that user.
Consuming from a Compacted Topic
Consumers must handle the fact that intermediate historical events could be missing after compaction.
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
Properties consumerProps = new Properties();
consumerProps.put("bootstrap.servers", "localhost:9092");
consumerProps.put("key.deserializer", StringDeserializer.class.getName());
consumerProps.put("value.deserializer", StringDeserializer.class.getName());
consumerProps.put("group.id", "user-event-consumer-group");
consumerProps.put("auto.offset.reset", "earliest");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
consumer.subscribe(Collections.singletonList("user-events"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("Processing userId=%s, offset=%d, event=%s%n", record.key(), record.offset(), record.value());
// Apply event to rebuild or update entity state
}
consumer.commitSync();
}
This consumer ensures it sees the latest events after compaction. Use state snapshotting to optimize startup/recovery.
Verification and Expected Results
To verify log compaction:
- Produce multiple events with the same key (e.g.,
user-123) but different values. - Use the Kafka console consumer or your application consumer to consume the topic.
- Observe that after compaction:
- Only the newest event for
user-123remains visible. - Older duplicates are removed from Kafka log segments.
Example CLI consumer to verify:
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic user-events --from-beginning --property print.key=true
You should see only one event per key (latest value) after compaction runs.
Production Failure Modes and Troubleshooting
Common Issues
- Unexpected missing events: Ensure keys are set correctly (non-null and stable).
- Compaction not happening: Verify
cleanup.policy=compact,log.cleaner.enable=truein broker configs, and cleaner thread activity. - High CPU load from cleaner threads: Tune
log.cleaner.threads,min.cleanable.dirty.ratio, and backoff settings. - Consumer sees fewer events than expected: This is normal after compaction; consumers must handle potential gaps.
Troubleshooting Steps
- Check Kafka broker logs for cleaner thread errors.
- Monitor JMX metrics such as
kafka.log:type=LogCleanerManager. - Confirm topic configs via
kafka-topics.sh --describe. - Test with a dedicated compacted topic with controlled inputs.
Security Considerations
- Use Kafka's TLS and SASL mechanisms to secure producer and consumer connections.
- Restrict topic creation and configuration permissions to prevent unauthorized cleanup policy changes.
- Monitor audit logs for producer and consumer activity on compacted topics.
Performance and Operational Safeguards
- Configure sufficient cleaner threads (
log.cleaner.threads) to handle compaction without impacting overall throughput. - Combine compaction with retention policies (
compact,delete) to limit storage growth. - Partition topics thoughtfully (e.g., one partition per entity group) for parallel compaction and consumption.
- Monitor disk usage, lag, and throughput metrics continuously.
Limitations
- Log compaction does not guarantee atomic state snapshots across partitions.
- Intermediate events may be lost, requiring consumers to support idempotency and state rebuilding.
- Compaction is asynchronous and may lag behind event production.
Summary
Kafka log compaction is a critical tool for event sourcing architectures to retain up-to-date entity states efficiently. By carefully designing topics with compacted policy, producing keyed events, and consuming with awareness of the compaction process, you can build resilient and scalable event-driven systems.
FAQ
Can log compaction delete old events?
No, log compaction does not delete all old events indiscriminately; it keeps the latest record per key indefinitely, removing only obsolete duplicates.
Are messages without keys compacted?
No, only records with non-null keys participate in log compaction. Records with null keys are excluded and managed by retention policies instead.
How should consumers handle compacted topics?
Consumers must be prepared that intermediate events may be missing due to compaction and should use idempotent processors and snapshotting as needed.
Can log compaction be combined with time-based retention?
Yes, setting cleanup.policy to compact,delete allows both compaction and time-based deletion to operate.
What are good key selection criteria?
Choose stable, unique identifiers representing aggregates or entities to ensure meaningful and effective compaction.
Sources and further reading
- Apache Kafka Official Documentation – Log Compaction
- Kafka AdminClient API Reference
- Best Practices for Kafka Producer and Consumer
