Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions crates/consistent-choose-k/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,36 @@ Why replication matters
- Distributes read/write load across multiple owners, reducing hotspots.
- Enables fast recovery and higher tail-latency resilience.

## Permutation APIs and membership semantics

The existing `ConsistentPermutation` preserves **survivor list order** when
nodes are appended or removed from the end of `0..n`. The additional,
experimental `VirtualPermutation` instead preserves **replica slots**: on a
single-node append, at most one old slot changes, to the new node; on removal,
only a surviving slot that named that node changes. It does not preserve list
restriction and is not a drop-in replacement for the existing iterator or
its failover policies. Both yield distinct nodes and stable `k` prefixes.

`VirtualPermutation::new(n, seed)` supports `1..=u64::MAX`, allocates no state,
and offers both an iterator and absolute `replica_at(slot)` lookup. Expected
`O(k)` enumeration follows under ideal independent uniform permutations with
constant-cost forward/inverse evaluation, **not** as a worst-case guarantee.
The implemented noncryptographic, 64-bit seeded Feistel family approximates
that randomness model; exact uniformity and independence are not claimed.

See the [algorithm, proof assumptions and API guide](docs/virtual-permutation.md)
and the [reproducible comparison with the existing algorithm](docs/virtual-permutation-performance.md).
The existing APIs and their mappings remain unchanged.

