subms-merge-iterator 0.10.0

submillisecond.com cookbook recipe - storage: subms-merge-iterator. N-way merge of sorted streams via min-heap.
Documentation
//! K-way merge iterator. Min-heap of stream heads; pop the minimum, advance
//! that stream, push its next value if any.
//!
//! Streams must be sorted ascending. Output is the global sorted union.
//!
//! ```
//! use subms_merge_iterator::MergeIterator;
//! let streams: Vec<Box<dyn Iterator<Item = i32>>> = vec![
//!     Box::new(vec![1, 4, 7].into_iter()),
//!     Box::new(vec![2, 5, 8].into_iter()),
//!     Box::new(vec![3, 6, 9].into_iter()),
//! ];
//! let mut merge = MergeIterator::new(streams);
//! assert_eq!(merge.peek(), Some(&1));
//! assert_eq!(merge.live_streams(), 3);
//! let merged: Vec<_> = merge.collect();
//! assert_eq!(merged, (1..=9).collect::<Vec<_>>());
//! ```
//!
//! The iterator is a single-threaded cursor. It owns its sources, so it is
//! `Send` when they are, and there is no interior mutability to share across
//! threads: one merge per consumer.
//!
//! Full writeup, design notes and measured benchmarks:
//! <https://www.submillisecond.com/cookbook/recipes/subms-merge-iterator>

use std::cmp::Reverse;
use std::collections::BinaryHeap;

pub struct MergeIterator<T: Ord, I: Iterator<Item = T>> {
    streams: Vec<I>,
    /// Min-heap on (value, stream index). `Reverse` flips to min ordering.
    heap: BinaryHeap<Reverse<(T, usize)>>,
}

impl<T: Ord, I: Iterator<Item = T>> MergeIterator<T, I> {
    pub fn new<S: IntoIterator<Item = I>>(streams: S) -> Self {
        let mut streams: Vec<I> = streams.into_iter().collect();
        let mut heap = BinaryHeap::with_capacity(streams.len());
        for (i, s) in streams.iter_mut().enumerate() {
            if let Some(v) = s.next() {
                heap.push(Reverse((v, i)));
            }
        }
        Self { streams, heap }
    }

    /// The value the next `next()` will yield, without consuming it.
    pub fn peek(&self) -> Option<&T> {
        self.heap.peek().map(|Reverse((value, _))| value)
    }

    /// Streams still holding a head in the heap. Drops to zero exactly when the
    /// merge is exhausted, so this doubles as the RocksDB-style `valid()` check.
    pub fn live_streams(&self) -> usize {
        self.heap.len()
    }

    /// Streams the merge was constructed over, live or not.
    pub fn num_streams(&self) -> usize {
        self.streams.len()
    }
}

impl<T: Ord, I: Iterator<Item = T>> Iterator for MergeIterator<T, I> {
    type Item = T;
    fn next(&mut self) -> Option<T> {
        let Reverse((value, idx)) = self.heap.pop()?;
        if let Some(next_value) = self.streams[idx].next() {
            self.heap.push(Reverse((next_value, idx)));
        }
        Some(value)
    }
}

#[cfg(feature = "harness")]
pub mod recipe;

// Opt-in feature modules. Each is independent of the base merge iterator
// and gated by its own Cargo feature; the bare `cargo add
// subms-merge-iterator` shape stays zero-dep + std-only.
#[cfg(any(
    feature = "seek-to",
    feature = "reverse",
    feature = "tombstones",
    feature = "dedup",
    feature = "priority"
))]
pub mod features;

#[cfg(feature = "dedup")]
pub use features::dedup::{DedupEntry, DedupMergeIterator};
#[cfg(feature = "priority")]
pub use features::priority::{PriorityEntry, PriorityMergeIterator, PrioritySource};
#[cfg(feature = "reverse")]
pub use features::reverse::ReverseMergeIterator;
#[cfg(feature = "seek-to")]
pub use features::seek::SeekableMergeIterator;
#[cfg(feature = "tombstones")]
pub use features::tombstones::{TombstoneEntry, TombstoneMergeIterator};

#[cfg(test)]
#[path = "merge_tests.rs"]
mod merge_tests;

#[cfg(test)]
#[path = "sample_app_tests.rs"]
mod sample_app_tests;