How it Works
At-Least-Once: How it Works
Mechanism
1. Producer sends message to queue
2. Consumer receives message
3. Consumer processes message
4. Consumer ACKNOWLEDGES message (after processing)
Timeline:
Receive → Process → ACK
↑
If crash here, message redelivered!
Implementation
# RabbitMQ manual ack
channel.basic_consume(
queue='tasks',
on_message_callback=callback,
auto_ack=False # Manual acknowledgment
)
def callback(ch, method, properties, body):
try:
process_message(body)
ch.basic_ack(delivery_tag=method.delivery_tag) # Ack after
except Exception as e:
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
Kafka Manual Commit
// Kafka manual commit
props.put("enable.auto.commit", "false");
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
processMessage(record.value());
}
consumer.commitSync(); // Commit after processing
}
// Crash before commit = message redelivered
Data Flow
Producer → Queue → Consumer → Process → ACK
↓
If crash: message redelivered
↓
Consumer may process again
Idempotency Requirement
Idempotency Requirement
Why Idempotency?
Problem:
1. Consumer receives message A
2. Consumer processes A
3. Consumer crashes before ACK
4. Message A redelivered
5. Consumer processes A again (DUPLICATE!)
Solution: Idempotent processing
- Same message processed multiple times = same result
Idempotency Patterns
# Pattern 1: Database unique constraint
def process_order(order):
try:
db.insert('orders', order, conflict='ignore')
except DuplicateKey:
pass # Already processed
# Pattern 2: Check-then-act
def process_payment(payment):
if db.exists('payments', payment['id']):
return # Already processed
db.insert('payments', payment)
# Pattern 3: Optimistic locking
def update_inventory(item):
current = db.get('inventory', item['id'])
if current['version'] != item['expected_version']:
return # Already updated
db.update('inventory', item['id'], {
'quantity': current['quantity'] - item['amount'],
'version': current['version'] + 1
})
Idempotency Key
class IdempotentProcessor:
def __init__(self, db):
self.db = db
def process(self, message):
idempotency_key = message['id']
# Check if already processed
if self.db.is_processed(idempotency_key):
return self.db.get_result(idempotency_key)
# Process and store result atomically
with self.db.transaction():
result = self.handle(message)
self.db.mark_processed(idempotency_key, result)
return result
Idempotency by Operation
| Operation | Idempotency Strategy |
|---|---|
| INSERT | Unique constraint, ignore duplicate |
| UPDATE | Version check, conditional update |
| DELETE | Delete if exists (no error) |
| UPSERT | Natural idempotency |
| CALCULATION | Cache result by input hash |
Deduplication
Deduplication Strategies
Where to Deduplicate
1. Producer Level:
- Idempotent producer
- Unique message ID
- Broker deduplicates
2. Broker Level:
- Kafka idempotent producer
- Sequence numbers
- Duplicate detection
3. Consumer Level:
- Check before processing
- Idempotent operations
- Result caching
Consumer-Level Deduplication
class DeduplicatingConsumer:
def __init__(self, cache_client, db_client):
self.cache = cache_client
self.db = db_client
def process(self, message):
msg_id = message['id']
# Check cache (fast)
if self.cache.get(f"processed:{msg_id}"):
return # Already processed
# Check database (slower)
if self.db.is_processed(msg_id):
self.cache.set(f"processed:{msg_id}", True, ttl=3600)
return
# Process
result = self.handle(message)
# Mark as processed (atomic)
with self.db.transaction():
self.db.mark_processed(msg_id, result)
self.cache.set(f"processed:{msg_id}", True, ttl=3600)
return result
Time-Based Deduplication
class TimeWindowDeduplicator:
def __init__(self, window_seconds=300):
self.window = window_seconds
self.seen = {} # message_id -> timestamp
def is_duplicate(self, message_id):
now = time.time()
# Clean old entries
self.seen = {k: v for k, v in self.seen.items()
if now - v < self.window}
# Check if seen
if message_id in self.seen:
return True
self.seen[message_id] = now
return False
Deduplication Trade-offs
| Strategy | Pros | Cons |
|---|---|---|
| Database check | Reliable | Slow |
| Cache check | Fast | May miss |
| Time window | Simple | Approximate |
| Bloom filter | Memory efficient | False positives |
Best Practices
- Always implement idempotency for at-least-once
- Use unique message IDs
- Check before processing
- Atomic mark-as-processed
- Monitor duplicate rates
Practice Problems
Design a scalable At Least 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 At Least 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 At Least 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. Why does at-least-once need idempotency?
2. What is an idempotent operation?
3. How to make database INSERT idempotent?
4. What is a deduplication window?
5. Why mark as processed atomically?
Flashcards
Question
At-least-once ack timing?
Click to reveal answer
Answer
Acknowledge AFTER processing - crash before ack = message redelivered
Question
Why idempotency required?
Click to reveal answer
Answer
Messages may be redelivered on failure; idempotent processing ensures same result on multiple calls
Question
How to make INSERT idempotent?
Click to reveal answer
Answer
Add unique constraint and ignore duplicate key error on insert
Question
Deduplication strategies?
Click to reveal answer
Answer
1) Database check, 2) Cache check, 3) Time window, 4) Bloom filter
Question
Atomic mark-as-processed?
Click to reveal answer
Answer
Ensures message is marked as processed in same transaction as handling, preventing double-processing
Revision Notes
Key Takeaways
- 1.At-least-once acks after processing - crash means redelivery
- 2.Idempotency is required to handle duplicates safely
- 3.Use unique constraints for INSERT idempotency
- 4.Atomic mark-as-processed prevents double-processing
- 5.Most common delivery semantic for critical operations
Interview Tips
- •Explain why idempotency is required
- •Give examples of idempotent operations
- •Discuss deduplication strategies
- •Compare with at-most-once trade-offs
Cheat Sheet
Cheat Sheet: At-Least-Once
Mechanism
Process before ACK
Crash before ACK = redelivery
Idempotency Required
- INSERT: Unique constraint
- UPDATE: Version check
- DELETE: Delete if exists
Deduplication
- Check cache (fast)
- Check DB (reliable)
- Mark atomically
Best Practices
- Unique message IDs
- Check before processing
- Atomic mark-as-processed
- Monitor duplicate rates