nmbrs-rate
An async token-bucket rate limiter built on tokio::sync::Semaphore. It uses
time-scaled permits and burst recovery. nmbrs uses it to cap op throughput,
and it reports how long callers were held back so that coordinated omission
is visible in the results. You can retarget it while it runs.
Where it sits in nmbrs
- Used by
nmbrs-runtime:- Each activity (phase) with a
rategets oneRateLimiter. The limiter is registered as the target of that activity'sratedynamic control. - An op-level
rate:field adds a per-op limiter through the runtime'sratewrapper.
- Each activity (phase) with a
- Also used by the
nmbrsCLI. - Depends on
nmbrs-metricsfor theControlAppliertrait, and ontokioandfutures.
End users normally install the nmbrs CLI
and set rates in a workload or on the command line (rate=1000):
# excerpt from crates/nmbrs/examples/workloads/controls/error_rate_circuit_breaker.yaml
params:
rate: "2000"
concurrency: "8"
phases:
steady:
cycles: 3000
rate: "{rate}"
concurrency: "{concurrency}"
What it provides
RateSpecis the limiter configuration. Its public fields areops_per_sec,burst_ratio,verb: Verbandunit: TimeUnit.-
RateSpec::new(ops)uses a burst ratio of 1.1. -
RateSpec::with_burst(ops, ratio)sets the burst ratio. -
RateSpec::parse(&str)reads the comma-separated form used by params and CLI flags:1000 # 1000 ops/s, burst 1.1, verb start 1000,1.5 # 1000 ops/s, burst 1.5 1000,1.1,restart # with a verb: start | configure | restart | stopA rate that isn't positive, or an unknown verb, is an error.
-
TimeUnitis the tick unit (Nanos,Micros,Millis,Seconds).TimeUnit::for_ratechooses it from the target rate so that ticks per op fit in au32.RateLimiter:RateLimiter::start(spec)spawns a refill task on the current tokio runtime, so it must be called from inside one. The task tops up permits every 10 ms. Ticks that overflow the active pool go to a waiting pool. Burst recovery moves ticks back from the waiting pool, up toburst_ratio.acquire().awaitwaits for one op's worth of permits and returns the current backlog in ticks.wait_time_nanos()returns that backlog in nanoseconds. This is the coordinated-omission signal.total_blocks()countsacquirecalls.rate()andspec()report the current configuration.stop().awaitends the refill task. Dropping the limiter also stops it.reconfigure(spec)swaps the rate, burst ratio and unit in place, without restarting the refill task. The nextacquireuses the new cost, and the existing backlog is kept. It rejectsops_per_sec <= 0andburst_ratio < 1.0.
RateLimiterApplierimplementsnmbrs_metrics::controls::ControlApplier<RateSpec>. Register it on aControl<RateSpec>, and every successfulseton that control callsreconfigureon the limiter.
use ;
// Must run inside a tokio runtime: `start` spawns the refill task.
async
Driving the limiter through a dynamic control (requires nmbrs-metrics):
use Arc;
use ;
use ;
let limiter = new;
let control: = new.build;
control.register_applier;
control.set.await?;
assert_eq!;
Cargo features
None.
Links
- Repository: https://github.com/nosqlbench/nmbrs
- API docs: https://docs.rs/nmbrs-rate
- Design (SRD 06, rate limiter): https://github.com/nosqlbench/nmbrs/blob/main/docs/SRD/06_rate_limiter.md
- Design notes and coordinated-omission rationale: https://github.com/nosqlbench/nmbrs/blob/main/docs/SRD/notes/19_rate_limiter.md
License
Apache-2.0