Kafka Message Loss and Duplication: A Complete Guide

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:

  1. Producer → Broker
  2. Broker itself
  3. 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

  1. Producer retries after send failure
  2. Consumer crashes after processing but before committing offset → repeats consumption on restart
  3. 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

  1. Producer: acks=all + retry mechanism + idempotency
  2. Broker: Multiple replicas + minimum ISR configuration
  3. 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

  1. Loss: Configurable solutions exist
  2. Duplication: Configuration helps reduce but cennot eliminate — business idempotency is required
  3. Optimal Reliability: Producer idempotency + multi-replica setup + manual offset commit + business idempotency

Tags: Kafka message-reliability at-least-once idempotency distributed-systems

Posted on Wed, 19 Aug 2026 16:55:19 +0000 by figo2476