`BalancedVirtualPermutation` is a matched-network experiment: it gives the
same cycle/slot semantics as `VirtualPermutation`, but uses **exactly** the
existing `ConsistentPermutation` Feistel, with even widths and two bits per
lift, over `1..=2^30`. Its state is also allocation-free. The
[three-way comparison](docs/virtual-permutation-performance.md#matched-network-follow-up)
repeats performance and primary/held-out statistical diagnostics. This variant
has repeatable small-domain distribution bias and is not the default or a
statistically equivalent replacement for the stronger mixer.

## Applications beyond replication

The `ConsistentChooseK` iterator produces a per-key ranking of all `n` nodes in priority order — consistently and with zero memory overhead. This ranking is a strict superset of simple replication and enables drop-in replacements for several well-known algorithms that traditionally require maintaining expensive data structures such as hash rings.
Expand Down
6 changes: 6 additions & 0 deletions crates/consistent-choose-k/benchmarks/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,12 @@ path = "performance.rs"
harness = false
test = false

[[bench]]
name = "replica_comparison"
path = "replica_comparison.rs"
harness = false
test = false

[dependencies]
consistent-choose-k = { path = "../" }

Expand Down
271 changes: 271 additions & 0 deletions crates/consistent-choose-k/benchmarks/replica_comparison.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,271 @@
//! Paired workloads only in the existing iterator's supported domain.
//! See ../docs/virtual-permutation-performance.md for methodology and results.

use std::{
hash::{DefaultHasher, Hash, Hasher},
hint::black_box,
time::Duration,
};

use consistent_choose_k::{BalancedVirtualPermutation, ConsistentPermutation, VirtualPermutation};
use criterion::{criterion_group, criterion_main, BatchSize, BenchmarkId, Criterion, Throughput};
use rand::{rngs::StdRng, RngExt, SeedableRng};

const WORKLOAD_SEED: u64 = 0x7065_726d_7574_6531;
const KEY_COUNT: usize = 128;
const NODES: &[u32] = &[
1,
3,
7,
8,
9,
15,
16,
17,
31,
32,
33,
255,
256,
257,
1000,
1023,
1024,
1025,
65535,
65536,
65537,
1_000_000,
(1 << 30) - 1,
1 << 30,
];

fn hash_key(key: u64) -> u64 {
let mut hasher = DefaultHasher::new();
key.hash(&mut hasher);
hasher.finish()
}

fn keys() -> Vec<u64> {
StdRng::seed_from_u64(WORKLOAD_SEED)
.random_iter()
.take(KEY_COUNT)
.collect()
}

fn counts(n: u32) -> Vec<usize> {
let mut counts = vec![1, 2, 3, 8, 16];
if n <= 1024 {
counts.extend([n as usize / 4, n as usize]);
}
counts.retain(|&k| k > 0 && k <= n as usize);
counts.sort_unstable();
counts.dedup();
counts
}

fn consume(iter: impl Iterator<Item = impl Into<u64>>, k: usize) {
let sum = iter
.take(k)
.fold(0u64, |sum, node| sum.wrapping_add(node.into()));
black_box(sum);
}

fn end_to_end(c: &mut Criterion) {
let keys = keys();
let seeds: Vec<_> = keys.iter().copied().map(hash_key).collect();
for mode in ["fresh", "seeded"] {
let mut group = c.benchmark_group(format!("replicas/{mode}"));
// Both execute exactly one complete query per key, including
// construction, streaming consumption and state destruction.
group.throughput(Throughput::Elements(KEY_COUNT as u64));
for &n in NODES {
for k in counts(n) {
let input = if mode == "fresh" { &keys } else { &seeds };
group.bench_function(BenchmarkId::new("layered", format!("n{n}_k{k}")), |b| {
b.iter(|| {
for &key in black_box(input) {
let seed = if mode == "fresh" { hash_key(key) } else { key };
consume(ConsistentPermutation::new(black_box(n), seed), black_box(k));
}
})
});
group.bench_function(BenchmarkId::new("virtual", format!("n{n}_k{k}")), |b| {
b.iter(|| {
for &key in black_box(input) {
let seed = if mode == "fresh" { hash_key(key) } else { key };
consume(
VirtualPermutation::new(u64::from(black_box(n)), seed),
black_box(k),
);
}
})
});
group.bench_function(BenchmarkId::new("balanced", format!("n{n}_k{k}")), |b| {
b.iter(|| {
for &key in black_box(input) {
let seed = if mode == "fresh" { hash_key(key) } else { key };
consume(
BalancedVirtualPermutation::new(black_box(n), seed),
black_box(k),
);
}
})
});
}
}
group.finish();
}
}

fn cost_components(c: &mut Criterion) {
let keys = keys();
let seeds: Vec<_> = keys.iter().copied().map(hash_key).collect();
let mut setup = c.benchmark_group("replicas/setup");
setup.throughput(Throughput::Elements(KEY_COUNT as u64));
setup.bench_function("hash_u64", |b| {
b.iter(|| {
for &key in black_box(&keys) {
black_box(hash_key(key));
}
})
});
for &n in &[17, 257, 1000, 65537, 1 << 30] {
setup.bench_function(BenchmarkId::new("layered", n), |b| {
b.iter(|| {
for &seed in black_box(&seeds) {
black_box(ConsistentPermutation::new(black_box(n), seed));
}
})
});
setup.bench_function(BenchmarkId::new("virtual", n), |b| {
b.iter(|| {
for &seed in black_box(&seeds) {
black_box(VirtualPermutation::new(u64::from(black_box(n)), seed));
}
})
});
setup.bench_function(BenchmarkId::new("balanced", n), |b| {
b.iter(|| {
for &seed in black_box(&seeds) {
black_box(BalancedVirtualPermutation::new(black_box(n), seed));
}
})
});
}
setup.finish();

for mode in ["stream_only", "collect", "slot"] {
let mut group = c.benchmark_group(format!("replicas/{mode}"));
group.throughput(Throughput::Elements(KEY_COUNT as u64));
for &n in &[17, 257, 1000, 65537, 1 << 30] {
for k in counts(n) {
group.bench_function(BenchmarkId::new("layered", format!("n{n}_k{k}")), |b| {
match mode {
"stream_only" => b.iter_batched_ref(
|| {
seeds
.iter()
.map(|&seed| ConsistentPermutation::new(n, seed))
.collect::<Vec<_>>()
},
|iterators| {
for iter in black_box(iterators) {
consume(iter, black_box(k));
}
},
BatchSize::SmallInput,
),
_ => b.iter(|| {
for &seed in black_box(&seeds) {
let mut iter = ConsistentPermutation::new(black_box(n), seed);
if mode == "collect" {
// Same output width and allocation policy for both.
let mut out = Vec::with_capacity(black_box(k));
out.extend(iter.take(k).map(u64::from));
black_box(out);
} else {
// No random-access API in the baseline: nth must replay.
black_box(iter.nth(black_box(k - 1)));
}
}
}),
}
});
group.bench_function(BenchmarkId::new("virtual", format!("n{n}_k{k}")), |b| {
match mode {
"stream_only" => b.iter_batched_ref(
|| {
seeds
.iter()
.map(|&seed| VirtualPermutation::new(u64::from(n), seed))
.collect::<Vec<_>>()
},
|iterators| {
for iter in black_box(iterators) {
consume(iter, black_box(k));
}
},
BatchSize::SmallInput,
),
_ => b.iter(|| {
for &seed in black_box(&seeds) {
let iter = VirtualPermutation::new(u64::from(black_box(n)), seed);
if mode == "collect" {
let mut out = Vec::with_capacity(black_box(k));
out.extend(iter.take(k));
black_box(out);
} else {
black_box(iter.replica_at(black_box(k as u64 - 1)));
}
}
}),
}
});
group.bench_function(BenchmarkId::new("balanced", format!("n{n}_k{k}")), |b| {
match mode {
"stream_only" => b.iter_batched_ref(
|| {
seeds
.iter()
.map(|&seed| BalancedVirtualPermutation::new(n, seed))
.collect::<Vec<_>>()
},
|iterators| {
for iter in black_box(iterators) {
consume(iter, black_box(k));
}
},
BatchSize::SmallInput,
),
_ => b.iter(|| {
for &seed in black_box(&seeds) {
let iter = BalancedVirtualPermutation::new(black_box(n), seed);
if mode == "collect" {
let mut out = Vec::with_capacity(black_box(k));
out.extend(iter.take(k));
black_box(out);
} else {
black_box(iter.replica_at(black_box(k as u64 - 1)));
}
}
}),
}
});
}
}
group.finish();
}
}

criterion_group! {
name = benches;
config = Criterion::default()
.sample_size(20)
.warm_up_time(Duration::from_millis(100))
.measurement_time(Duration::from_millis(300))
.nresamples(1000)
.without_plots();
targets = end_to_end, cost_components
}
criterion_main!(benches);
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
#!/usr/bin/env python3
"""Convert Criterion's batch estimates to ns/query and ns/replica CSV.

Usage: python3 crates/consistent-choose-k/benchmarks/summarize_replica_comparison.py \
target/criterion > comparison.csv
"""

import csv
import json
from pathlib import Path
import re
import sys


def summarize(root):
rows = []
for metadata in root.glob("**/new/benchmark.json"):
benchmark = json.loads(metadata.read_text())
group = benchmark["group_id"]
if not group.startswith("replicas/"):
continue
mode = group.removeprefix("replicas/")
algorithm = benchmark["function_id"]
value = benchmark.get("value_str") or ""
match = re.fullmatch(r"n(\d+)_k(\d+)", value)
n, k = map(int, match.groups()) if match else (int(value or 0), 0)
estimates = json.loads((metadata.parent / "estimates.json").read_text())
mean = estimates["mean"]
interval = mean["confidence_interval"]
divisor = benchmark["throughput"]["Elements"]
rows.append(
(
mode,
algorithm,
n,
k,
mean["point_estimate"] / divisor,
interval["lower_bound"] / divisor,
interval["upper_bound"] / divisor,
estimates["std_dev"]["point_estimate"] / divisor,
mean["point_estimate"] / divisor / (k if k and mode != "slot" else 1),
)
)
if not rows:
raise SystemExit(f"No replica_comparison results found in {root}")
writer = csv.writer(sys.stdout)
writer.writerow(
["mode", "algorithm", "n", "k", "ns_query", "ci95_low", "ci95_high",
"stddev_ns", "ns_replica"]
)
for row in sorted(rows):
writer.writerow([*row[:4], *(f"{value:.3f}" for value in row[4:])])


if __name__ == "__main__":
if len(sys.argv) != 2:
raise SystemExit(__doc__)
summarize(Path(sys.argv[1]))
Loading
Loading