Skip to main content

subms_merge_iterator/
lib.rs

1//! K-way merge iterator. Min-heap of stream heads; pop the minimum, advance
2//! that stream, push its next value if any.
3//!
4//! Streams must be sorted ascending. Output is the global sorted union.
5//!
6//! ```
7//! use subms_merge_iterator::MergeIterator;
8//! let streams: Vec<Box<dyn Iterator<Item = i32>>> = vec![
9//!     Box::new(vec![1, 4, 7].into_iter()),
10//!     Box::new(vec![2, 5, 8].into_iter()),
11//!     Box::new(vec![3, 6, 9].into_iter()),
12//! ];
13//! let mut merge = MergeIterator::new(streams);
14//! assert_eq!(merge.peek(), Some(&1));
15//! assert_eq!(merge.live_streams(), 3);
16//! let merged: Vec<_> = merge.collect();
17//! assert_eq!(merged, (1..=9).collect::<Vec<_>>());
18//! ```
19//!
20//! The iterator is a single-threaded cursor. It owns its sources, so it is
21//! `Send` when they are, and there is no interior mutability to share across
22//! threads: one merge per consumer.
23//!
24//! Full writeup, design notes and measured benchmarks:
25//! <https://www.submillisecond.com/cookbook/recipes/subms-merge-iterator>
26
27use std::cmp::Reverse;
28use std::collections::BinaryHeap;
29
30pub struct MergeIterator<T: Ord, I: Iterator<Item = T>> {
31    streams: Vec<I>,
32    /// Min-heap on (value, stream index). `Reverse` flips to min ordering.
33    heap: BinaryHeap<Reverse<(T, usize)>>,
34}
35
36impl<T: Ord, I: Iterator<Item = T>> MergeIterator<T, I> {
37    pub fn new<S: IntoIterator<Item = I>>(streams: S) -> Self {
38        let mut streams: Vec<I> = streams.into_iter().collect();
39        let mut heap = BinaryHeap::with_capacity(streams.len());
40        for (i, s) in streams.iter_mut().enumerate() {
41            if let Some(v) = s.next() {
42                heap.push(Reverse((v, i)));
43            }
44        }
45        Self { streams, heap }
46    }
47
48    /// The value the next `next()` will yield, without consuming it.
49    pub fn peek(&self) -> Option<&T> {
50        self.heap.peek().map(|Reverse((value, _))| value)
51    }
52
53    /// Streams still holding a head in the heap. Drops to zero exactly when the
54    /// merge is exhausted, so this doubles as the RocksDB-style `valid()` check.
55    pub fn live_streams(&self) -> usize {
56        self.heap.len()
57    }
58
59    /// Streams the merge was constructed over, live or not.
60    pub fn num_streams(&self) -> usize {
61        self.streams.len()
62    }
63}
64
65impl<T: Ord, I: Iterator<Item = T>> Iterator for MergeIterator<T, I> {
66    type Item = T;
67    fn next(&mut self) -> Option<T> {
68        let Reverse((value, idx)) = self.heap.pop()?;
69        if let Some(next_value) = self.streams[idx].next() {
70            self.heap.push(Reverse((next_value, idx)));
71        }
72        Some(value)
73    }
74}
75
76#[cfg(feature = "harness")]
77pub mod recipe;
78
79// Opt-in feature modules. Each is independent of the base merge iterator
80// and gated by its own Cargo feature; the bare `cargo add
81// subms-merge-iterator` shape stays zero-dep + std-only.
82#[cfg(any(
83    feature = "seek-to",
84    feature = "reverse",
85    feature = "tombstones",
86    feature = "dedup",
87    feature = "priority"
88))]
89pub mod features;
90
91#[cfg(feature = "dedup")]
92pub use features::dedup::{DedupEntry, DedupMergeIterator};
93#[cfg(feature = "priority")]
94pub use features::priority::{PriorityEntry, PriorityMergeIterator, PrioritySource};
95#[cfg(feature = "reverse")]
96pub use features::reverse::ReverseMergeIterator;
97#[cfg(feature = "seek-to")]
98pub use features::seek::SeekableMergeIterator;
99#[cfg(feature = "tombstones")]
100pub use features::tombstones::{TombstoneEntry, TombstoneMergeIterator};
101
102#[cfg(test)]
103#[path = "merge_tests.rs"]
104mod merge_tests;
105
106#[cfg(test)]
107#[path = "sample_app_tests.rs"]
108mod sample_app_tests;