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.

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 di Celery memakai parameter pada task dan exception handling:
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.
Idempotent berarti menjalankan task beberapa kali menghasilkan hasil yang sama seperti sekali. Ini syarat retry yang aman. Contoh pola:
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.
Celery punya primitif komposisi task:
chain — jalankan berurutan; output task pertama jadi input task kedua.group — jalankan paralel.chord — group + callback setelah semua selesai.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.
Saat worker berjumlah banyak, konfigurasi menentukan keandalan:
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| Setting | Efek |
|---|---|
ACKS_LATE=True | Task baru di-ack saat selesai (bukan saat diterima) → task tidak hilang saat worker crash |
PREFETCH_MULTIPLIER=1 | Worker ambil 1 task per saat — distribusi merata antar worker |
REJECT_ON_WORKER_LOST=True | Task dari worker mati dikembalikan ke antrian |
TIME_LIMIT/SOFT_TIME_LIMIT | Hard/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.
Flower adalah dashboard web untuk memantau Celery:
pip install flower
celery -A devblog flower --port=5555Buka http://127.0.0.1:5555. Flower menampilkan:
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).
Rangkuman kesalahan yang paling sering merusak task di produksi:
| Anti-pattern | Masalah | Solusi |
|---|---|---|
| Task raksasa (semua fitur dalam satu task) | Gagal total jika satu bagian error | Pecah jadi chain + status per-step |
| Retry tanpa backoff | Memperparah server down | retry_backoff=True + jitter |
| Menyimpan objek model ke argumen | Basi saat worker jalan; tidak serializable | Kirim pk, query ulang |
Task mengakses request/session | Konteks request tidak ada di worker | Lewati data eksplisit sebagai argumen |
| Logging rahasia di task | Info bocor di log | Struktur log JSON (episode 15), tanpa PII |
Inti yang harus dibawa pulang:
bind=True + self.retry() + retry_backoff; tangkap exception spesifik.get_or_create, cek status, transaction.atomic() — task boleh jalan 2x tanpa efek ganda.chain (berurutan), group (paralel), chord (paralel + callback).ACKS_LATE, prefetch 1, time limits, late-ack.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!