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

第 3 課:大型電子郵件系統的設計模式

扇出模式、生產者-消費者模式、優先權佇列模式。用於通知工作流程的斷路器、隔板、寄件箱模式、Saga 模式。

🏗️ 建築 — 第 3 課 第 3 課:電子郵件系統的設計模式 大規模

設計一個通知系統來發送數百萬封電子郵件

第 1 部分:基礎 — 了解大規模通知問題

亞洲開發網

簡介

設計大型電子郵件系統不僅僅是編寫程式碼,而是應用正確的設計模式來使系統可靠、可擴展和可維護。本文將介紹 6 個最重要的模式。


1. 扇出模式

數學問題

發送給 1000 萬收件人的行銷活動 → 需要 扇出 成 1000 萬條單獨的訊息。

實作

Campaign Created (1 event)
    │
    ▼
Fan-Out Service
    │
    ├──▶ Batch 1: recipients[0:1000]     → 1000 messages to Queue
    ├──▶ Batch 2: recipients[1000:2000]   → 1000 messages to Queue
    ├──▶ Batch 3: recipients[2000:3000]   → 1000 messages to Queue
    │    ...
    └──▶ Batch 10000: recipients[9999000:10000000] → 1000 messages
class FanOutService:
    BATCH_SIZE = 1000

    async def fan_out_campaign(self, campaign_id: str):
        campaign = await self.db.get_campaign(campaign_id)
        segment = await self.segment_service.get(campaign.segment_id)

        # Stream recipients in batches (không load tất cả vào memory)
        async for batch in segment.stream_recipients(batch_size=self.BATCH_SIZE):
            # Filter suppressed emails
            active = await self.suppression_service.filter_batch(batch)

            # Enqueue batch
            messages = [
                EmailMessage(
                    campaign_id=campaign_id,
                    recipient=r,
                    template_id=campaign.template_id,
                    priority=campaign.priority,
                )
                for r in active
            ]
            await self.queue.publish_batch("email-send", messages)

            # Update progress
            await self.db.increment_campaign_queued(
                campaign_id, len(messages)
            )

要點

  • 串流傳輸,不載入全部:使用遊標/分頁串流傳輸收件人
  • 批量發布:將大量商品發送到佇列,而不是單獨發布
  • 進度追蹤:更新扇出進度

2. 生產者-消費者模式

設計

Producers (Notification Services)
    │
    ▼
┌──────────────────────────────────┐
│         Message Queue            │
│  ┌─────┐ ┌─────┐ ┌─────┐       │
│  │ P0  │ │ P1  │ │ P2  │  ...  │
│  └─────┘ └─────┘ └─────┘       │
└──────────────┬───────────────────┘
               │
    ┌──────────┼──────────┐
    ▼          ▼          ▼
┌────────┐ ┌────────┐ ┌────────┐
│Worker 1│ │Worker 2│ │Worker 3│  ... (Consumer Group)
└────────┘ └────────┘ └────────┘

消費者群組配置

# Kafka consumer config
consumer:
  group_id: email-workers
  auto_offset_reset: earliest
  enable_auto_commit: false  # Manual commit sau khi send thành công
  max_poll_records: 100      # Process 100 messages per poll
  session_timeout_ms: 30000
  heartbeat_interval_ms: 10000

縮放規則

Throughput target: 2000 emails/giây
Single worker capacity: ~100 emails/giây
Workers needed: 2000 / 100 = 20 workers

Kafka partitions >= Workers count
→ 24 partitions (headroom for scaling)

3. 優先隊列模式

多級優先級

class PriorityQueueRouter:
    TOPICS = {
        "critical": "email-send-critical",   # OTP, password reset
        "high": "email-send-high",           # Order confirmation
        "normal": "email-send-normal",       # Marketing
        "low": "email-send-low",             # Newsletter digest
    }

    # Worker allocation
    WORKER_ALLOCATION = {
        "critical": 0.30,  # 30% workers cho critical
        "high": 0.30,      # 30% cho high
        "normal": 0.30,    # 30% cho normal
        "low": 0.10,       # 10% cho low
    }

    async def route(self, message: EmailMessage):
        topic = self.TOPICS[message.priority]
        await self.kafka.publish(topic, message)

預防飢餓

class WeightedConsumer:
    """Đảm bảo low-priority messages không bị starve"""

    async def poll(self):
        # Weighted round-robin
        messages = []
        messages += await self.poll_topic("critical", max=30)
        messages += await self.poll_topic("high", max=30)
        messages += await self.poll_topic("normal", max=30)
        messages += await self.poll_topic("low", max=10)
        return messages

4. 斷路器模式

數學問題

ESP(電子郵件服務提供者)可能會降級或限制您→需要暫時斷開連線。

States: CLOSED → OPEN → HALF_OPEN → CLOSED
                  │                      ▲
                  └──────────────────────┘
                     (sau timeout, thử lại)

實作

