Penumpukan unacked message pada RabbitMQ terjadi ketika consumer mengambil pesan dari queue tetapi tidak mengirimkan konfirmasi (ACK atau NACK) kembali ke broker. Pada arsitektur Actix Web dan Rust, masalah ini umumnya muncul akibat runtime panic pada task worker, pemakaian operator ? yang keluar lebih awal tanpa blok acknowledgment, atau ketiadaan batas prefetch count. Akibatnya, broker menahan pesan dalam RAM, metrik konsumsi macet, dan throughput antrean menurun drastis.

Akar Masalah: Mengapa Unacked Message Menumpuk?

RabbitMQ secara default mendistribusikan pesan menggunakan mode push. Ketika koneksi consumer aktif, broker akan terus mengirimkan pesan yang tersedia ke memori lokal worker jika batas kapasitas tidak ditentukan.

  • Default Prefetch Tak Terbatas: Tanpa basic_qos, RabbitMQ mengirim seluruh isi antrean ke satu worker instance. Jika instance tersebut mengalami deadlock pada panggilan database (misal SQLx pool exhaustion), ribuan pesan tertahan dalam status Unacknowledged dan worker lain tidak dapat mengambilnya.
  • Unhandled Error Exit: Penggunaan error propagation (?) di dalam worker loop yang keluar dari scope pemrosesan sebelum memanggil method ACK atau NACK. Pesan tetap berstatus unacked hingga TCP connection ditutup.
  • Infinite Requeue Loop: Penggunaan basic_nack dengan opsi requeue: true pada error deterministik (seperti payload JSON rusak) menyebabkan pesan masuk kembali ke queue, langsung ditarik ulang, gagal kembali, dan membebani resource I/O tanpa henti.

Arsitektur Worker Loop: tokio::spawn Berdampingan dengan Actix Web

Actix Web menjalankan runtime multi-threaded untuk HTTP server. Worker RabbitMQ tidak boleh memblokir thread HTTP worker Actix. Pola yang benar adalah memisahkan connection initialization ke RabbitMQ, menjalankan worker loop melalui tokio::spawn, dan menyematkan status koneksi ke dalam web::Data untuk kebutuhan endpoint pemantauan.

use actix_web::{web, App, HttpServer, HttpResponse, Responder};
use lapin::{Connection, ConnectionProperties, options::*, types::FieldTable};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;

#[derive(Clone)]
pub struct AppState {
    pub rabbit_connected: Arc<AtomicBool>,
}

async fn health_check(state: web::Data<AppState>) -> impl Responder {
    if state.rabbit_connected.load(Ordering::Relaxed) {
        HttpResponse::Ok().json(serde_json::json!({"status": "healthy", "consumer": "up"}))
    } else {
        HttpResponse::ServiceUnavailable().json(serde_json::json!({"status": "unhealthy", "consumer": "down"}))
    }
}

Solusi Teknis 1: Membatasi Prefetch Count dengan basic_qos

Atur batas prefetch pada AMQP channel sebelum consumer dibuat. Prefetch count menentukan berapa banyak pesan unacked yang diizinkan berada di memori worker secara bersamaan. Nilai ideal berkisar antara 10 hingga 50 per channel, bergantung pada durasi pemrosesan tiap unit pesan.

use lapin::options::BasicQosOptions;

// Batasi hanya menerima maksimal 20 unacked message sekaligus
channel
    .basic_qos(20, BasicQosOptions::default())
    .await
    .expect("Gagal mengonfigurasi basic_qos");

Tindakan ini mencegah worker menimbun pesan di RAM ketika terjadi lonjakan beban atau pelambatan eksternal pada query SQLx.

Solusi Teknis 2 & 3: Crash-Safe ACK/NACK dan Dead-Letter Exchange (DLX)

Untuk menghindari unacked message akibat error internal sekaligus mencegah poison message loop, pisahkan penanganan error menjadi dua kategori:

  1. Transient Error: Kegagalan sementara (koneksi DB putus sesaat, timeout jaringan). NACK dengan requeue: true dapat digunakan secara terbatas, atau kirim ke delay queue.
  2. Permanent / Poison Error: Payload korup, validasi schema gagal, constraint violation permanen. Kirim NACK dengan requeue: false dan arahkan ke Dead Letter Exchange (DLX) via konfigurasi queue.

Deklarasi Queue dengan Argument Dead-Letter

let mut queue_args = FieldTable::default();
queue_args.insert(
    "x-dead-letter-exchange".into(),
    lapin::types::AMQPValue::LongString("dlx.events".into()),
);
queue_args.insert(
    "x-dead-letter-routing-key".into(),
    lapin::types::AMQPValue::LongString("events.dead".into()),
);

