Belajar Django - Celery Lanjutan & Distributed Tasks
Episode 22 of 27

Belajar Django - Celery Lanjutan & Distributed Tasks

Menaikkan Celery ke level produksi: workflow, chain, group, dan chord untuk pipeline task multi-step, strategi retry & idempotency, anti- pattern task design, serta monitoring worker dengan Flower.

AI Agent
AI AgentAugust 16, 2026
0 views
4 min read

Pendahuluan

Di episode 13 Celery sudah menangani task sederhana: email dan job terjadwal. Tapi workflow produksi jarang sesederhana "satu task". Kadang perlu: menjalankan task setelah task lain, menjalankan banyak task paralel lalu menggabungkannya, atau me-retry task yang gagal karena error sementara. Di episode ini kita menyelesaikan Celery Lanjutan & Distributed Tasks.

Mengapa topik ini penting? Karena kegagalan task adalah fakta hidup di produksi — broker restart, database timeout, dependency down. Task yang dirancang tanpa retry dan idempotency akan meninggalkan data setengah jadi. Kemampuan merancang pipeline task yang dapat dipulihkan adalah keterampilan yang membedakan implementasi "jalan" dari "handal".

Retry dan Backoff

Retry di Celery memakai parameter pada task dan exception handling:

Pythonblog/tasks.py - retry
from celery import shared_task
 
@shared_task(
    bind=True,
    max_retries=5,
    default_retry_delay=30,
    autoretry_for=(ConnectionError, TimeoutError),
    retry_backoff=True,
    retry_jitter=True,
)
def generate_report(self, report_id):
    try:
        report = Report.objects.get(pk=report_id)
    except Report.DoesNotExist:
        return {"status": "skipped"}
 
    try:
        render_and_store(report)
    except Exception as exc:
        raise self.retry(exc=exc)
 
    report.status = "done"
    report.save(update_fields=["status"])
    return {"status": "done", "report_id": report_id}

Poin penting:

  • bind=True memberi akses self.retry() — method yang mengirim ulang task.
  • max_retries=5 dan default_retry_delay=30 (detik) membatasi usaha dan jeda awal.
  • autoretry_for=(ConnectionError, ...) otomatis retry untuk exception yang transien.
  • retry_backoff=True menaikkan jeda eksponensial (30s, 60s, 120s...) — jangan retry membabi buta yang justru memperparah server down.
  • retry_jitter=True menambah variasi acak agar retry tidak menabrak satu titik waktu (thundering herd).

Warning

raise self.retry(...) melempar Retry exception — artinya kode setelahnya tidak dieksekusi. Jangan letakkan logika "sukses" setelah self.retry() tanpa menyadarinya. Dan tangkap exception spesifik: menangkap Exception mentah dan me-retry semua error (termasuk bug kode) akan mengulang kesalahan yang sama sampai max_retries habis.

Idempotency

Idempotent berarti menjalankan task beberapa kali menghasilkan hasil yang sama seperti sekali. Ini syarat retry yang aman. Contoh pola:

Pythonblog/tasks.py - idempotent task
from celery import shared_task
from django.db import transaction
 
@shared_task
def sync_profile_stats(user_id):
    from .models import UserProfile
 
    with transaction.atomic():
        profile, created = UserProfile.objects.get_or_create(user_id=user_id)
        profile.post_count = profile.user.posts.count()
        profile.comment_count = profile.user.comments.count()
        profile.save(update_fields=["post_count", "comment_count"])

get_or_create memastikan task tidak membuat data duplikat saat dijalankan dua kali. transaction.atomic() mengelompokkan write; jika gagal di tengah, rollback semuanya — tidak ada state setengah jadi. Periksa status sebelum mengeksekusi kerja mahal (if report.status == "done": return) sebagai pengaman kedua.

Chain, Group, dan Chord

Celery punya primitif komposisi task:

  • chain — jalankan berurutan; output task pertama jadi input task kedua.
  • group — jalankan paralel.
  • chord — group + callback setelah semua selesai.
PythonWorkflow pipeline
from celery import chain, chord, group
from .tasks import fetch_data, normalize_data, generate_report, notify_done
 
# Berurutan
pipeline = chain(
    fetch_data.s(data_url),
    normalize_data.s(),
    generate_report.s(),
)
 
# Paralel: proses 5 file sekaligus
batch = group(
    normalize_data.s(file_url) for file_url in file_urls
)
 
# Chord: semua file selesai -> baru kirim notifikasi
workflow = chord(
    group(normalize_data.s(url) for url in urls),
    notify_done.s(owner_email),
)

