Introduction
Message Queue is the backbone of any large-scale notification system. It helps decouple producers (services that create notifications) and consumers (workers that send emails), while ensuring durability when system components fail.
1. Why do we need a Message Queue?
There is no Message Queue
Client → API → Send Email → Response
│
└── Nếu ESP chậm 5s → API response chậm 5s
Nếu ESP down → API return 500
10K concurrent requests → 10K connections to ESP
There is a Message Queue
Client → API → Enqueue → Response (< 100ms)
│
└── Queue buffer messages
Workers consume at their own pace
ESP down → messages wait in queue
Spike traffic → queue absorbs burst
Key benefits
| Benefit | Description |
|---|---|
| Decoupling | API does not need to know about ESP |
| Buffering | Absorb traffic spikes |
| Durability | Messages persist via restarts |
| Ordering | Guaranteed order (per partition) |
| Retry | Failed messages re-process |
| Scaling | Add independent workers |
2. Compare Message Queue Technologies
Feature Matrix
| Features | Kafka | RabbitMQ | Amazon SQS | Redis Streams |
|---|---|---|---|---|
| Throughput | 1M+ msg/s | 50K msg/s | 3K msg/s/queue | 100K+ msg/s |
| Latency | ms | μs-ms | ms-100ms | μs-ms |
| Ordering | Per partition | Per queue | Best effort (FIFO available) | Per stream |
| Retention | Configurable (days/unlimited) | Until consumed | 14 days max | Configurable |
| Replay | ✅ Offset-based | ❌ | ❌ | ✅ ID-based |
| Dead Letter | Manual | ✅ Built-in | ✅ Built-in | Manual |
| Scaling | Partitions | Limited | Unlimited | Sharding |
| Operational | Complex | Moderate | Managed | Simple |
| Cost (10M msg/day) | ~$500/mo (self-hosted) | ~$300/mo | ~$150/mo | ~$200/mo |
Recommendation by use case
Email system gửi triệu email:
├── High throughput + replay needed → Kafka ✅
├── Low latency + simple setup → RabbitMQ
├── Managed + no ops overhead → Amazon SQS
└── Already have Redis + moderate scale → Redis Streams
Choose Kafka for million email notification system because:
- Highest throughput
- Message replay for debugging
- Partition-based parallel processing
- Event sourcing / audit trail
3. Kafka Deep Dive for Email System
Topic & Partition Design
Topic: email-send
├── Partition 0: critical priority emails
├── Partition 1: high priority emails
├── Partition 2-7: normal priority emails (6 partitions)
├── Partition 8-9: low priority emails
Topic: email-status
├── Partition 0-3: delivery status updates
Topic: email-dlq
├── Partition 0: failed messages for manual review
Topic: email-webhook
├── Partition 0-1: ESP webhook events (bounce, complaint, delivery)
Producer Configuration
from confluent_kafka import Producer
producer_config = {
'bootstrap.servers': 'kafka-1:9092,kafka-2:9092,kafka-3:9092',
# Durability: đảm bảo message không mất
'acks': 'all', # Wait for all replicas
'enable.idempotence': True, # Exactly-once per partition
# Performance: batch gửi messages
'batch.size': 65536, # 64KB batch
'linger.ms': 10, # Wait 10ms to fill batch
'compression.type': 'lz4', # Compress for throughput
# Reliability
'retries': 3,
'retry.backoff.ms': 100,
'max.in.flight.requests.per.connection': 5, # With idempotence=true
}
producer = Producer(producer_config)
def publish_email(campaign_id: str, recipient: dict, priority: str):
message = {
'campaign_id': campaign_id,
'recipient': recipient,
'priority': priority,
'created_at': datetime.utcnow().isoformat(),
}
# Partition key = priority để route đúng partition
producer.produce(
topic='email-send',
key=priority.encode('utf-8'),
value=json.dumps(message).encode('utf-8'),
callback=delivery_callback,
)
def delivery_callback(err, msg):
if err:
logger.error(f"Delivery failed: {err}")
# Save to local fallback queue
else:
logger.debug(f"Delivered to {msg.topic()}[{msg.partition()}]")
Consumer Configuration
from confluent_kafka import Consumer
consumer_config = {
'bootstrap.servers': 'kafka-1:9092,kafka-2:9092,kafka-3:9092',
'group.id': 'email-workers',
# Offset management
'auto.offset.reset': 'earliest',
'enable.auto.commit': False, # Manual commit after processing
# Performance
'max.poll.interval.ms': 300000, # 5 min max processing time
'session.timeout.ms': 30000,
'heartbeat.interval.ms': 10000,
'fetch.min.bytes': 1024, # Wait for 1KB before fetching
'fetch.max.wait.ms': 500,
# Batch consumption
'max.poll.records': 100, # 100 messages per poll
}
class EmailConsumer:
def __init__(self):
self.consumer = Consumer(consumer_config)
self.consumer.subscribe(['email-send'])
self.email_sender = EmailSender()
def run(self):
try:
while True:
messages = self.consumer.consume(
num_messages=100,
timeout=1.0
)
for msg in messages:
if msg.error():
self.handle_error(msg)
continue
try:
self.process_message(msg)
except Exception as e:
self.send_to_dlq(msg, str(e))
# Commit offset after successful processing
self.consumer.commit(asynchronous=False)
except KeyboardInterrupt:
pass
finally:
self.consumer.close()
def process_message(self, msg):
email_data = json.loads(msg.value())
self.email_sender.send(email_data)
4. Partition Strategy for Email Workload
Strategy 1: Priority-Based Partitioning
class PriorityPartitioner:
PARTITION_MAP = {
'critical': [0], # Dedicated partition
'high': [1], # Dedicated partition
'normal': [2, 3, 4, 5, 6, 7], # 6 partitions
'low': [8, 9], # 2 partitions
}
def partition(self, priority: str, num_partitions: int) -> int:
partitions = self.PARTITION_MAP[priority]
# Round-robin within priority's partitions
return random.choice(partitions)
Strategy 2: Domain-Based Partitioning
class DomainPartitioner:
"""Group emails by recipient domain for rate limiting"""
def partition(self, recipient_email: str, num_partitions: int) -> int:
domain = recipient_email.split('@')[1]
return hash(domain) % num_partitions
# → All @gmail.com emails go to same partition
# → Easier per-domain rate limiting
Strategy 3: Hybrid Approach (Recommended)
class HybridPartitioner:
def partition(self, message: dict, num_partitions: int) -> int:
priority = message['priority']
if priority == 'critical':
return 0 # Dedicated
# Hash by campaign_id for normal/marketing
# → All emails of same campaign go to same partition set
campaign_hash = hash(message['campaign_id'])
return 1 + (campaign_hash % (num_partitions - 1))
5. Exactly-Once vs At-Least-Once
At-Least-Once (Recommended for email)
# Consumer processes message → sends email → commits offset
# If crash between send and commit → re-delivery → duplicate email
# Solution: idempotency check before sending
async def process(self, message):
key = f"email:{message.campaign_id}:{message.recipient}"
# Idempotency check
if await redis.exists(key):
return # Already sent
await self.send_email(message)
await redis.set(key, "sent", ex=7*86400) # 7 days TTL
Exactly-Once (Kafka Transactions)
# More complex, higher latency, but guaranteed
producer = Producer({
**producer_config,
'transactional.id': f'email-worker-{worker_id}',
})
producer.init_transactions()
try:
producer.begin_transaction()
# Process and produce in transaction
result = process_email(message)
producer.produce('email-status', result)
producer.send_offsets_to_transaction(
consumer.position(consumer.assignment()),
consumer.consumer_group_metadata()
)
producer.commit_transaction()
except Exception:
producer.abort_transaction()
6. Monitoring Kafka Health
Key Metrics
# Consumer lag = messages in queue chưa processed
kafka_consumer_lag:
warning: > 10000
critical: > 100000
# Throughput
kafka_messages_per_second:
target: > 2000
# Consumer group health
kafka_consumer_group_members:
expected: 20 # Number of workers
# Partition balance
kafka_partition_assignment:
check: even distribution across consumers
Summary
- Message Queue is required for large-scale email systems
- Kafka is the best choice for high throughput and audit trails
- Partition strategy directly affects performance and ordering
- At-least-once + idempotency is the most pragmatic approach
Next article: Build a complete Event-Driven Notification Pipeline with Kafka.