Core Principles
Kafka provides At-Least-Once delivery semantics by default, not exactly-once. This means:
- Messages may be duplicated but will never be lost
- Achieving no-duplication requires business-level idempotency
- Achieving no-loss requires proper configuration and mechanisms
Part One: Message Loss Scenarios and Solutions
Message loss occurs in three stages:
- Producer → Broker
- Broker itself
- Broker → Consumer
Producer-Side Message Loss (Most Common)
Cause:
The producer sends a message but doesn't receive an acknowledgment. Assuming the send failed, it doesn't retry, resulting in message loss. Kafka defaults to asynchronous sending without waiting for confirmation.
Solution:
# Wait for all replicas to sync before acknowledgment
acks=all
# Infinite retries on failure (production recommended)
retries=2147483647
# Ensure message ordering during retries
max.in.flight.requests.per.connection=1
# Enable producer idempotency for automatic deduplication
enable.idempotence=true
Broker-Side Message Loss
Causes:
- Data written to PageCache only; broker crashes before flush → data loss
- Replicas not fully synced; leader fails → data loss
Solution:
# Topic must have at least 3 replicas
replication.factor=3
# Write succeeds only after minimum replicas sync
min.insync.replicas=2
# Disable async flush, use forced sync flush (not recommended in production due to performance impact)
# log.flush.interval.messages=1
# log.flush.interval.ms=1000
Best Practice: Rely on replication, not disk flush, to ensure no data loss.
Consumer-Side Message Loss (Most Problematic)
Cause:
Committing offset before processing message → if processing fails, the message is permanently lost.
Solution:
Process message first → commit offset only after successful processing
# Disable auto-commit
enable.auto.commit=false
Consumer logic:
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
processBusinessLogic(record);
persistToDatabase(record);
consumer.commitSync();
}
}
Part Two: Message Duplication Scenarios and Solutions
Duplication is more common than loss and cannot be completely avoided. The solution lies in business idempotency.
Three Causes of Duplication
- Producer retries after send failure
- Consumer crashes after processing but before committing offset → repeats consumption on restart
- Rebalance triggers duplicate pulls
Ultimate Solution: Business Idempotency (Deduplication with Unique Keys)
Option 1: Database Unique Constraint (Simplest)
- Generate a unique message ID for each message
- Consumer query:
INSERT ... ON DUPLICATE KEY UPDATE
public void consume(Message msg) {
String idempotentKey = msg.getMessageId();
try {
jdbcTemplate.update(
"INSERT INTO consumed_messages(msg_id, content, created_at) VALUES (?, ?, NOW()) " +
"ON DUPLICATE KEY UPDATE content = VALUES(content)",
idempotentKey, msg.getPayload()
);
} catch (DuplicateKeyException e) {
log.info("Duplicate message skipped: {}", idempotentKey);
}
}
Option 2: Distributed Lock (Redis)
- Use message ID as the lock key
- Acquire lock before processing, release after completion
- Duplicate messages are skipped directly
public void consume(Message msg) {
String lockKey = "msg:lock:" + msg.getMessageId();
Boolean acquired = redisTemplate.opsForValue().setIfAbsent(lockKey, "1", Duration.ofMinutes(5));
if (Boolean.FALSE.equals(acquired)) {
log.info("Duplicate message detected, skipping: {}", msg.getMessageId());
return;
}
try {
processMessage(msg);
} finally {
redisTemplate.delete(lockKey);
}
}
Option 3: Kafka Native Idempotency + Transactions (Exactly-Once)
# Enable producer idempotency
enable.idempotence=true
# Producer transaction (ordered across partitions)
transactional.id=transaction-001
producer.initTransactions();
producer.beginTransaction();
producer.send(record);
producer.sendOffsetsToTransaction(offsets, consumerGroupMetadata);
producer.commitTransaction();
Note: Kafka transactions only guarantee no duplication within Kafka itself, not at the bussiness layer. Business idempotency is ultimately required.
Part Three: Production-Ready Configuration
Producer Configuration (No Loss, No Duplication)
acks=all
retries=2147483647
max.in.flight.requests.per.connection=1
enable.idempotence=true
linger.ms=5
batch.size=16384
compression.type=lz4
Consumer Configuration (No Loss, Controlled Duplication)
enable.auto.commit=false
auto.offset.reset=earliest
max.poll.records=500
Broker Configuration (No Loss)
min.insync.replicas=2
unclean.leader.election.enable=false
Summary
Handling Message Loss
- Producer:
acks=all+ retry mechanism + idempotency - Broker: Multiple replicas + minimum ISR configuration
- Consumer: Disable auto-commit, process first then commit
Handling Message Duplication
Kafka cannot eliminate duplication entirely. Business idempotency is mandatory:
- Database unique constraints
- Redis distributed locks
- State machine validation
Key Takeaways
- Loss: Configurable solutions exist
- Duplication: Configuration helps reduce but cennot eliminate — business idempotency is required
- Optimal Reliability: Producer idempotency + multi-replica setup + manual offset commit + business idempotency