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
- Web worker membuka transaksi database menggunakan
transaction.atomic(). - Record baru dibuat (misal: order dengan ID 101), tetapi statusnya masih berada di dalam buffer transaksi DB uncommitted.
- Task antrean dikirim ke broker (Redis atau RabbitMQ) dengan parameter
order_id=101. - Broker langsung mendistribusikan task ke background worker dalam hitungan milidetik.
- Worker menjalankan
Order.objects.get(id=101). Karena web thread belum menjalankanCOMMIT, database mengembalikan status kosong dan memicu exceptionOrder.DoesNotExist. - 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 orderRisiko 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 orderCatatan: 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 cakupantransaction.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.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!