はじめに
すべてのメールが同じ優先度を持つわけではなく、常にすぐに送信されるわけではありません。この記事は、OTP が常に電子メール マーケティングよりも優先されるようにするための 優先キュー と、最適なタイミングで電子メールを送信するための スケジュール エンジン を設計するのに役立ちます。
1. 複数レベルの優先順位システム
優先レベル
┌─────────────────────────────────────────────┐
│ CRITICAL (P0) — SLA: < 10 giây │
│ OTP, Password Reset, Security Alerts │
│ Resources: Dedicated workers + ESP pool │
├─────────────────────────────────────────────┤
│ HIGH (P1) — SLA: < 1 phút │
│ Order Confirmation, Payment Receipt │
│ Resources: Shared priority pool │
├─────────────────────────────────────────────┤
│ NORMAL (P2) — SLA: < 15 phút │
│ Marketing Campaign, Promotions │
│ Resources: Standard worker pool │
├─────────────────────────────────────────────┤
│ LOW (P3) — SLA: < 1 giờ │
│ Weekly Digest, Newsletter, Reports │
│ Resources: Background workers │
└─────────────────────────────────────────────┘
個別のトピックによる実装
class PriorityRouter:
TOPIC_MAP = {
'critical': 'email-send-p0',
'high': 'email-send-p1',
'normal': 'email-send-p2',
'low': 'email-send-p3',
}
WORKER_CONFIG = {
'critical': {
'min_workers': 5,
'max_workers': 20,
'poll_interval_ms': 10, # Poll rất nhanh
'batch_size': 1, # Process từng cái một
},
'high': {
'min_workers': 5,
'max_workers': 30,
'poll_interval_ms': 50,
'batch_size': 10,
},
'normal': {
'min_workers': 10,
'max_workers': 100,
'poll_interval_ms': 100,
'batch_size': 100,
},
'low': {
'min_workers': 2,
'max_workers': 20,
'poll_interval_ms': 500,
'batch_size': 500,
},
}
async def route(self, message: dict):
priority = message.get('priority', 'normal')
topic = self.TOPIC_MAP[priority]
await self.kafka.produce(topic, message)
2. スケジューリング エンジン
使用例
- 予定されているキャンペーン: 「月曜午前9時にフラッシュセールメールを送信」
- タイムゾーン対応: 「各ユーザーのタイムゾーンに従って午前 9 時に送信」
- 定期: 「毎週金曜日の午後 5 時に毎週のダイジェストを送信」
- 遅延: 「ユーザーがアクションを完了していない場合は、24 時間後にリマインダーを送信します」
アーキテクチャ
┌──────────────┐ ┌────────────────┐ ┌──────────────┐
│ Schedule │────▶│ Redis Sorted │────▶│ Scheduler │
│ API │ │ Set (ZSET) │ │ Worker │
│ │ │ │ │ │
│ POST /schedule│ │ Score = Unix │ │ Poll every │
│ {time, data} │ │ timestamp │ │ 1 second │
└──────────────┘ └────────────────┘ └──────┬───────┘
│
┌───────▼──────┐
│ Kafka Topic │
│ email-send │
└──────────────┘
Redis ソート セットの実装
import redis
import json
import time
class SchedulingEngine:
SCHEDULE_KEY = "email:schedule"
def __init__(self):
self.redis = redis.Redis(host='localhost', port=6379, db=0)
self.kafka = KafkaProducer()
async def schedule(self, message: dict, send_at: datetime):
"""Schedule an email for future delivery"""
score = send_at.timestamp() # Unix timestamp as score
member = json.dumps({
'id': message['id'],
'data': message,
'scheduled_at': send_at.isoformat(),
})
self.redis.zadd(self.SCHEDULE_KEY, {member: score})
async def schedule_timezone_aware(
self, campaign_id: str, recipients: list, local_time: str
):
"""Schedule email at local_time in each recipient's timezone"""
for recipient in recipients:
tz = recipient.get('timezone', 'Asia/Ho_Chi_Minh')
local_dt = parse_time(local_time, tz)
utc_dt = local_dt.astimezone(UTC)
await self.schedule(
message={
'campaign_id': campaign_id,
'recipient': recipient,
},
send_at=utc_dt,
)
async def poll_and_dispatch(self):
"""Background worker: check for due emails every second"""
while True:
now = time.time()
# Get all messages due for delivery
# ZRANGEBYSCORE: score <= now
due_messages = self.redis.zrangebyscore(
self.SCHEDULE_KEY,
'-inf',
now,
start=0,
num=1000, # Process max 1000 per tick
)
if due_messages:
pipe = self.redis.pipeline()
for raw in due_messages:
message = json.loads(raw)
# Publish to Kafka for immediate delivery
await self.kafka.produce(
'email-send',
message['data'],
)
# Remove from schedule
pipe.zrem(self.SCHEDULE_KEY, raw)
pipe.execute()
await asyncio.sleep(1) # Poll every second
3. Cron を使用した定期的なスケジュール
from croniter import croniter
class RecurringScheduler:
async def create_recurring(
self,
schedule_id: str,
cron_expression: str,
campaign_template: dict,
):
"""
Examples:
- "0 9 * * 1" → Mỗi thứ 2 lúc 9h
- "0 17 * * 5" → Mỗi thứ 6 lúc 17h
- "0 8 1 * *" → Ngày 1 hàng tháng lúc 8h
"""
await self.db.insert('recurring_schedules', {
'id': schedule_id,
'cron_expression': cron_expression,
'campaign_template': campaign_template,
'is_active': True,
'last_run_at': None,
'next_run_at': self._next_run(cron_expression),
})
def _next_run(self, cron_expr: str) -> datetime:
cron = croniter(cron_expr, datetime.utcnow())
return cron.get_next(datetime)
async def check_recurring(self):
"""Run every minute via system cron or scheduler"""
now = datetime.utcnow()
due_schedules = await self.db.query(
"SELECT * FROM recurring_schedules "
"WHERE is_active = true AND next_run_at <= %s",
[now]
)
for schedule in due_schedules:
# Create campaign from template
campaign = await self.create_campaign(
schedule.campaign_template
)
# Update next run
await self.db.update('recurring_schedules', {
'id': schedule.id,
'last_run_at': now,
'next_run_at': self._next_run(schedule.cron_expression),
})
4. スマートな送信時間の最適化
概念
全員に同時に送信するのではなく、各ユーザーがメールを開く可能性が最も高いときに送信します。
class SendTimeOptimizer:
async def get_optimal_time(self, user_id: str) -> int:
"""Trả về giờ tối ưu (0-23) dựa trên lịch sử open"""
# Query email open history
opens = await self.analytics.query(
"SELECT EXTRACT(HOUR FROM opened_at) as hour, COUNT(*) as cnt "
"FROM email_events "
"WHERE recipient_id = %s AND event_type = 'opened' "
"AND opened_at > NOW() - INTERVAL '90 days' "
"GROUP BY hour ORDER BY cnt DESC LIMIT 1",
[user_id]
)
if opens:
return opens[0]['hour']
# Fallback: industry best practices
# Theo SendGrid data: 10h sáng local time
return 10
async def schedule_optimized(self, campaign_id: str, recipients: list):
for recipient in recipients:
optimal_hour = await self.get_optimal_time(recipient['id'])
tz = recipient.get('timezone', 'Asia/Ho_Chi_Minh')
# Schedule at optimal hour in user's timezone
send_at = next_occurrence(optimal_hour, tz)
await self.scheduler.schedule(
message={'campaign_id': campaign_id, 'recipient': recipient},
send_at=send_at,
)
5. 分散スケジューリング — 重複の回避
問題
複数のスケジューラ インスタンス → 同じメッセージが 2 回ディスパッチされます。
解決策: Redis 分散ロック
class DistributedScheduler:
LOCK_KEY = "scheduler:lock"
LOCK_TTL = 30 # seconds
async def poll_with_lock(self):
while True:
# Try to acquire lock
acquired = self.redis.set(
self.LOCK_KEY,
self.instance_id,
nx=True,
ex=self.LOCK_TTL,
)
if acquired:
try:
await self.poll_and_dispatch()
finally:
# Release lock
# Lua script: only delete if we own the lock
self.redis.eval(
"""
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end
""",
1, self.LOCK_KEY, self.instance_id
)
else:
# Another instance is processing
pass
await asyncio.sleep(1)
6. スケジュール ダッシュボード API
# GET /api/v1/schedules
# Response:
{
"schedules": [
{
"id": "sched_001",
"campaign_id": "camp_weekly_digest",
"type": "recurring",
"cron": "0 17 * * 5",
"next_run": "2026-04-05T17:00:00Z",
"last_run": "2026-03-29T17:00:00Z",
"status": "active",
"pending_count": 0
},
{
"id": "sched_002",
"campaign_id": "camp_flash_sale",
"type": "scheduled",
"scheduled_at": "2026-04-02T09:00:00Z",
"status": "pending",
"pending_count": 5000000
}
],
"queue_snapshot": {
"critical": {"pending": 12, "processing": 5},
"high": {"pending": 230, "processing": 45},
"normal": {"pending": 1500000, "processing": 2000},
"low": {"pending": 50000, "processing": 100}
}
}
概要
- 4 つの優先度レベル による重要な電子メール専用のリソース
- スケジューリング エンジン は Redis Sorted Set を使用します - シンプルで信頼性が高い
- タイムゾーン対応 および 送信時間の最適化 により最高の UX を実現
- 分散ロック により重複したディスパッチを防止します
次の記事: SMTP の詳細 — 電子メール配信を根本から理解します。