簡介
在上一課中,您了解了發送數百萬封電子郵件的問題。現在,我們將設計整體架構—整個系統的藍圖。這是系統設計中最重要的一步。
1. 高層架構
元件概述
┌─────────────┐ ┌──────────────────┐ ┌─────────────────┐
│ Client │────▶│ API Gateway │────▶│ Notification │
│ (Web/API) │ │ (Rate Limit, │ │ Service │
│ │ │ Auth, Route) │ │ (Validate, │
└─────────────┘ └──────────────────┘ │ Enqueue) │
└────────┬────────┘
│
┌────────▼────────┐
│ Message Queue │
│ (Kafka/SQS) │
│ │
│ ┌────────────┐ │
│ │ Priority Q │ │
│ │ Critical │ │
│ │ High │ │
│ │ Normal │ │
│ │ Low │ │
│ └────────────┘ │
└────────┬────────┘
│
┌────────▼────────┐
│ Worker Pool │
│ (N workers) │
│ │
│ ┌──┐┌──┐┌──┐ │
│ │W1││W2││W3│...│
│ └──┘└──┘└──┘ │
└────────┬────────┘
│
┌─────────────────┬──────┴───────┐
▼ ▼ ▼
┌──────────────┐ ┌──────────────┐ ┌────────────┐
│ Amazon SES │ │ SendGrid │ │ Mailgun │
│ (Primary) │ │ (Secondary) │ │ (Backup) │
└──────┬───────┘ └──────┬───────┘ └─────┬──────┘
│ │ │
└─────────┬───────┘ │
▼ │
┌──────────────┐ │
│ Webhook │◀───────────────┘
│ Receiver │
└──────┬───────┘
│
┌──────▼───────┐
│ Status │
│ Tracker │
└──────┬───────┘
│
┌──────▼───────┐
│ Analytics │
│ Dashboard │
└──────────────┘
2. 詳細元件
2.1 API網關
負責:
- 身份驗證:API 金鑰/OAuth2 驗證
- 速率限制:限制每個客戶端的請求速率
- 請求驗證:架構驗證
- 路由:路由到正確的服務
// POST /api/v1/notifications/email
{
"campaign_id": "camp_2026_flash_sale",
"template_id": "tmpl_flash_sale_v2",
"recipients": {
"type": "segment",
"segment_id": "seg_active_users_30d"
},
"schedule": {
"type": "immediate" // hoặc "scheduled", "recurring"
},
"priority": "high",
"metadata": {
"utm_source": "email",
"utm_campaign": "flash_sale_april"
}
}
2.2 通知服務
核心業務邏輯:
class NotificationService:
def create_campaign(self, request):
# 1. Validate request
self.validate(request)
# 2. Resolve recipients
recipients = self.recipient_service.resolve(request.recipients)
# → Returns: List[RecipientInfo] with email, name, preferences
# 3. Check suppression list
recipients = self.suppression_service.filter(recipients)
# → Remove: unsubscribed, bounced, complained
# 4. Create campaign record
campaign = self.db.create_campaign(
id=generate_ulid(),
total_recipients=len(recipients),
status="QUEUED",
template_id=request.template_id,
)
# 5. Enqueue messages (chunked)
for chunk in chunked(recipients, size=1000):
self.queue.publish(
topic="email-send",
messages=[
{
"campaign_id": campaign.id,
"recipient": r,
"template_id": request.template_id,
"priority": request.priority,
}
for r in chunk
],
partition_key=request.priority,
)
# 6. Return campaign ID for tracking
return {"campaign_id": campaign.id, "queued": len(recipients)}
2.3 訊息佇列(Kafka)
主題設計:
email-send # Main send queue
├── partition-0 # Critical priority
├── partition-1 # High priority
├── partition-2..N # Normal priority
email-status # Delivery status events
email-dlq # Dead letter queue
email-webhook # ESP webhook events
2.4 工作池
class EmailWorker:
def __init__(self):
self.template_engine = TemplateEngine()
self.rate_limiter = RateLimiter()
self.email_provider = MultiProviderClient()
async def process(self, message):
# 1. Rate limit check
await self.rate_limiter.acquire(
provider=self.email_provider.current,
domain=message.recipient.domain
)
# 2. Render template
html = self.template_engine.render(
template_id=message.template_id,
context=message.recipient.to_dict()
)
# 3. Send email
result = await self.email_provider.send(
to=message.recipient.email,
subject=message.subject,
html=html,
headers=message.tracking_headers
)
# 4. Publish status
await self.status_publisher.publish({
"message_id": message.id,
"status": result.status,
"provider": result.provider,
"timestamp": now()
})
3. 資料庫架構設計
核心表
-- Campaigns
CREATE TABLE campaigns (
id UUID PRIMARY KEY,
name VARCHAR(255) NOT NULL,
template_id UUID NOT NULL REFERENCES templates(id),
status VARCHAR(20) DEFAULT 'DRAFT',
priority VARCHAR(10) DEFAULT 'normal',
total_count INTEGER DEFAULT 0,
sent_count INTEGER DEFAULT 0,
failed_count INTEGER DEFAULT 0,
scheduled_at TIMESTAMPTZ,
started_at TIMESTAMPTZ,
completed_at TIMESTAMPTZ,
created_at TIMESTAMPTZ DEFAULT NOW(),
created_by UUID NOT NULL
);
-- Individual email records
CREATE TABLE email_messages (
id UUID PRIMARY KEY,
campaign_id UUID REFERENCES campaigns(id),
recipient VARCHAR(255) NOT NULL,
status VARCHAR(20) DEFAULT 'QUEUED',
provider VARCHAR(50),
provider_id VARCHAR(255),
retry_count INTEGER DEFAULT 0,
sent_at TIMESTAMPTZ,
delivered_at TIMESTAMPTZ,
opened_at TIMESTAMPTZ,
clicked_at TIMESTAMPTZ,
bounced_at TIMESTAMPTZ,
error_message TEXT,
created_at TIMESTAMPTZ DEFAULT NOW()
);
-- Partitioned by date for performance
CREATE INDEX idx_email_messages_campaign ON email_messages(campaign_id);
CREATE INDEX idx_email_messages_status ON email_messages(status);
CREATE INDEX idx_email_messages_recipient ON email_messages(recipient);
-- Suppression list
CREATE TABLE suppressions (
email VARCHAR(255) PRIMARY KEY,
reason VARCHAR(50) NOT NULL, -- 'bounce', 'complaint', 'unsubscribe'
source VARCHAR(100),
created_at TIMESTAMPTZ DEFAULT NOW()
);
-- Templates
CREATE TABLE templates (
id UUID PRIMARY KEY,
name VARCHAR(255) NOT NULL,
subject VARCHAR(500) NOT NULL,
html_body TEXT NOT NULL,
text_body TEXT,
variables JSONB DEFAULT '[]',
version INTEGER DEFAULT 1,
is_active BOOLEAN DEFAULT true,
created_at TIMESTAMPTZ DEFAULT NOW()
);
4. 資料流:從觸發器到收件匣
1. Client gửi API request
│
2. API Gateway: auth, rate limit, validate
│
3. Notification Service:
│ ├─ Resolve recipients (query database/segment service)
│ ├─ Filter suppression list
│ ├─ Create campaign record
│ └─ Enqueue messages to Kafka (chunked by 1000)
│
4. Kafka: buffer messages, distribute to partitions
│
5. Worker Pool (N consumers):
│ ├─ Consume from Kafka
│ ├─ Rate limit check (per provider, per domain)
│ ├─ Render template with recipient data
│ ├─ Send via Email Provider (SES/SendGrid)
│ └─ Publish status to email-status topic
│
6. Email Provider → recipient's mail server
│
7. Webhook callbacks: delivered, opened, clicked, bounced
│
8. Status Tracker: update email_messages table
│
9. Analytics: aggregate metrics, update dashboard
5. 冪等性 — 只送一次
問題
Worker 在發送電子郵件後但在提交偏移之前崩潰 → Kafka 重新交付 → 電子郵件發送兩次。
解決方案:冪等性金鑰
async def process_with_idempotency(self, message):
idempotency_key = f"{message.campaign_id}:{message.recipient.email}"
# Check if already processed
if await self.redis.exists(f"sent:{idempotency_key}"):
logger.info(f"Skipping duplicate: {idempotency_key}")
return
# Send email
result = await self.send_email(message)
# Mark as processed (TTL 7 days)
await self.redis.set(
f"sent:{idempotency_key}",
result.provider_id,
ex=7 * 86400
)
6. 關注點分離
| 組件 | 責任 | 縮放 |
|---|---|---|
| API閘道 | 身分驗證、速率限制、路由 | 橫向、無國籍 |
| 通知服務 | 業務邏輯、排隊 | 橫向、無國籍 |
| 訊息佇列 | 緩衝、排序、耐用性 | 卡夫卡分區 |
| 工人池 | 渲染、發送、重試 | 水平,自動縮放 |
| 狀態追蹤器 | 更新交貨狀態 | 水平、非同步 |
| 分析 | 聚合指標 | 批量處理 |
總結
您現在擁有完整的藍圖:
- 6 個主要組成部分 以及每個組成部分的職責
- 通知系統的資料庫架構
- 資料流從觸發器到收件匣的端到端
- 冪等性確保只發送一次
下一篇文章: 深入研究大型電子郵件系統的重要設計模式。