use std::io::{self, Write};
use subms::{SubMsPerfHarness, SubMsStageKind, SubMsTimer, summarize, summary_to_json};
use subms_merge_iterator::{
DedupEntry, DedupMergeIterator, MergeIterator, PriorityEntry, PriorityMergeIterator,
PrioritySource, SeekableMergeIterator, TombstoneEntry, TombstoneMergeIterator,
};
const ENTRIES: usize = 50_000;
const STREAMS: usize = 16;
const SEED: u64 = 0;
const PER_STREAM: usize = ENTRIES / STREAMS;
const TOTAL: usize = PER_STREAM * STREAMS;
fn plain_streams() -> Vec<std::vec::IntoIter<u64>> {
(0..STREAMS)
.map(|s| {
(0..PER_STREAM)
.map(move |i| (s + i * STREAMS) as u64)
.collect::<Vec<u64>>()
.into_iter()
})
.collect()
}
fn main() -> io::Result<()> {
let mut h = SubMsPerfHarness::new("merge-iterator-features", "rust");
h.input("entries", &ENTRIES.to_string());
h.input("streams", &STREAMS.to_string());
h.input("seed", &SEED.to_string());
h.add_meta("per_stream", &PER_STREAM.to_string());
h.add_meta("subms.recipe.slug", "subms-merge-iterator");
h.add_meta("subms.recipe.category", "storage");
{
h.add_meta("subms.workload.feature", "base");
let mut iter = MergeIterator::new(plain_streams());
let s = h
.stage("base_next", TOTAL)
.with_kind(SubMsStageKind::HotPath);
for _ in 0..TOTAL {
let t0 = SubMsTimer::tick();
let _ = iter.next();
s.record(t0.elapsed_ns());
}
}
{
h.add_meta("subms.workload.feature", "seek-to");
const SEEKS: usize = 4_000;
let step = (TOTAL / SEEKS).max(1) as u64;
let targets: Vec<u64> = (0..SEEKS).map(|i| i as u64 * step).collect();
let mut iter = SeekableMergeIterator::new(plain_streams());
let s_seek = h.stage("seek", SEEKS).with_kind(SubMsStageKind::HotPath);
let mut seek_times: Vec<u64> = Vec::with_capacity(SEEKS);
let mut after_times: Vec<u64> = Vec::with_capacity(SEEKS);
for target in &targets {
let t0 = SubMsTimer::tick();
iter.seek(target);
seek_times.push(t0.elapsed_ns());
let t1 = SubMsTimer::tick();
let _ = iter.next();
after_times.push(t1.elapsed_ns());
}
for ns in seek_times {
s_seek.record(ns);
}
let s_after = h
.stage("next_after_seek", SEEKS)
.with_kind(SubMsStageKind::HotPath);
for ns in after_times {
s_after.record(ns);
}
}
{
h.add_meta("subms.workload.feature", "tombstones");
let streams: Vec<std::vec::IntoIter<TombstoneEntry<u64, u64>>> = (0..STREAMS)
.map(|s| {
(0..PER_STREAM)
.map(move |i| {
let key = (s + i * STREAMS) as u64;
if key % 8 == 0 {
TombstoneEntry::tombstone(key)
} else {
TombstoneEntry::live(key, key)
}
})
.collect::<Vec<_>>()
.into_iter()
})
.collect();
let mut iter = TombstoneMergeIterator::new(streams);
let s = h
.stage("tombstones_next", TOTAL)
.with_kind(SubMsStageKind::HotPath);
for _ in 0..TOTAL {
let t0 = SubMsTimer::tick();
let yielded = iter.next();
s.record(t0.elapsed_ns());
if yielded.is_none() {
break;
}
}
}
{
h.add_meta("subms.workload.feature", "dedup");
let streams: Vec<std::vec::IntoIter<DedupEntry<u64, u64>>> = (0..STREAMS)
.map(|s| {
(0..PER_STREAM)
.map(move |i| {
let key = ((s + i * STREAMS) as u64) / 2;
DedupEntry::new(key, key)
})
.collect::<Vec<_>>()
.into_iter()
})
.collect();
let mut iter = DedupMergeIterator::new(streams);
let s = h
.stage("dedup_next", TOTAL)
.with_kind(SubMsStageKind::HotPath);
for _ in 0..TOTAL {
let t0 = SubMsTimer::tick();
let yielded = iter.next();
s.record(t0.elapsed_ns());
if yielded.is_none() {
break;
}
}
}
{
h.add_meta("subms.workload.feature", "priority");
let sources: Vec<PrioritySource<std::vec::IntoIter<PriorityEntry<u64, u64>>>> = (0
..STREAMS)
.map(|s| {
let stream = (0..PER_STREAM)
.map(move |i| {
let key = ((s + i * STREAMS) as u64) / 2;
PriorityEntry::new(key, key)
})
.collect::<Vec<_>>()
.into_iter();
PrioritySource::new((STREAMS - s) as i32, stream)
})
.collect();
let mut iter = PriorityMergeIterator::new(sources);
let s = h
.stage("priority_next", TOTAL)
.with_kind(SubMsStageKind::HotPath);
for _ in 0..TOTAL {
let t0 = SubMsTimer::tick();
let yielded = iter.next();
s.record(t0.elapsed_ns());
if yielded.is_none() {
break;
}
}
}
let summary = summarize(&h);
let mut stdout = io::stdout();
summary_to_json(&summary, &mut stdout)?;
writeln!(stdout)?;
Ok(())
}