quantile-sketch 0.1.0

A fast, concurrent DDSketch for relative-error quantiles.
Documentation
use loom::{sync::Arc, thread};

use super::{assert_estimate_eq, ConcurrentDDSketch};

const ERRORS: [f64; 4] = [0.005, 0.01, 0.02, 0.05];

fn check_model(max_branches: usize, model: impl Fn() + Send + Sync + 'static) {
    let mut builder = loom::model::Builder::new();
    builder.max_branches = max_branches;
    builder.check(model);
}

fn check_models(max_branches: usize, model: impl Fn(f64) + Copy + Send + Sync + 'static) {
    for error in ERRORS {
        check_model(max_branches, move || model(error));
    }
}

#[test]
fn quantiles_are_independent_of_concurrent_insert_order() {
    const QUANTILES: [f64; 5] = [0.0, 0.25, 0.5, 0.75, 1.0];

    check_models(160_000, |error| {
        let serial = ConcurrentDDSketch::with_err_and_range(error, 1.0, 4.0);
        for value in [1.0, 2.0, 3.0, 2.0, 3.0, 4.0] {
            serial.insert(value);
        }
        let expected = QUANTILES.map(|q| serial.quantile(q));
        drop(serial);

        let concurrent = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 4.0));
        let a = Arc::clone(&concurrent);
        let b = Arc::clone(&concurrent);

        let first = thread::spawn(move || {
            a.insert(1.0);
            a.insert(2.0);
            a.insert(3.0);
        });
        let second = thread::spawn(move || {
            b.insert(2.0);
            b.insert(3.0);
            b.insert(4.0);
        });
        first.join().unwrap();
        second.join().unwrap();

        for (actual, expected) in QUANTILES.map(|q| concurrent.quantile(q)).into_iter().zip(expected) {
            assert_estimate_eq(actual, expected);
        }
    });
}

#[test]
fn merges_can_initialize_the_same_block_concurrently() {
    check_models(160_000, |error| {
        let serial = ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0);
        for value in [1.0, 1.2, 1.8, 2.0] {
            serial.insert(value);
        }
        let expected = [serial.quantile(0.0), serial.quantile(1.0)];
        drop(serial);

        let merged = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0));
        let low = ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0);
        low.insert(1.0);
        low.insert(1.2);
        let high = ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0);
        high.insert(1.8);
        high.insert(2.0);

        let a = Arc::clone(&merged);
        let b = Arc::clone(&merged);
        let first = thread::spawn(move || a.merge(&low).unwrap());
        let second = thread::spawn(move || b.merge(&high).unwrap());
        first.join().unwrap();
        second.join().unwrap();

        for (actual, expected) in [merged.quantile(0.0), merged.quantile(1.0)].into_iter().zip(expected) {
            assert_estimate_eq(actual, expected);
        }
        assert_eq!(merged.percentile_rank(1.2), Some(0.5));
    });
}

#[test]
fn merge_may_include_a_concurrent_source_insert() {
    check_models(80_000, |error| {
        let serial = ConcurrentDDSketch::with_err_and_range(error, 1.0, 1.0);
        serial.insert(1.0);
        let expected = serial.quantile(0.5);
        drop(serial);

        let source = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 1.0));
        let merged = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 1.0));
        let insert_source = Arc::clone(&source);
        let merge_source = Arc::clone(&source);
        let merge_target = Arc::clone(&merged);

        let insert = thread::spawn(move || insert_source.insert(1.0));
        let merge = thread::spawn(move || merge_target.merge(&merge_source).unwrap());
        insert.join().unwrap();
        merge.join().unwrap();

        assert_estimate_eq(source.quantile(0.5), expected);
        let result = merged.quantile(0.5);
        if result.is_some() {
            assert_estimate_eq(result, expected);
        }
    });
}

#[test]
fn queries_tolerate_a_concurrent_insert() {
    check_models(80_000, |error| {
        let sketch = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0));
        sketch.insert(1.0);

        let inserter = Arc::clone(&sketch);
        let insert = thread::spawn(move || inserter.insert(2.0));

        let quantile = sketch.quantile(1.0).unwrap();
        assert!((1.0..=2.0).contains(&quantile));
        let rank = sketch.percentile_rank(1.0).unwrap();
        assert!(rank == 0.5 || rank == 1.0);

        insert.join().unwrap();
        let maximum = sketch.quantile(1.0).unwrap();
        assert!((maximum - 2.0).abs() / 2.0 <= error * 1.001);
        assert_eq!(sketch.percentile_rank(1.0), Some(0.5));
    });
}

#[test]
fn bulk_quantiles_tolerate_a_concurrent_insert() {
    const QUANTILES: [f64; 4] = [0.0, 0.5, 0.5, 1.0];

    check_models(80_000, |error| {
        let sketch = Arc::new(ConcurrentDDSketch::with_err_and_range(error, 1.0, 2.0));
        sketch.insert(1.0);

        let inserter = Arc::clone(&sketch);
        let insert = thread::spawn(move || inserter.insert(2.0));

        let mut estimates = [None; 4];
        sketch.quantiles(&QUANTILES, &mut estimates);
        let estimates = estimates.map(Option::unwrap);
        assert!(estimates.windows(2).all(|pair| pair[0] <= pair[1]));
        assert!(estimates.iter().all(|value| (1.0..=2.0).contains(value)));

        insert.join().unwrap();
    });
}