Skip to content
advancedPhase 47 · Messaging

Exactly Once

Achieve exactly-once semantics with idempotent producers and transactions.

45m
0 problems
Topic Progress0%

How it Works

Exactly-Once: How it Works

Components

Exactly-Once = Idempotent Producer + Transactional Consumer

1. Idempotent Producer:
   - Unique sequence ID per message
   - Broker deduplicates retries
   - Prevents producer-caused duplicates

2. Transactional Consumer:
   - Read committed isolation
   - Atomic processing + offset commit
   - Prevents consumer-caused duplicates

Kafka Implementation

// 1. Idempotent Producer
props.put("enable.idempotence", "true");
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
props.put("max.in.flight.requests.per.connection", 5);

// 2. Transactional Producer
props.put("transactional.id", "my-transaction");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();

// 3. Consumer with isolation
props.put("isolation.level", "read_committed");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);

// 4. Transactional processing
producer.beginTransaction();
try {
    // Process message
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        // Process
        String result = processMessage(record.value());
        
        // Produce output
        producer.send(new ProducerRecord<>("output-topic", result));
    }
    
    // Commit offset atomically
    producer.sendOffsetsToTransaction(offsets, consumerGroupId);
    
    producer.commitTransaction();
} catch (Exception e) {
    producer.abortTransaction();
}

Guarantees

1. No duplicates from producer retries
   - Broker tracks sequence numbers
   - Rejects duplicate sequence

2. No duplicates from consumer processing
   - Atomic: process + commit offset
   - If crash, transaction aborted, reprocessed

3. No message loss
   - Transactional delivery
   - Committed messages only

Transactional Messaging

Transactional Messaging

Transaction Boundaries

1. Begin Transaction
2. Process messages (read from input)
3. Produce messages (write to output)
4. Commit consumer offsets
5. Commit transaction (atomic)

All or nothing!

Transaction States

Transaction Lifecycle:

Init → Begin → Process → Commit/Abort
  ↓       ↓        ↓          ↓
Ready  Active   Active    Complete/Abort

Exactly-Once Pipeline

class ExactlyOnceProcessor:
    def __init__(self, kafka_config):
        self.producer = KafkaProducer(**kafka_config)
        self.producer.init_transactions()
    
    def process(self, consumer, input_topic, output_topic):
        self.producer.beginTransaction()
        
        try:
            records = consumer.poll(100)
            
            for record in records:
                # Process
                result = self.transform(record.value())
                
                # Produce to output
                self.producer.send(output_topic, result)
            
            # Commit offsets atomically
            self.producer.sendOffsetsToTransaction(
                self.get_offsets(records),
                consumer.group_metadata()
            )
            
            self.producer.commitTransaction()
            
        except Exception as e:
            self.producer.abortTransaction()
            raise

Transaction Isolation Levels

Level Description
read_uncommitted See all messages (default)
read_committed Only see committed messages

Best Practices

  1. Use unique transactional.id per producer instance
  2. Set appropriate isolation.level for consumers
  3. Handle transaction timeouts gracefully
  4. Monitor transaction metrics
  5. Test failure scenarios thoroughly

Idempotent Producers

Idempotent Producers

How Idempotent Producers Work

Producer assigns sequence to each message:

Message 1: seq=1
Message 2: seq=2
Message 3: seq=3

Broker tracks last sequence per producer:

Receive seq=1 → Accept, last=1
Receive seq=2 → Accept, last=2
Receive seq=2 → Reject (duplicate!)
Receive seq=3 → Accept, last=3

Configuration

// Enable idempotence
props.put("enable.idempotence", "true");

// Required settings
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
props.put("max.in.flight.requests.per.connection", 5);

// What it does:
- Assigns sequence numbers
- Broker deduplicates by sequence
- Prevents retries causing duplicates

Sequence Number Tracking

# Broker-side tracking concept
producer_sequences = {
    'producer-1': {
        'topic-orders': {
            0: 100,  # partition 0, last seq = 100
            1: 50,   # partition 1, last seq = 50
        }
    }
}

