Tantangan Error Handling pada Queue Worker Rust
Membangun sistem pemrosesan antrean (queue worker) berbasis Rust menuntut kontrol konkurensi dan keandalan data yang tinggi. Masalah umum yang kerap terjadi di lingkungan produksi adalah hilangnya pesan (message drop) akibat pengakuan prematur (early ack), serta kegagalan sistemik akibat panic yang menyebar ke seluruh worker thread atau meracuni shared state (Mutex poisoning).
Untuk mengatasi kendala ini, worker Rust harus membedakan klasifikasi kegagalan secara eksplisit antara kegagalan sementara (transient error) yang memerlukan mekanisme retry backoff, dan kegagalan permanen yang harus langsung dialihkan ke Dead-Letter Queue (DLQ). Pendekatan terstruktur ini menjamin ketersediaan thread pool dan integritas antrean data.
Domain Error Terstruktur: Memisahkan Transient dan Permanent Error
Penggunaan string mentah atau penanganan error generik menyulitkan pengambilan keputusan di level konsumen queue. Gunakan crate thiserror untuk mendefinisikan tipe error domain yang secara tegas membedakan sifat kegagalan.
use thiserror::Error;
#[derive(Error, Debug)]
pub enum JobError {
#[error("Transient network error: {0}")]
TransientNetwork(String),
#[error("Database connection timed out")]
DatabaseTimeout,
#[error("Validation failure: {0}")]
ValidationError(String),
#[error("Corrupted payload schema")]
InvalidPayload,
}
impl JobError {
pub fn is_retryable(&self) -> bool {
match self {
JobError::TransientNetwork(_) | JobError::DatabaseTimeout => true,
JobError::ValidationError(_) | JobError::InvalidPayload => false,
}
}
}Dengan mengekspos method is_retryable(), logika orkestrator queue dapat menentukan status pesan tanpa perlu melakukan parsing string manual.
Mencegah Mutex Poisoning dan Mengisolasi Task Crash
Di Rust, jika sebuah thread mengalami panic! saat memegang std::sync::MutexGuard, mutex tersebut akan masuk ke status poisoned. Pemanggilan .lock().unwrap() pada thread lain secara berantai akan memicu cascade failure di seluruh worker.
1. Isolasi Task Menggunakan Tokio Spawn
Pada runtime asynchronous (Tokio), gunakan tokio::spawn untuk mengisolasi eksekusi tugas. Kegagalan panic dalam task terisolasi dapat ditangkap melalui JoinHandle tanpa mematikan thread runner utama.
use tokio::task;
let task_handle = task::spawn(async move {
process_payload(data).await
});
match task_handle.await {
Ok(job_result) => match job_result {
Ok(_) => ack_message().await,
Err(err) => handle_job_error(err).await,
},
Err(join_err) if join_err.is_panic() => {
tracing::error!("Worker task panicked! Isolating error.");
route_to_dlq_raw(raw_data).await;
}
Err(join_err) => {
tracing::error!(%join_err, "Task cancelled or aborted");
}
}2. Menangani Mutex Poisoning dengan Safe Recovery
Jika state bersama berbasis std::sync::Mutex tetap dibutuhkan, hindari penggunaan .unwrap(). Gunakan metode recovery fallback saat kondisi poisoning terdeteksi.
use std::sync::{Arc, Mutex};
fn update_metrics(shared_counter: &Arc<Mutex<u64>>) {
let mut counter = match shared_counter.lock() {
Ok(guard) => guard,
Err(poisoned) => {
tracing::warn!("Shared mutex was poisoned. Recovering inner state.");
poisoned.into_inner()
}
};
*counter += 1;
}Catatan: Untuk beban kerja async, prioritaskan penggunaantokio::sync::Mutexuntuk alur kerja yang melintasi titik.await, atau gunakan atomics (AtomicU64) yang bebas dari risiko poisoning.
Pipeline Eksekusi: Safe Routing ke DLQ dan Backoff
Berikut adalah implementasi pipeline pemrosesan pesan berbasis Redis Stream atau AMQP yang menangani retry secara terukur, routing DLQ, dan instrumentasi log via tracing.
use std::time::Duration;
use tracing::{error, info, instrument, warn};
struct JobContext {
id: String,
attempt: u32,
max_retries: u32,
payload: Vec<u8>,
}
#[instrument(skip(job, db_pool), fields(job_id = %job.id, attempt = job.attempt))]
async fn process_job_pipeline(job: JobContext, db_pool: &DbPool) -> Result<(), eyre::Report> {
match execute_job(&job.payload, db_pool).await {
Ok(_) => {
info!("Job processed successfully");
ack_message(&job.id).await?;
}
Err(err) if err.is_retryable() && job.attempt < job.max_retries => {
let backoff = Duration::from_secs(2_u64.pow(job.attempt));
warn!(%err, retry_after_secs = backoff.as_secs(), "Transient failure, retrying");
tokio::time::sleep(backoff).await;
requeue_message(&job.id, job.attempt + 1).await?;
}
Err(err) => {
error!(%err, "Job failed permanently or exceeded retry limit. Sending to DLQ");
send_to_dlq(&job.id, &job.payload, &err.to_string()).await?;
ack_message(&job.id).await?;
}
}
Ok(())
}Audit Log dan Konteks Tracing
Gunakan macro tracing::instrument untuk melacak jalur eksekusi pesan di seluruh worker thread. Ini memastikan correlation ID tetap terbawa saat terjadi error fatal tanpa kehilangan konteks eksekusi.
- Catat alasan DLQ: Jangan kirim payload mentah ke DLQ tanpa menyertakan metadata error asli dan jumlah percobaan terakhir.
- Hindari Auto-Ack: Lakukan acknowledgment ke broker queue hanya setelah status akhir job ditentukan (sukses atau teralihkan ke DLQ). Membuka ACK di awal siklus merupakan penyebab utama message drop saat worker crash.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!