Aplikasi pemrosesan instruksi paralel yang sering berganti target konteks—seperti switch-instruct execution engine pada desktop dev-tools—sering mengalami degradasi throughput ekstrem saat beban kerja meningkat. Penyebab utamanya bukan kekurangan core CPU, melainkan lock thrashing dan cache line bouncing akibat perebutan akses ke state context bersama.
Anatomi Masalah: Naive RwLock dan Cache Churn
Pola implementasi awal yang umum adalah menempatkan registry context di dalam Arc<RwLock<HashMap<ContextId, ContextData>>>. Setiap kali worker menerima task:
- Worker mengambil read lock untuk membaca state context aktif.
- Jika task memerlukan mutasi atau context reload, worker meminta write lock.
- Permintaan write lock memblokir seluruh worker lain yang membaca context yang tidak terkait sekalipun.
Ketika ratusan task berganti konteks secara agresif, siklus komputasi habis di kernel synchronization primitives. Lebih buruk lagi, L1/L2/L3 CPU cache line terus-menerus diinvalir di antara core processor (cache line bouncing), mengakibatkan latency spike dan penurunan drastis pada IPC (Instructions Per Cycle).
Arsitektur Solusi: Affinity Partitioning & Epoch Invalidation
Pendekatan ini memisahkan konkurensi menjadi dua mekanisme:
- Partitioned Task Channel: Dispatcher merutekan task ke worker channel tertentu berdasarkan konsistensi hash dari
context_id. Konteks yang sama selalu dieksekusi oleh worker yang sama jika memungkinkan, memaksimalkan cache locality. - Atomic Epoch Counter: State update tidak memblokir pembacaan. Setiap modifikasi konteks menaikkan nilai
AtomicU64global khusus context tersebut. Worker membandingkan versi epoch lokal miliknya dengan epoch global secara atomic tanpa lock contention.
Alternatif teringkas: Gunakan model multi-process OS terisolasi per konteks jika isolasi memori mutlak diperlukan tanpa runtime Rust kustom.
Implementasi Rust Minimal
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{mpsc, Arc};
use std::thread;
#[derive(Clone, Debug, PartialEq)]
pub struct InstructionTask {
pub context_id: usize,
pub payload: String,
}
struct LocalContext {
epoch: u64,
data: String,
}
pub struct WorkerPool {
dispatchers: Vec<mpsc::Sender<InstructionTask>>,
epochs: Arc<HashMap<usize, AtomicU64>>,
workers: Vec<thread::JoinHandle<()>>,
}
impl WorkerPool {
pub fn new(num_workers: usize, context_ids: Vec<usize>) -> Self {
let mut dispatchers = Vec::with_capacity(num_workers);
let mut workers = Vec::with_capacity(num_workers);
let mut epoch_map = HashMap::new();
for id in context_ids {
epoch_map.insert(id, AtomicU64::new(1));
}
let epochs = Arc::new(epoch_map);
for _ in 0..num_workers {
let (tx, rx) = mpsc::channel::<InstructionTask>();
dispatchers.push(tx);
let thread_epochs = Arc::clone(&epochs);
let handle = thread::spawn(move || {
// Cache lokal per worker; zero lock contention saat pembacaan
// ponytail: memory context unbounded. Upgrade ke LRU eviction jika ID bertambah dinamis.
let mut local_cache: HashMap<usize, LocalContext> = HashMap::new();
while let Ok(task) = rx.recv() {
let current_global_epoch = thread_epochs
.get(&task.context_id)
.map(|e| e.load(Ordering::Acquire))
.unwrap_or(0);
let ctx = local_cache.entry(task.context_id).or_insert_with(|| {
LocalContext {
epoch: 0,
data: format!("Init State Context {}", task.context_id),
}
});
if ctx.epoch < current_global_epoch {
// Invalidation terjadi: reload state lokal worker
ctx.data = format!("Reloaded Context {} @ Epoch {}", task.context_id, current_global_epoch);
ctx.epoch = current_global_epoch;
}
// Eksekusi task dengan context lokal terisolasi
assert!(ctx.data.contains(&format!("Context {}", task.context_id)));
}
});
workers.push(handle);
}
WorkerPool { dispatchers, epochs, workers }
}
pub fn dispatch(&self, task: InstructionTask) {
// ponytail: modulo routing rentan skew jika beban context tidak rata. Upgrade ke dynamic work stealing.
let target_worker = task.context_id % self.dispatchers.len();
let _ = self.dispatchers[target_worker].send(task);
}
pub fn invalidate_context(&self, context_id: usize) {
if let Some(epoch) = self.epochs.get(&context_id) {
epoch.fetch_add(1, Ordering::Release);
}
}
}
fn main() {
// Assertion test konkurensi
let pool = WorkerPool::new(4, vec![101, 102, 103, 104]);
// Dispatch awal
for _ in 0..50 {
pool.dispatch(InstructionTask {
context_id: 101,
payload: "EXECUTE_INSTRUCT_A".into(),
});
}
// Simulasi invalidasi state context secara concurrently
pool.invalidate_context(101);
for _ in 0..50 {
pool.dispatch(InstructionTask {
context_id: 101,
payload: "EXECUTE_INSTRUCT_B".into(),
});
}
let epoch_val = pool.epochs.get(&101).unwrap().load(Ordering::Acquire);
assert_eq!(epoch_val, 2, "Epoch harus ter-increment tepat satu kali");
}
[code] → skipped: dynamic channel resizing, unbounded LRU context eviction, add when context count exceeds worker memory limits.
Analisis Trade-Off: Memori vs. Contention
| Metrik | Naive Shared RwLock | Partitioned Channel + Epoch |
|---|---|---|
| Contention Overhead | Tinggi (lock contention pada setiap switch context) | Hampir Nol (hanya atomic load pada hot path) |
| Konsumsi Memori | Rendah (hanya ada 1 salinan state di heap) | Moderat (potensi replikasi context lokal pada beberapa worker) |
| Cache Locality | Buruk (cache lines terus berganti di core CPU) | Optimal (context tetap di L1/L2 core worker terkait) |
| Latensi Mutasi | Tinggi (menunggu semua read lock selesai) | Rendah (hanya 1 instruksi atomic increment) |
Metrik Operasional yang Perlu Dipantau
Saat menerapkan pola ini di production, ukur efisiensi sistem melalui telemetri berikut:
- Worker Channel Queue Depth: Deteksi ketimpangan distribusi beban jika satu
context_idmembebani satu worker secara berlebihan. - Epoch Invalidation Frequency: Mutasi context yang terlalu sering akan mengurangi efektivitas cache lokal worker.
- Memory Footprint Per Worker: Pastikan replikasi context lokal tidak memicu Out-Of-Memory (OOM) pada mesin dengan RAM terbatas.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!