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
- Use unique transactional.id per producer instance
- Set appropriate isolation.level for consumers
- Handle transaction timeouts gracefully
- Monitor transaction metrics
- 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
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 & reliabilityHow 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 decompositionAnalyze 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 degradationQuiz
1. What two components make exactly-once possible?
2. How does idempotent producer prevent duplicates?
3. What is transactional messaging?
4. What isolation level sees only committed messages?
5. What is a limitation of idempotent producer?
Flashcards
Question
Exactly-once = ? + ?
Click to reveal answer
Answer
Idempotent Producer (no retry duplicates) + Transactional Consumer (atomic processing + offset commit)
Question
How idempotent producer works?
Click to reveal answer
Answer
Assigns sequence numbers to messages; broker rejects duplicates with already-seen sequences
Question
What is transactional messaging?
Click to reveal answer
Answer
Atomic operation: process messages + produce output + commit offsets, all or nothing
Question
read_committed vs read_uncommitted?
Click to reveal answer
Answer
read_committed: only sees committed messages. read_uncommitted: sees all messages including uncommitted.
Question
Limitation of idempotent producer?
Click to reveal answer
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
- Idempotent Producer
- Sequence numbers
- Broker deduplication
- 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