Anatomi Masalah: Cache Stampede dan Thundering Herd

Cache stampede (dikenal juga sebagai thundering herd) terjadi saat data populer (hot key) pada cache in-memory atau Redis kedaluwarsa (TTL expired) persis di puncak lalu lintas. Ketika ratusan atau ribuan request masuk secara bersamaan, seluruh request tersebut mendapati kondisi cache miss secara serentak.

Dampak operasionalnya langsung terasa pada infrastruktur backend:

  • Database Pool Exhaustion: Setiap request konkuren membuka koneksi DB untuk menjalankan query yang identik. Connection pool (misal: sqlx, deadpool, atau r2d2) terkuras habis dalam hitungan milidetik.
  • Lonjakan Latensi P99: Request baru tertahan mengantre koneksi pool hingga menyentuh batas HTTP timeout.
  • Cascading Failure: Database mengalami lonjakan CPU 100%, IOPS habis, dan memicu penolakan koneksi sistemik yang menjatuhkan service lain.

Kelemahan Solusi Naif: Global Mutex di Async Runtime

Upaya umum untuk mencegah duplikasi fetch adalah mengunci blok pembaruan cache dengan mutex. Pada async runtime seperti Tokio di Actix Web, strategi ini menimbulkan masalah baru:

  • Pemblokiran Thread (std::sync::Mutex): Menggunakan blocking mutex standar di dalam future async akan menahan thread worker OS milik Tokio runtime. Ini memicu starvation pada task-task lain yang dijadwalkan di thread yang sama.
  • Serialisasi Global (tokio::sync::Mutex): Penggunaan async mutex tunggal secara naif memaksa seluruh request—bahkan untuk key yang berbeda—menunggu dalam antrean FIFO. Throughput HTTP drop drastis dan latensi melonjak karena konkurensi dimatikan.

Pola Singleflight: Request Coalescing

Pola Singleflight (request coalescing) menyelesaikan masalah ini pada tingkat per-key:

  1. Request pertama untuk key yang expired bertindak sebagai leader. Leader mencatat key ke dalam in-flight registry dan mengeksekusi query database.
  2. Request berikutnya untuk key yang sama mendeteksi bahwa proses fetch sedang berjalan, lalu bertindak sebagai follower. Follower tidak menyentuh database, melainkan subscribe ke channel hasil milik leader.
  3. Setelah query leader selesai, hasilnya di-broadcast ke semua follower secara serentak, lalu key dihapus dari in-flight registry.

Implementasi Singleflight dengan Tokio Broadcast

Implementasi berikut menggunakan tokio::sync::broadcast untuk mendistribusikan hasil query satu kali ke seluruh subscriber konkuren.

use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::{broadcast, Mutex};
use tokio::time::{timeout, Duration};

#[derive(Clone)]
pub struct SingleFlight<T> {
    in_flight: Arc<Mutex<HashMap<String, broadcast::Sender<Result<T, String>>>>>,
}

impl<T: Clone + Send + Sync + 'static> SingleFlight<T> {
    pub fn new() -> Self {
        Self {
            in_flight: Arc::new(Mutex::new(HashMap::new())),
        }
    }

    pub async fn execute<F, Fut>(&self, key: &str, fetcher: F) -> Result<T, String>
    where
        F: FnOnce() -> Fut,
        Fut: std::future::Future<Output = Result<T, String>>,
    {
        let mut rx = {
            let mut map = self.in_flight.lock().await;

            if let Some(sender) = map.get(key) {
                // Task lain sedang fetch data: jadilah subscriber (follower)
                sender.subscribe()
            } else {
                // Task pertama: buat channel broadcast dan bertindak sebagai leader
                let (tx, rx) = broadcast::channel(1);
                map.insert(key.to_string(), tx);
                drop(map); // Lepas lock secepatnya sebelum I/O database

                // Eksekusi fetcher dengan timeout proteksi
                let fetch_result = match timeout(Duration::from_secs(3), fetcher()).await {
                    Ok(res) => res,
                    Err(_) => Err("Database query timeout".to_string()),
                };

                // Bersihkan registry dan kirimkan hasil ke semua subscriber
                let mut map = self.in_flight.lock().await;
                if let Some(sender) = map.remove(key) {
                    let _ = sender.send(fetch_result.clone());
                }

                return fetch_result;
            }
        };

        // Follower menunggu hasil dari leader
        match rx.recv().await {
            Ok(res) => res,
            Err(_) => Err("Leader fetcher failed prematurely".to_string()),
        }
    }
}

