karawaci.kode

← Semua snippet

Python Menengah Otomasi

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=True per 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