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

Lesson 6: Priority Queue and Scheduling Engine

Multi-priority queue design, delayed/scheduled email delivery, cron-based vs event-based scheduling, time-zone aware sending, Redis Sorted Set for scheduling.

🏗️ Architecture — Lesson 6 Lesson 6: Priority Queue and Scheduling Engine

Design a Notification System to send millions of Emails

Part 2: Message Queue & Event-Driven Architecture

xdev.asia

Introduction

Not all emails have the same priority, and they don't always send immediately. This article will help you design Priority Queue to ensure OTP always comes before email marketing, and Scheduling Engine to send emails at the optimal time.


1. Multi-Level Priority System

Priority Levels

┌─────────────────────────────────────────────┐
│ 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                │
└─────────────────────────────────────────────┘

Implementation with Separate Topics

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. Scheduling Engine

Use Cases

  1. Scheduled Campaign: "Send flash sale email at 9am on Monday"
  2. Timezone-aware: "Send at 9am according to each user's timezone"
  3. Recurring: "Send weekly digest every Friday at 5pm"
  4. Delayed: "Send reminder after 24 hours if user has not completed action"

Architecture

┌──────────────┐     ┌────────────────┐     ┌──────────────┐
│  Schedule    │────▶│  Redis Sorted  │────▶│  Scheduler   │
│  API         │     │  Set (ZSET)    │     │  Worker      │
│              │     │                │     │              │
│ POST /schedule│    │ Score = Unix   │     │ Poll every   │
│ {time, data} │     │ timestamp      │     │ 1 second     │
└──────────────┘     └────────────────┘     └──────┬───────┘
                                                    │
                                            ┌───────▼──────┐
                                            │ Kafka Topic  │
                                            │ email-send   │
                                            └──────────────┘

Redis Sorted Set Implementation

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. Recurring Schedule with 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. Smart Send Time Optimization

Concepts

Instead of sending it to everyone at the same time, send it at the time each user is most likely to open emails.

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. Distributed Scheduling — Avoiding Duplicates

Problem

Multiple scheduler instances → same messages are dispatched twice.

Solution: Redis Distributed Lock

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. Schedule Dashboard 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}
    }
}

Summary

  • 4 priority levels with dedicated resources for critical emails
  • Scheduling Engine uses Redis Sorted Set — simple & reliable
  • Timezone-aware and send time optimization for the best UX
  • Distributed lock prevents duplicate dispatch

Next article: SMTP Deep Dive — understand email delivery from the ground up.