Integrasi pada Actix Web Shared State

Daftarkan SingleFlight ke dalam state aplikasi menggunakan web::Data. Handler menerima request, mencoba membaca cache, lalu menggunakan singleflight jika terjadi miss.

use actix_web::{get, web, App, HttpResponse, HttpServer, Responder};

struct AppState {
    flight: SingleFlight<String>,
}

// Simulasi operasi database lambat
async fn fetch_from_db(product_id: &str) -> Result<String, String> {
    tokio::time::sleep(tokio::time::Duration::from_millis(150)).await;
    Ok(format!("{{\"id\": \"{}\", \"name\": \"Item Hot\"}}", product_id))
}

#[get("/products/{id}")]
async fn get_product(
    path: web::Path<String>,
    data: web::Data<AppState>,
) -> impl Responder {
    let product_id = path.into_inner();
    let cache_key = format!("product:{}", product_id);

    // ponytail: implementasikan fast-path in-memory read (cth: moka) di sini sebelum singleflight

    let result = data
        .flight
        .execute(&cache_key, || fetch_from_db(&product_id))
        .await;

    match result {
        Ok(body) => HttpResponse::Ok().content_type("application/json").body(body),
        Err(err) if err.contains("timeout") => HttpResponse::GatewayTimeout().body(err),
        Err(err) => HttpResponse::InternalServerError().body(err),
    }
}

#[actix_web::main]
async fn main() -> std::io::Result<()> {
    let state = web::Data::new(AppState {
        flight: SingleFlight::new(),
    });

    HttpServer::new(move || {
        App::new()
            .app_data(state.clone())
            .service(get_product)
    })
    .bind(("127.0.0.1", 8080))?
    .run()
    .await
}

Runnable Test: Validasi Eksekusi Tunggal

Unit test berikut memvalidasi bahwa dari 50 pemanggilan konkuren untuk key yang sama, operasi fetch hanya dipanggil 1 kali.

#[cfg(test)]
mod tests {
    use super::*;
    use std::sync::atomic::{AtomicUsize, Ordering};

    #[tokio::test]
    async fn test_singleflight_coalesces_concurrent_calls() {
        let flight = SingleFlight::new();
        let db_hits = Arc::new(AtomicUsize::new(0));
        let mut handles = vec![];

        for _ in 0..50 {
            let flight = flight.clone();
            let db_hits = db_hits.clone();

            handles.push(tokio::spawn(async move {
                flight
                    .execute("hot_key", || async {
                        tokio::time::sleep(Duration::from_millis(50)).await;
                        db_hits.fetch_add(1, Ordering::SeqCst);
                        Ok("payload".to_string())
                    })
                    .await
            }));
        }

        for h in handles {
            let res = h.await.expect("Task panicked");
            assert_eq!(res.unwrap(), "payload");
        }

        // Pastikan DB hanya diakses 1 kali meskipun ada 50 pemanggil serentak
        assert_eq!(db_hits.load(Ordering::SeqCst), 1);
    }
}

Alternatif Praktis: Moka Cache Built-in Coalescing

Jika tidak ingin me-maintain custom singleflight, gunakan crate moka. Tipe moka::future::Cache menyediakan request deduplication secara native menggunakan metode try_get_with.

use moka::future::Cache;
use std::time::Duration;

// moka secara otomatis menggabungkan concurrent loader untuk key yang sama
let cache: Cache<String, String> = Cache::builder()
    .time_to_live(Duration::from_secs(60))
    .build();

let value = cache
    .try_get_with("hot_key".to_string(), async {
        fetch_from_db("hot_key").await
    })
    .await;
Pilih custom singleflight jika ingin request coalescing murni tanpa memori cache lokal (misal: layer di depan Redis cluster). Pilih Moka jika membutuhkan in-memory cache dengan TTL dan eviction policy otomatis.