use rustis::{
Result,
cache::Cache,
client::Client,
commands::{
ClientTrackingOptions, ClientTrackingStatus, ConnectionCommands, FlushingMode,
ServerCommands, StringCommands,
},
resp::cmd,
};
use std::time::Instant;
async fn server_get_calls(control: &Client) -> Result<u64> {
let info: String = control.send(cmd("INFO").arg("commandstats"), None).await?;
for line in info.lines() {
if let Some(rest) = line.strip_prefix("cmdstat_get:")
&& let Some(calls) = rest.split(',').find_map(|kv| kv.strip_prefix("calls="))
{
return Ok(calls.trim().parse().unwrap_or(0));
}
}
Ok(0)
}
#[tokio::main]
async fn main() -> Result<()> {
let control = Client::connect("redis://127.0.0.1?connection_name=stampede_control").await?;
let cache_client = Client::connect("redis://127.0.0.1?connection_name=stampede_cache").await?;
control
.client_tracking(ClientTrackingStatus::Off, ClientTrackingOptions::default())
.await?;
println!(
"{:>5} | {:>12} | {:>15} | {:>13} | {:>12}",
"N", "server GETs", "stampede factor", "cold burst", "warm burst"
);
println!("{}", "-".repeat(70));
for &concurrency in &[1usize, 2, 4, 8, 16, 32, 64, 128] {
control.flushall(FlushingMode::Sync).await?;
let _: () = control.set("hot", "value").await?;
let cache = Cache::new(cache_client.clone(), 60, ClientTrackingOptions::default()).await?;
control
.send::<()>(cmd("CONFIG").arg("RESETSTAT"), None)
.await?;
let before = server_get_calls(&control).await?;
let start = Instant::now();
let mut handles = Vec::with_capacity(concurrency);
for _ in 0..concurrency {
let cache = cache.clone();
handles.push(tokio::spawn(
async move { cache.get::<String>("hot").await },
));
}
for handle in handles {
let value = handle.await.expect("task panicked")?;
assert_eq!(value, "value");
}
let cold = start.elapsed();
let after = server_get_calls(&control).await?;
let server_gets = after.saturating_sub(before);
let start = Instant::now();
let mut handles = Vec::with_capacity(concurrency);
for _ in 0..concurrency {
let cache = cache.clone();
handles.push(tokio::spawn(
async move { cache.get::<String>("hot").await },
));
}
for handle in handles {
handle.await.expect("task panicked")?;
}
let warm = start.elapsed();
println!(
"{concurrency:>5} | {server_gets:>12} | {:>14.1}x | {:>10.2?} | {:>10.2?}",
server_gets as f64 / concurrency as f64,
cold,
warm,
);
}
println!(
"\nReading: a stampede factor near 1.0x means the burst is already coalesced;\n\
near N (server GETs ~= N) means every concurrent miss hits the server. The\n\
cold-vs-warm gap is the round-trip latency a single-flight coalescer could\n\
remove for the redundant callers."
);
Ok(())
}