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

レッスン 4: メッセージ キュー — 通知システムのバックボーン

電子メール システムにメッセージ キューが必要な理由。 Kafka、RabbitMQ、Amazon SQS、Redis ストリームを比較します。パーティショニング戦略、コンシューマ グループ、1 回限りのセマンティクス。

🏗️ アーキテクチャ — レッスン 4 レッスン 4: メッセージ キュー — のバックボーン 通知システム

数百万の電子メールを送信する通知システムを設計する

パート 2: メッセージ キューとイベント駆動型アーキテクチャ

xdev.asia

はじめに

メッセージ キューは、大規模な通知システムのバックボーンです。これは、システム コンポーネントに障害が発生した場合の耐久性を確保しながら、プロデューサー (通知を作成するサービス) とコンシューマー (電子メールを送信するワーカー) を分離するのに役立ちます。


1. なぜメッセージキューが必要なのでしょうか?

メッセージキューはありません

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

メッセージキューがあります

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

主な利点

メリット説明
デカップリングAPI は ESP について知る必要はありません
バッファリングトラフィックの急増を吸収
耐久性メッセージは再起動しても持続します。
注文保証された順序 (パーティションごと)
再試行失敗したメッセージは再処理されます
スケーリング独立したワーカーを追加する

2. メッセージ キュー テクノロジの比較

機能マトリックス

特長カフカラビットMQアマゾンSQSRedis ストリーム
スループット1M+ メッセージ/秒50K メッセージ/秒3K メッセージ/秒/キュー100K+ メッセージ/秒
レイテンシミリ秒μs-msms-100msμs-ms
注文パーティションごとキューごとベストエフォート (FIFO 使用可能)ストリームごと
保持設定可能 (日数/無制限)消費されるまで最長 14 日間設定可能
リプレイ✅ オフセットベース❌❌✅ IDベース
デッドレターマニュアル✅ 内蔵✅ 内蔵マニュアル
スケーリングパーティション限定無制限シャーディング
運用中複雑な中程度管理シンプル
コスト (1,000 万メッセージ/日)~$500/月 (自己ホスト型)~$300/月~$150/月~$200/月

ユースケース別の推奨事項

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

100 万件の電子メール通知システムとして Kafka を選択する理由は次のとおりです。

  • 最高のスループット
  • デバッグ用のメッセージの再生
  • パーティションベースの並列処理
  • イベントソーシング/監査証跡

3. 電子メール システムの Kafka の詳細

トピックとパーティションの設計

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)

プロデューサーの構成

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()}]")

コンシューマー構成

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. 電子メールワークロードのパーティション戦略

戦略 1: 優先順位ベースのパーティショニング

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)

戦略 2: ドメインベースのパーティショニング

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

戦略 3: ハイブリッド アプローチ (推奨)

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. 必ず 1 回と少なくとも 1 回

少なくとも 1 回 (電子メールに推奨)

# 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

必ず 1 回 (Kafka トランザクション)

# 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. Kafka の健全性の監視

主要な指標

# 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

概要

  • 大規模な電子メール システムにはメッセージ キューが 必須
  • Kafka は、高スループットと監査証跡に最適な選択肢です
  • パーティション戦略は パフォーマンス と 順序に直接影響します
  • 少なくとも 1 回 + べき等性 は最も実用的なアプローチです

次の記事: Kafka を使用して完全なイベント駆動型通知パイプラインを構築します。