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

第 6 課:優先權佇列和調度引擎

多重優先權佇列設計、延遲/預定電子郵件傳送、基於 cron 與基於事件的調度、時區感知發送、用於調度的 Redis 排序集。

🏗️ 建築 — 第 6 課 第 6 課:優先權佇列和調度引擎

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

第 2 部分:訊息佇列和事件驅動架構

亞洲開發網

簡介

並非所有電子郵件都具有相同的優先級,而且它們並不總是立即發送。本文將協助您設計優先佇列以確保 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.調度引擎

用例

  1. 預定活動:“週一上午 9 點發送閃購電子郵件”
  2. 時區感知:“根據每個使用者的時區在上午 9 點發送”
  3. 重複:“每週五下午 5 點發送每週摘要”
  4. 延遲:“如果使用者未完成操作,則在 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. 分散式調度-避免重複

問題

多個調度程序實例→相同的訊息被調度兩次。

###解決方案: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 - 簡單可靠
  • 時區感知和發送時間優化以獲得最佳用戶體驗
  • 分散式鎖防止重複調度

下一篇文章: SMTP 深入研究 — 從頭開始了解電子郵件傳送。