Galat ObjectDoesNotExist atau DoesNotExist pada background worker seperti Celery atau RQ sering terjadi bukan karena bug pada query, melainkan akibat race condition antara transaksi database Django dan message broker. Worker mengambil pesan dari antrean dan mengeksekusinya sebelum database selesai memproses perintah COMMIT dari web thread.

Akar Masalah: Asinkronitas Antrean vs Transaksi Database

Saat aplikasi memproses mutasi data di dalam blok transaksi, perubahan data tersebut belum persisten bagi thread atau proses lain hingga transaksi tuntas dilakukan (level isolasi database standar seperti Read Committed). Masalah muncul ketika pesan diserahkan ke broker antrean sebelum transaksi tersebut selesai.

Kronologi Race Condition

  1. Web worker membuka transaksi database menggunakan transaction.atomic().
  2. Record baru dibuat (misal: order dengan ID 101), tetapi statusnya masih berada di dalam buffer transaksi DB uncommitted.
  3. Task antrean dikirim ke broker (Redis atau RabbitMQ) dengan parameter order_id=101.
  4. Broker langsung mendistribusikan task ke background worker dalam hitungan milidetik.
  5. Worker menjalankan Order.objects.get(id=101). Karena web thread belum menjalankan COMMIT, database mengembalikan status kosong dan memicu exception Order.DoesNotExist.
  6. Web thread akhirnya menyelesaikan eksekusi blok dan melakukan COMMIT, terlambat beberapa milidetik setelah worker gagal.

Anti-Pattern: Dispatch Task di Dalam Blok atomic()

Memanggil fungsi .delay() atau apply_async() secara langsung di dalam konteks transaction.atomic() merupakan anti-pattern yang berbahaya:

from django.db import transaction
from orders.models import Order
from orders.tasks import process_order

def create_order(user, data):
    with transaction.atomic():
        order = Order.objects.create(user=user, **data)
        # ANTI-PATTERN: Task dikirim sebelum commit
        process_order.delay(order.id)
        
        # Operasi DB tambahan berpotensi menunda commit
        generate_invoice_records(order)
    return order

Risiko lain dari pola ini adalah pembatalan transaksi (rollback). Jika generate_invoice_records() melempar error, transaksi database di-rollback sehingga record order tidak pernah tersimpan. Namun, task Celery sudah terlanjur masuk ke broker dan tetap dieksekusi, menghasilkan operasi pada data hantu.

Solusi: Menggunakan transaction.on_commit()

Django menyediakan fungsi bawaan django.db.transaction.on_commit(). Fungsi ini menerima callable yang eksekusinya ditunda hingga transaksi terluar berhasil di-commit ke database. Jika transaksi di-rollback, callback akan otomatis dibatalkan.

from django.db import transaction
from orders.models import Order
from orders.tasks import process_order

def create_order(user, data):
    with transaction.atomic():
        order = Order.objects.create(user=user, **data)
        generate_invoice_records(order)
        
        # SOLUSI: Menjamin commit sukses sebelum push ke broker
        transaction.on_commit(lambda: process_order.delay(order.id))
        
    return order
Catatan: Jika transaction.on_commit() dipanggil di luar transaksi aktif, callback akan langsung dieksekusi seketika.

Implementasi Produksi: Views dan Django Signals

Berikut adalah implementasi clean architecture yang memisahkan dispatch logic pada layer views dan signals dengan aman.

1. Implementasi pada View / Service Layer

from rest_framework.views import APIView
from rest_framework.response import Response
from rest_framework import status
from django.db import transaction
from orders.models import Order
from orders.tasks import process_order
import logging

logger = logging.getLogger(__name__)

class OrderCreateView(APIView):
    def post(self, request):
        serializer = OrderSerializer(data=request.data)
        serializer.is_valid(raise_exception=True)
        
        with transaction.atomic():
            order = serializer.save()
            # Registrasi dispatch task setelah DB commit
            transaction.on_commit(
                lambda: self._dispatch_task(order.id)
            )
            
        return Response(serializer.data, status=status.HTTP_201_CREATED)

    @staticmethod
    def _dispatch_task(order_id):
        try:
            process_order.delay(order_id)
        except Exception as exc:
            # Menangani potensi kegagalan broker (misal: Redis timeout)
            logger.error("Gagal mendistribusikan task untuk Order %s: %s", order_id, exc)