channel
    .queue_declare(
        "main_task_queue",
        QueueDeclareOptions { durable: true, ..Default::QueueDeclareOptions() },
        queue_args,
    )
    .await?;

Implementasi Worker Loop Crash-Safe

use lapin::options::{BasicAckOptions, BasicNackOptions, BasicConsumeOptions};
use futures_util::StreamExt;
use serde::Deserialize;

#[derive(Deserialize)]
struct TaskPayload {
    id: i64,
    action: String,
}

enum ProcessResult {
    Success,
    TransientError(String),
    PoisonError(String),
}

async fn process_delivery(data: &[u8]) -> ProcessResult {
    // 1. Parsing JSON (Poison Error jika format salah)
    let payload: TaskPayload = match serde_json::from_slice(data) {
        Ok(p) => p,
        Err(e) => return ProcessResult::PoisonError(e.to_string()),
    };

    // 2. Simulasi eksekusi database/business logic
    if payload.id == 0 {
        return ProcessResult::TransientError("Database timeout".into());
    }

    ProcessResult::Success
}

pub async fn start_consumer_worker(
    channel: lapin::Channel,
    connected_flag: Arc<AtomicBool>,
) {
    let mut consumer = channel
        .basic_consume(
            "main_task_queue",
            "actix_worker_tag",
            BasicConsumeOptions::default(),
            FieldTable::default(),
        )
        .await
        .expect("Gagal mendaftarkan consumer");

    connected_flag.store(true, Ordering::Relaxed);

    while let Some(delivery_result) = consumer.next().await {
        match delivery_result {
            Ok(delivery) => {
                let result = process_delivery(&delivery.data).await;

                match result {
                    ProcessResult::Success => {
                        if let Err(e) = delivery.ack(BasicAckOptions::default()).await {
                            eprintln!("Gagal mengirim ACK: {e}");
                        }
                    }
                    ProcessResult::TransientError(reason) => {
                        eprintln!("Transient error ({reason}), requeue message");
                        // Requeue hanya untuk error sementara
                        let _ = delivery
                            .nack(BasicNackOptions { requeue: true, ..Default::BasicNackOptions() })
                            .await;
                    }
                    ProcessResult::PoisonError(reason) => {
                        eprintln!("Poison message ({reason}), route to DLX");
                        // Requeue false mengarahkan pesan ke DLX karena x-dead-letter-exchange disetel
                        let _ = delivery
                            .nack(BasicNackOptions { requeue: false, ..Default::BasicNackOptions() })
                            .await;
                    }
                }
            }
            Err(err) => {
                eprintln!("Consumer delivery stream error: {err}");
                connected_flag.store(false, Ordering::Relaxed);
                break;
            }
        }
    }
    connected_flag.store(false, Ordering::Relaxed);
}

Panduan Integrasi Main Binary

Jalankan consumer sebelum memulai HTTP server pada thread main:

#[actix_web::main]
async fn main() -> std::io::Result<()> {
    let rabbit_uri = std::env::var("AMQP_ADDR").unwrap_or_else(|_| "amqp://guest:[email protected]:5672".into());
    let conn = Connection::connect(&rabbit_uri, ConnectionProperties::default())
        .await
        .expect("Koneksi AMQP gagal");

    let channel = conn.create_channel().await.expect("Gagal membuat channel");
    channel.basic_qos(10, lapin::options::BasicQosOptions::default()).await.unwrap();

    let rabbit_status = Arc::new(AtomicBool::new(false));

    // Spawn background worker di runtime Tokio
    let worker_status = rabbit_status.clone();
    tokio::spawn(async move {
        start_consumer_worker(channel, worker_status).await;
    });

    let state = web::Data::new(AppState { rabbit_connected: rabbit_status });

    HttpServer::new(move || {
        App::new()
            .app_data(state.clone())
            .route("/health", web::get().to(health_check))
    })
    .bind(("0.0.0.0", 8080))?
    .run()
    .await
}

Observabilitas dan Monitoring RabbitMQ

Gunakan tooling observabilitas berikut untuk mendeteksi message leak sebelum mempengaruhi SLA produksi:

  • Metrik Kunci: Pantau rabbitmq_queue_messages_unacknowledged via Prometheus. Jika nilainya sama persis dengan total batas prefetch_count * jumlah_worker dan bertahan konstan, worker mengalami hung/deadlock.
  • Grafana Alert: Pasang alert ketika persentase pesan unacked melampaui 80% dari total prefetch limit selama lebih dari 3 menit berturut-turut.
  • Koneksi TCP & Heartbeat: Pastikan parameter heartbeat pada ConnectionProperties diaktifkan. Jika worker process mati mendadak (SIGKILL / Out of Memory), broker akan mendeteksi hilangnya heartbeat dan secara otomatis me-requeue unacked message yang tertahan pada channel terkait.