# On receive:
def handle_message(producer_id, topic, partition, sequence):
    last_seq = producer_sequences[producer_id][topic][partition]
    
    if sequence <= last_seq:
        return 'REJECT'  # Duplicate
    
    producer_sequences[producer_id][topic][partition] = sequence
    return 'ACCEPT'

Idempotent Producer Guarantees

Guarantee Description
Per partition Duplicates detected per partition
Per producer Each producer has own sequence
No cross-producer Different producers not deduplicated
Per session Sequences reset on producer restart

Limitations

1. Per producer instance:
   - Different producer IDs = no dedup
   - Need transactional.id for cross-session

2. Per partition:
   - Only within same partition
   - Cross-partition not deduplicated

3. Session-bound:
   - Sequences reset on restart
   - Need transactional.id for durability

When to Use

  • High-throughput producers
  • Retry-prone networks
  • When duplicates from retries are problematic
  • Combine with transactional consumer for true exactly-once

Practice Problems

0/3solved
Design Exactly Once System

Design a scalable Exactly Once system. Cover high-level architecture, data model, and API design.

Solution
// Complete system design:
// - Functional + Non-functional requirements
// - Capacity estimation
// - Data model (SQL/NoSQL choice)
// - API endpoints
// - Component architecture
// - Scaling strategy
// - Monitoring & reliability
Exactly Once Scaling

How would you scale Exactly Once to handle 10x the current load? Identify bottlenecks and solutions.

Solution
// Scaling approach:
// 1. Load balancing
// 2. Database sharding/replication
// 3. Cache layer (Redis)
// 4. CDN for static assets
// 5. Async processing (queues)
// 6. Microservices decomposition
Exactly Once Failure Modes

Analyze potential failure modes for Exactly Once and design mitigation strategies.

Solution
// Failure mitigation:
// 1. Redundancy (multi-AZ)
// 2. Circuit breakers
// 3. Retry with backoff
// 4. Dead letter queues
// 5. Health checks
// 6. Graceful degradation

Quiz

1. What two components make exactly-once possible?

Question 1 options

2. How does idempotent producer prevent duplicates?

Question 2 options

3. What is transactional messaging?

Question 3 options

4. What isolation level sees only committed messages?

Question 4 options

5. What is a limitation of idempotent producer?

Question 5 options

Flashcards

Question

Exactly-once = ? + ?

Answer

Idempotent Producer (no retry duplicates) + Transactional Consumer (atomic processing + offset commit)

Question

How idempotent producer works?

Answer

Assigns sequence numbers to messages; broker rejects duplicates with already-seen sequences

Question

What is transactional messaging?

Answer

Atomic operation: process messages + produce output + commit offsets, all or nothing

Question

read_committed vs read_uncommitted?

Answer

read_committed: only sees committed messages. read_uncommitted: sees all messages including uncommitted.

Question

Limitation of idempotent producer?

Answer

Only deduplicates within same producer instance; different producers or after restart not guaranteed

Revision Notes

Key Takeaways

  • 1.Exactly-once = idempotent producer + transactional consumer
  • 2.Idempotent producer uses sequence numbers for deduplication
  • 3.Transactional messaging atomically processes and commits
  • 4.read_committed isolation sees only committed messages
  • 5.Strongest delivery guarantee with highest complexity

Interview Tips

  • Explain both components clearly
  • Discuss sequence number tracking
  • Know Kafka configuration for exactly-once
  • Understand limitations and when to use

Cheat Sheet

Cheat Sheet: Exactly-Once

Components

  1. Idempotent Producer
    • Sequence numbers
    • Broker deduplication
  2. Transactional Consumer
    • Atomic processing
    • Offset commit

Kafka Config

enable.idempotence=true
acks=all
isolation.level=read_committed

Transaction Flow

Begin → Process → Produce → Commit/Abort

Guarantees

  • No producer retry duplicates
  • No consumer processing duplicates
  • Atomic processing + commit

Limitations

  • Per producer instance
  • Requires Kafka 0.11+
  • Performance overhead