Chuyển đến nội dung chính

Lesson 4: Message Queue — The backbone of the Notification System

Why is a message queue needed for email system. Compare Kafka vs RabbitMQ vs Amazon SQS vs Redis Streams. Partitioning strategies, consumer groups, exactly-once semantics.

🏗️ Architecture — Lesson 4 Lesson 4: Message Queue — The backbone of Notification System

Design a Notification System to send millions of Emails

Part 2: Message Queue & Event-Driven Architecture

xdev.asia

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

BenefitDescription
DecouplingAPI does not need to know about ESP
BufferingAbsorb traffic spikes
DurabilityMessages persist via restarts
OrderingGuaranteed order (per partition)
RetryFailed messages re-process
ScalingAdd independent workers

2. Compare Message Queue Technologies

Feature Matrix

FeaturesKafkaRabbitMQAmazon SQSRedis Streams
Throughput1M+ msg/s50K msg/s3K msg/s/queue100K+ msg/s
Latencymsμs-msms-100msμs-ms
OrderingPer partitionPer queueBest effort (FIFO available)Per stream
RetentionConfigurable (days/unlimited)Until consumed14 days maxConfigurable
Replay✅ Offset-based❌❌✅ ID-based
Dead LetterManual✅ Built-in✅ Built-inManual
ScalingPartitionsLimitedUnlimitedSharding
OperationalComplexModerateManagedSimple
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.