Penyebab Utama Job Stalling pada BullMQ

Job stalling terjadi ketika Redis kehilangan sinyal aktif dari worker yang sedang memproses antrean. BullMQ mengunci setiap job menggunakan token unik dan TTL (Time-To-Live) tertentu pada Redis key. Jika pembaruan token gagal dilakukan sebelum TTL habis, master process mengasumsikan worker telah mati, menandai status job sebagai stalled, dan menjadwalkan ulang job tersebut ke worker lain. Pada backend Nitro Nuxt 3, anomali ini dipicu oleh dua akar masalah operasional:

  • Terminasi Container Tanpa Koordinasi (SIGTERM): Saat rolling deployment di Kubernetes atau Docker, orchestrator mengirim sinyal SIGTERM. Jika proses Nitro langsung mati tanpa menunggu siklus job selesai, lock Redis tidak dilepas secara bersih dan proses terhenti di tengah jalan.
  • Beban Event Loop dan Default Lock Duration: BullMQ menggunakan nilai bawaan lockDuration sebesar 30 detik. Jika task melakukan kalkulasi CPU-heavy atau operasi I/O yang memblokir event loop Node.js melebihi durasi tersebut, perpanjangan lock otomatis (lock renewal) gagal dikirim ke Redis.

Konfigurasi Worker BullMQ pada Nitro Plugin

Inisialisasi BullMQ worker harus ditempatkan pada server plugin Nitro agar siklus hidupnya terikat pada runtime server Nuxt 3. Konfigurasikan parameter lockDuration, stalledInterval, dan maxStalledCount sesuai karakteristik beban kerja.

// server/plugins/bullmq.ts
import { Worker, type Job } from 'bullmq'
import IORedis from 'ioredis'

interface EmailPayload {
  to: string
  subject: string
  idempotencyKey: string
}

export default defineNitroPlugin((nitroApp) => {
  const config = useRuntimeConfig()

  const connection = new IORedis(config.redisUrl, {
    maxRetriesPerRequest: null,
    enableReadyCheck: false
  })

  // ponytail: lockDuration 60s covers slow SMTP gateways; upgrade to progress heartbeats if tasks exceed 2m.
  const emailWorker = new Worker<EmailPayload>(
    'email-queue',
    async (job: Job<EmailPayload>) => {
      await processEmailJob(job)
    },
    {
      connection,
      concurrency: 5,
      lockDuration: 60000,   // 60 detik batas lock TTL
      stalledInterval: 30000, // Cek status stalled setiap 30 detik
      maxStalledCount: 2      // Maksimal toleransi retry jika job stalled
    }
  )

  emailWorker.on('stalled', (jobId) => {
    console.warn(`[BullMQ] Job ${jobId} terdeteksi stalled dan dijadwalkan ulang.`)
  })

  emailWorker.on('failed', (job, err) => {
    console.error(`[BullMQ] Job ${job?.id} gagal permanen:`, err.message)
  })
})

Implementasi Graceful Shutdown

Saat Nitro menerima sinyal terminasi sistem operasi, koneksi tidak boleh diputus secara instan. Worker harus berhenti mengambil job baru dari queue, merampungkan job yang sedang berjalan di memori, melepaskan lock Redis, lalu menutup koneksi jaringan. Nitro menyediakan hook siklus hidup close untuk menangani skenario ini.

// server/plugins/bullmq.ts (lanjutan)
  nitroApp.hooks.hook('close', async () => {
    console.info('[Shutdown] Menutup BullMQ worker secara aman...')
    
    // worker.close() menunggu eksekusi job yang sedang aktif hingga selesai
    await emailWorker.close()
    
    // Putus koneksi Redis setelah worker sepenuhnya idle
    await connection.quit()
    
    console.info('[Shutdown] BullMQ worker dan koneksi Redis berhasil dimatikan.')
  })
Penting: Pastikan parameter terminationGracePeriodSeconds pada spesifikasi Pod Kubernetes disetel lebih besar daripada kombinasi lockDuration dan perkiraan eksekusi job terlama (misalnya 90 detik). Jika grace period habis, sistem akan mengirim SIGKILL yang mematikan proses secara paksa.

Idempotency Key: Mencegah Eksekusi Ganda Saat Auto-Retry

Ketika worker crash akibat Out-Of-Memory (OOM) atau kegagalan jaringan mendadak, mekanisme graceful shutdown tidak sempat dijalankan. BullMQ akan menjadwalkan ulang job yang stalled secara otomatis. Tanpa penanganan status idempotensi, tindakan eksternal seperti pemotongan saldo atau pengiriman email akan tereksekusi dua kali.

Gunakan Redis atomic lock berbasis payload idempotencyKey untuk memverifikasi eksekusi sebelum payload utama diproses:

// server/utils/queueProcessor.ts
import type { Job } from 'bullmq'
import IORedis from 'ioredis'

const redisClient = new IORedis(process.env.REDIS_URL!)

export async function processEmailJob(job: Job<{ to: string; idempotencyKey: string }>) {
  const { idempotencyKey, to } = job.data
  const lockKey = `idempotency:${idempotencyKey}`

  // Atomic set if not exists dengan TTL 24 jam
  const acquired = await redisClient.set(lockKey, 'PROCESSING', 'EX', 86400, 'NX')

  if (!acquired) {
    const currentStatus = await redisClient.get(lockKey)
    if (currentStatus === 'COMPLETED') {
      console.info(`[Idempotency] Job ${job.id} sudah pernah sukses. Lewati proses.`)
      return
    }
    if (currentStatus === 'PROCESSING') {
      throw new Error(`[Idempotency] Job ${job.id} sedang diproses worker lain.`)
    }
  }

  try {
    // Jalankan operasi non-idempoten (contoh: pemanggilan API eksternal)
    await sendMailViaProvider(to)

    // Tandai status selesai permanen
    await redisClient.set(lockKey, 'COMPLETED', 'KEEPTTL')
  } catch (error) {
    // Lepas lock jika gagal agar worker berikutnya bisa melakukan retry
    await redisClient.del(lockKey)
    throw error
  }
}

async function sendMailViaProvider(recipient: string) {
  // Implementasi provider SMTP atau third-party API
}

Panduan Debugging Job Stalled

Gunakan langkah-langkah berikut ketika job masih terdeteksi stalled berulang kali di monitoring Redis:

  1. Periksa Latensi Event Loop: Gunakan monitoring APM untuk memastikan waktu event loop delay di Node.js berada di bawah 100ms. Latensi tinggi mencegah pengiriman perintah lock renewal.
  2. Cek Alokasi Memori Container: Periksa log sistem menggunakan perintah dmesg -T | grep -i oom untuk memverifikasi apakah worker mati akibat pembatasan RAM container.
  3. Evaluasi Durasi Job: Jika job memerlukan waktu eksekusi dinamis yang sangat lama, panggil job.updateProgress() secara berkala di dalam kode prosesor untuk memicu pembaruan lock ke server Redis secara proaktif.