class ESPCircuitBreaker:
    def __init__(self, provider_name: str):
        self.provider = provider_name
        self.state = "CLOSED"
        self.failure_count = 0
        self.failure_threshold = 10        # Mở circuit sau 10 failures
        self.timeout = timedelta(minutes=5) # Thử lại sau 5 phút
        self.last_failure_time = None

    async def call(self, send_func, *args):
        if self.state == "OPEN":
            if datetime.now() - self.last_failure_time > self.timeout:
                self.state = "HALF_OPEN"
            else:
                raise CircuitOpenError(self.provider)

        try:
            result = await send_func(*args)
            if self.state == "HALF_OPEN":
                self.state = "CLOSED"
                self.failure_count = 0
            return result
        except ESPError as e:
            self.failure_count += 1
            self.last_failure_time = datetime.now()
            if self.failure_count >= self.failure_threshold:
                self.state = "OPEN"
                logger.critical(
                    f"Circuit OPEN for {self.provider}: "
                    f"{self.failure_count} failures"
                )
            raise

多提供者故障轉移

class MultiProviderClient:
    def __init__(self):
        self.providers = {
            "ses": AmazonSESProvider(),
            "sendgrid": SendGridProvider(),
            "mailgun": MailgunProvider(),
        }
        self.breakers = {
            name: ESPCircuitBreaker(name)
            for name in self.providers
        }
        self.primary = "ses"
        self.fallback_order = ["sendgrid", "mailgun"]

    async def send(self, email: Email) -> SendResult:
        # Try primary
        try:
            return await self.breakers[self.primary].call(
                self.providers[self.primary].send, email
            )
        except (CircuitOpenError, ESPError):
            pass

        # Try fallbacks
        for provider_name in self.fallback_order:
            try:
                return await self.breakers[provider_name].call(
                    self.providers[provider_name].send, email
                )
            except (CircuitOpenError, ESPError):
                continue

        raise AllProvidersDownError("No email provider available")

5.寄件匣模式

數學問題

需要保證:資料庫寫入+佇列發布是原子的。如果發布成功但資料庫崩潰→不一致。

解決方案:交易寄件箱

1. Write to DB + outbox trong cùng 1 transaction
2. Background process đọc outbox → publish to queue
3. Mark outbox entry as published
class OutboxPublisher:
    async def create_campaign_with_outbox(self, campaign_data):
        async with self.db.transaction():
            # 1. Create campaign
            campaign = await self.db.insert("campaigns", campaign_data)

            # 2. Write to outbox (same transaction!)
            await self.db.insert("outbox", {
                "id": generate_ulid(),
                "aggregate_type": "campaign",
                "aggregate_id": campaign.id,
                "event_type": "CampaignCreated",
                "payload": json.dumps(campaign_data),
                "status": "PENDING",
                "created_at": now(),
            })

        # Transaction committed → both campaign & outbox saved atomically

    async def poll_outbox(self):
        """Background job: relay outbox → Kafka"""
        entries = await self.db.query(
            "SELECT * FROM outbox WHERE status = 'PENDING' "
            "ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED"
        )
        for entry in entries:
            await self.kafka.publish(
                topic=f"notification-{entry.event_type}",
                message=entry.payload,
            )
            await self.db.update(
                "outbox",
                {"id": entry.id},
                {"status": "PUBLISHED", "published_at": now()}
            )

6.艙壁圖案

數學問題

行銷電子郵件氾濫 → 影響交易電子郵件(OTP 延遲) → 使用者無法登入。

解決方案:資源隔離

┌─────────────────────────────────────────┐
│           Notification System            │
│                                          │
│  ┌──────────────────┐ ┌───────────────┐ │
│  │ Transactional    │ │  Marketing    │ │
│  │ Bulkhead         │ │  Bulkhead     │ │
│  │                  │ │               │ │
│  │ Workers: 10      │ │ Workers: 20   │ │
│  │ Queue: dedicated │ │ Queue: shared │ │
│  │ ESP: dedicated   │ │ ESP: shared   │ │
│  │ Rate: unlimited  │ │ Rate: limited │ │
│  └──────────────────┘ └───────────────┘ │
└─────────────────────────────────────────┘
class BulkheadConfig:
    BULKHEADS = {
        "transactional": {
            "max_concurrent": 10,
            "queue_size": 1000,
            "timeout_seconds": 5,
            "dedicated_provider_pool": True,
        },
        "marketing": {
            "max_concurrent": 50,
            "queue_size": 100000,
            "timeout_seconds": 30,
            "dedicated_provider_pool": False,
        },
    }

總結

圖案解決何時使用
扇出1 個活動 → N 則訊息活動廣播
生產者-消費者從進程中解耦發送隨時
優先隊列批評與行銷多類型通知
斷路器ESP 故障外部服務電話
寄件匣DB + 佇列原子性資料一致性
艙壁資源隔離多租戶/優先權

下一篇文章: 深入探討訊息佇列-通知系統的支柱。