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
- Scheduled Campaign: "Send flash sale email at 9am on Monday"
- Timezone-aware: "Send at 9am according to each user's timezone"
- Recurring: "Send weekly digest every Friday at 5pm"
- 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.