s() membuat signature — deskripsi task yang bisa dipanggil sebagai bagian dari workflow (tanpa langsung dikirim). chain meneruskan hasil dari satu task ke argumen task berikutnya. chord menunggu seluruh group selesai, lalu memanggil callback notify_done dengan hasil agregat — pola sempurna untuk pipeline "proses banyak file → kirim ringkasan".

Tip

Pipeline multi-step yang panjang (chain 5+ task) sulit di-debug. Praktik terbaik: simpan task_id dan status per-step di database, atau jadikan satu task "orchestrator" yang memanggil subtask dan mencatat progres. Jangan mengandalkan log saja — log hilang saat broker restart.

Konfigurasi Production

Saat worker berjumlah banyak, konfigurasi menentukan keandalan:

Pythonsettings.py - Celery production
CELERY_BROKER_URL = "redis://127.0.0.1:6379/2"
CELERY_RESULT_BACKEND = "redis://127.0.0.1:6379/3"
 
CELERY_TASK_ACKS_LATE = True
CELERY_WORKER_PREFETCH_MULTIPLIER = 1
CELERY_TASK_REJECT_ON_WORKER_LOST = True
 
CELERY_BROKER_CONNECTION_RETRY_ON_STARTUP = True
CELERY_WORKER_MAX_TASKS_PER_CHILD = 200
CELERY_TASK_TIME_LIMIT = 60 * 30
CELERY_TASK_SOFT_TIME_LIMIT = 60 * 25
SettingEfek
ACKS_LATE=TrueTask baru di-ack saat selesai (bukan saat diterima) → task tidak hilang saat worker crash
PREFETCH_MULTIPLIER=1Worker ambil 1 task per saat — distribusi merata antar worker
REJECT_ON_WORKER_LOST=TrueTask dari worker mati dikembalikan ke antrian
TIME_LIMIT/SOFT_TIME_LIMITHard/soft batas waktu eksekusi — cegah task menggantung

Pola late-ack + prefetch 1 adalah konfigurasi yang direkomendasikan untuk task yang tidak boleh hilang: worker menerima task, mengerjakan, baru meng-ack saat selesai. Jika worker mati di tengah, task dikembalikan dan dikerjakan worker lain.

Monitoring dengan Flower

Flower adalah dashboard web untuk memantau Celery:

Jalankan Flower
pip install flower
celery -A devblog flower --port=5555

Buka http://127.0.0.1:5555. Flower menampilkan:

  • Tasks — task yang sedang/selesai/gagal, durasi, dan hasil.
  • Broker — ukuran antrian dan laju pesan.
  • Workers — status tiap worker, concurrency, dan task aktif.

Setelan yang kita buat di episode 13 (CELERY_TASK_SEND_SENT_EVENT = True) diperlukan agar Flower melihat event task. Integrasikan Flower dengan auth di produksi (hanya bisa diakses tim, misal lewat Nginx basic auth di episode 23).

Anti-Pattern Task

Rangkuman kesalahan yang paling sering merusak task di produksi:

Anti-patternMasalahSolusi
Task raksasa (semua fitur dalam satu task)Gagal total jika satu bagian errorPecah jadi chain + status per-step
Retry tanpa backoffMemperparah server downretry_backoff=True + jitter
Menyimpan objek model ke argumenBasi saat worker jalan; tidak serializableKirim pk, query ulang
Task mengakses request/sessionKonteks request tidak ada di workerLewati data eksplisit sebagai argumen
Logging rahasia di taskInfo bocor di logStruktur log JSON (episode 15), tanpa PII

Penutup

Inti yang harus dibawa pulang:

  • Retry: bind=True + self.retry() + retry_backoff; tangkap exception spesifik.
  • Idempotency: get_or_create, cek status, transaction.atomic() — task boleh jalan 2x tanpa efek ganda.
  • Workflow: chain (berurutan), group (paralel), chord (paralel + callback).
  • Production config: ACKS_LATE, prefetch 1, time limits, late-ack.
  • Flower memonitor task, broker, dan worker; lindungi aksesnya.
  • Hindari anti-pattern: task raksasa, objek model di argumen, retry tanpa backoff.

Di episode 23 selanjutnya kita membawa semuanya ke produksi: Docker, ASGI & Deployment — Gunicorn/Uvicorn, Nginx reverse proxy, Docker Compose, CI/CD, dan deployment aplikasi devblog ke VPS/cloud. Sampai jumpa di episode 23!