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_nackdengan opsirequeue: truepada 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:
- Transient Error: Kegagalan sementara (koneksi DB putus sesaat, timeout jaringan). NACK dengan
requeue: truedapat digunakan secara terbatas, atau kirim ke delay queue. - Permanent / Poison Error: Payload korup, validasi schema gagal, constraint violation permanen. Kirim NACK dengan
requeue: falsedan 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_unacknowledgedvia Prometheus. Jika nilainya sama persis dengan total batasprefetch_count * jumlah_workerdan 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
heartbeatpadaConnectionPropertiesdiaktifkan. 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.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!