2. Implementasi pada Signals (post_save)

Sinyal post_save tetap dipanggil di dalam konteks transaksi jika operasi penyimpanan terjadi di dalam atomic(). Jangan langsung mendispatch worker dari sinyal tanpa membungkusnya dengan on_commit.

from django.db.models.signals import post_save
from django.dispatch import receiver
from django.db import transaction
from orders.models import Order
from orders.tasks import process_order

@receiver(post_save, sender=Order)
def handle_order_post_save(sender, instance, created, **kwargs):
    if created:
        # Jangan langsung instance task di sini
        transaction.on_commit(lambda: process_order.delay(instance.id))

Menangani Kegagalan Broker pada Callback on_commit

Kelemahan pola on_commit: transaksi database sudah tuntas di-commit, namun broker pesan bisa saja mengalami downtime tepat saat callback dieksekusi. Callback gagal melempar task, sehingga order tidak pernah diproses oleh worker.

Pola Mitigasi: Status State Machine

Terapkan kolom status pada model untuk mendeteksi order yang menggantung akibat kegagalan broker, dipadukan dengan periodic task (Celery Beat) sebagai mekanisme rekonsiliasi cadangan.

# tasks.py
from celery import shared_task
from django.utils import timezone
from datetime import timedelta
from orders.models import Order

@shared_task
def reconcile_stale_orders():
    # Menangkap order yang tertahan lebih dari 5 menit tanpa diproses
    threshold = timezone.now() - timedelta(minutes=5)
    stale_orders = Order.objects.filter(status=Order.Status.PENDING, created_at__lte=threshold)
    
    for order in stale_orders:
        process_order.delay(order.id)

Idempotency Worker dengan Redis Lock (django-redis)

Mekanisme retry otomatis pada Celery atau rekonsiliasi periodik dapat memicu eksekusi task ganda untuk record yang sama. Worker harus dirancang idempoten menggunakan distributed lock dari django-redis.

from celery import shared_task
from django.core.cache import cache
from django.db import transaction
from orders.models import Order
import logging

logger = logging.getLogger(__name__)

@shared_task(bind=True, max_retries=3, default_retry_delay=10)
def process_order(self, order_id):
    lock_key = f"lock:order:process:{order_id}"
    # Lock timeout 60 detik mencegah deadlock jika worker crash mendadak
    acquire_lock = cache.add(lock_key, "locked", timeout=60)
    
    if not acquire_lock:
        logger.warning("Task untuk order_id %s sedang diproses oleh worker lain.", order_id)
        return

    try:
        with transaction.atomic():
            # Menggunakan select_for_update untuk mencegah concurrent updates di level database
            order = Order.objects.select_for_update().get(id=order_id)
            
            if order.status != Order.Status.PENDING:
                logger.info("Order %s sudah diproses sebelumnya (Idempotent bypass).", order_id)
                return
            
            # Jalankan logika bisnis
            order.process_payment()
            order.status = Order.Status.PROCESSED
            order.save()
            
    except Order.DoesNotExist:
        logger.error("Record Order %s tidak ditemukan.", order_id)
    except Exception as exc:
        logger.exception("Gagal memproses Order %s: %s", order_id, exc)
        raise self.retry(exc=exc)
    finally:
        # Lepaskan lock setelah proses selesai
        cache.delete(lock_key)

Ringkasan Best Practice

  • Hindari eksekusi .delay() langsung di dalam cakupan transaction.atomic().
  • Gunakan transaction.on_commit(callback) di views maupun signals untuk memastikan record database sudah stabil dan visible bagi worker.
  • Selalu kirim ID record primitif (integer atau UUID) melalui parameter task, bukan serialisasi model instance yang sudah stale.
  • Bungkus eksekusi worker kritis dengan distributed lock via Redis untuk menjaga idempotency dari duplikasi eksekusi task.