Celery priority queue dengan Redis backend
Celery worker dengan priority queue — urgent task (notification user) jalan duluan dari batch email. Pakai Redis broker dengan queue per priority.
Dipublikasikan 4 Juli 2026
Worker Celery yang campur kirim email batch promosi dengan OTP urgent itu kacau — OTP user bisa nunggu 10 menit karena di-queue belakang 50rb email promo. Solusi: queue terpisah per priority dan worker yang prefetch lebih agresif untuk queue urgent. Snippet ini setup 3-tier priority.
Kode
# celery_app.py
from celery import Celery
from kombu import Queue
# Redis broker
broker_url = "redis://redis:6379/0"
result_backend = "redis://redis:6379/1"
app = Celery("tokopedia_jobs", broker=broker_url, backend=result_backend)
app.conf.update(
# Routes — tentukan task mana ke queue mana
task_routes={
"tasks.kirim_otp": {"queue": "urgent"},
"tasks.kirim_notif_push": {"queue": "urgent"},
"tasks.kirim_email_konfirmasi_order": {"queue": "normal"},
"tasks.kirim_email_promo_batch": {"queue": "low"},
"tasks.generate_laporan_pdf": {"queue": "low"},
},
# Queue definition
task_queues=(
Queue("urgent", routing_key="urgent"),
Queue("normal", routing_key="normal"),
Queue("low", routing_key="low"),
),
task_default_queue="normal",
task_default_routing_key="normal",
# Reliability
task_acks_late=True, # Ack setelah task selesai (bukan saat received)
task_reject_on_worker_lost=True, # Re-queue kalau worker crash
task_track_started=True,
# Worker prefetch
worker_prefetch_multiplier=4, # Default — fetch 4 task per worker process
task_acks_on_failure_or_timeout=False,
# Time limit
task_soft_time_limit=300, # 5 menit warning
task_time_limit=600, # 10 menit hard kill
# Result expiration
result_expires=3600,
)
# tasks.py
import logging
from celery import shared_task
from celery.exceptions import SoftTimeLimitExceeded
logger = logging.getLogger(__name__)
@shared_task(
bind=True,
autoretry_for=(ConnectionError,),
retry_backoff=True,
retry_backoff_max=60,
max_retries=3,
)
def kirim_otp(self, user_id: int, kode_otp: str, channel: str = "sms") -> dict:
"""OTP HARUS cepat sampai. Queue: urgent."""
logger.info(f"Kirim OTP user_id={user_id} via {channel}")
try:
if channel == "sms":
send_sms_via_provider(user_id, f"Kode OTP: {kode_otp}. Berlaku 5 menit.")
elif channel == "whatsapp":
send_whatsapp(user_id, f"Kode OTP kamu: {kode_otp}")
else:
raise ValueError(f"Channel tidak dikenal: {channel}")
return {"user_id": user_id, "status": "sent", "channel": channel}
except ConnectionError:
# autoretry handle retry — re-raise
raise
@shared_task(bind=True)
def kirim_email_konfirmasi_order(self, order_id: int) -> dict:
"""Konfirmasi order — penting tapi tidak super urgent. Queue: normal."""
order = fetch_order(order_id)
send_email(
to=order.email,
subject=f"Pesanan #{order.id} berhasil",
body=render_template("email/konfirmasi_order.html", order=order),
)
return {"order_id": order_id, "email_sent_to": order.email}
@shared_task(bind=True, soft_time_limit=600)
def kirim_email_promo_batch(self, kampanye_id: int) -> dict:
"""Email promo ke ribuan user. Queue: low — jangan ganggu urgent task."""
kampanye = fetch_kampanye(kampanye_id)
user_list = fetch_target_users(kampanye)
sent = 0
failed = 0
try:
for user in user_list:
try:
send_email(to=user.email, subject=kampanye.subject, body=kampanye.body)
sent += 1
except Exception as e:
logger.warning(f"Email gagal ke {user.email}: {e}")
failed += 1
except SoftTimeLimitExceeded:
logger.warning("Time limit hampir habis, save progress dan re-queue")
# Re-queue sisanya dengan task baru
kirim_email_promo_batch.apply_async(
args=[kampanye_id],
queue="low",
countdown=60,
)
return {"sent": sent, "failed": failed}
@shared_task
def generate_laporan_pdf(periode: str, koperasi_id: int) -> str:
"""Generate laporan PDF — berat tapi tidak urgent."""
pdf_path = generate_pdf(periode, koperasi_id)
upload_to_s3(pdf_path, f"laporan/{koperasi_id}/{periode}.pdf")
notify_admin(koperasi_id, pdf_path)
return pdf_path
# Helper dummy
def send_sms_via_provider(user_id, msg): pass
def send_whatsapp(user_id, msg): pass
def send_email(**kwargs): pass
def fetch_order(id): pass
def fetch_kampanye(id): pass
def fetch_target_users(k): return []
def render_template(name, **ctx): return ""
def generate_pdf(p, k): return "/tmp/x.pdf"
def upload_to_s3(p, k): pass
def notify_admin(k, p): pass
Pemakaian
# Worker konfigurasi — pisah per queue untuk SLA berbeda
# Worker urgent: prefetch rendah, banyak worker — latency minimum
celery -A celery_app worker \
--queues=urgent \
--concurrency=8 \
--prefetch-multiplier=1 \
--loglevel=info \
-n urgent@%h
# Worker normal: prefetch normal
celery -A celery_app worker \
--queues=normal \
--concurrency=4 \
--prefetch-multiplier=4 \
--loglevel=info \
-n normal@%h
# Worker low: prefetch tinggi (batch friendly)
celery -A celery_app worker \
--queues=low,normal \
--concurrency=2 \
--prefetch-multiplier=16 \
--loglevel=info \
-n low@%h
# Trigger task dari API handler
from tasks import (
kirim_otp,
kirim_email_konfirmasi_order,
kirim_email_promo_batch,
)
# OTP — masuk queue urgent, di-pick worker dengan prefetch 1 → minimum delay
kirim_otp.delay(user_id=123, kode_otp="847291", channel="sms")
# Konfirmasi order — masuk queue normal
kirim_email_konfirmasi_order.delay(order_id=5678)
# Batch promo — masuk queue low
kirim_email_promo_batch.apply_async(
args=[kampanye_id],
countdown=60, # Delay 1 menit
)
# docker-compose.yml — worker per queue jadi container terpisah
services:
worker-urgent:
image: tokopedia-jobs:latest
command: celery -A celery_app worker -Q urgent -c 8 -P prefork --prefetch-multiplier=1
deploy:
replicas: 3
worker-normal:
image: tokopedia-jobs:latest
command: celery -A celery_app worker -Q normal -c 4 --prefetch-multiplier=4
worker-low:
image: tokopedia-jobs:latest
command: celery -A celery_app worker -Q low -c 2 --prefetch-multiplier=16
Kapan dipakai
- Aplikasi dengan mixed workload (real-time + batch).
- Notification system: OTP, push, email, dipisah priority.
- Report generation berat yang tidak boleh ganggu transactional task.
- Webhook delivery dengan retry — urgent webhook (payment) vs analytic webhook.
Catatan
- Queue separation lebih reliable dari priority field — Redis tidak benar-benar support priority dalam list. Multiple queue lebih predictable.
- prefetch_multiplier=1 untuk urgent — worker fetch 1 task at a time, fairness max. Default 4 bikin urgent task nunggu kalau worker lagi proses 3 task lainnya.
- acks_late=True wajib untuk reliability — task baru ack setelah selesai. Kalau worker crash, broker re-queue.
- soft_time_limit kirim signal sebelum hard kill. Worker bisa catch SoftTimeLimitExceeded dan cleanup gracefully.
- autoretry_for + backoff — built-in retry untuk transient error. Manual retry rumit, autoretry cukup untuk 90% kasus.
- Result backend opsional — kalau gak butuh result, set
ignore_result=Trueper task. Hemat Redis memory.
Hati-hati
worker_prefetch_multiplier=0— value-nya berarti “unlimited prefetch”, bukan 0. Bisa fetch jutaan task dan worker tidak release ke worker lain. Selalu set explicit angka.
# tags
celeryredisqueuebackground-jobpriority
Ditulis oleh Asti Larasati · 4 Juli 2026