簡介
設計大型電子郵件系統不僅僅是編寫程式碼,而是應用正確的設計模式來使系統可靠、可擴展和可維護。本文將介紹 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 + 佇列原子性 | 資料一致性 |
| 艙壁 | 資源隔離 | 多租戶/優先權 |
下一篇文章: 深入探討訊息佇列-通知系統的支柱。