Skip to main content

subms_merge_iterator/features/
dedup.rs

1//! Deduplicating k-way merge. Each input is a stream of (key, value)
2//! entries sorted by key. On key tie across sources, the highest-
3//! indexed source wins ("latest"). Output is one entry per distinct
4//! key.
5//!
6//! Same shape as the tombstone-aware merge, minus the tombstone flag -
7//! every entry is "live". Useful for log compaction over append-only
8//! shards where you want the freshest value per key.
9
10use std::cmp::Reverse;
11use std::collections::BinaryHeap;
12
13#[derive(Clone, Debug, PartialEq, Eq)]
14pub struct DedupEntry<K, V> {
15    pub key: K,
16    pub value: V,
17}
18
19impl<K, V> DedupEntry<K, V> {
20    pub fn new(key: K, value: V) -> Self {
21        Self { key, value }
22    }
23}
24
25pub struct DedupMergeIterator<K, V, I>
26where
27    K: Ord,
28    I: Iterator<Item = DedupEntry<K, V>>,
29{
30    streams: Vec<I>,
31    heap: BinaryHeap<Reverse<HeapItem<K, V>>>,
32}
33
34struct HeapItem<K, V> {
35    key: K,
36    source: usize,
37    value: V,
38}
39
40impl<K: Ord, V> PartialEq for HeapItem<K, V> {
41    fn eq(&self, other: &Self) -> bool {
42        self.key == other.key && self.source == other.source
43    }
44}
45impl<K: Ord, V> Eq for HeapItem<K, V> {}
46
47impl<K: Ord, V> PartialOrd for HeapItem<K, V> {
48    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
49        Some(self.cmp(other))
50    }
51}
52impl<K: Ord, V> Ord for HeapItem<K, V> {
53    fn cmp(&self, other: &Self) -> std::cmp::Ordering {
54        self.key
55            .cmp(&other.key)
56            .then(other.source.cmp(&self.source))
57    }
58}
59
60impl<K, V, I> DedupMergeIterator<K, V, I>
61where
62    K: Ord,
63    I: Iterator<Item = DedupEntry<K, V>>,
64{
65    pub fn new<S: IntoIterator<Item = I>>(streams: S) -> Self {
66        let mut streams: Vec<I> = streams.into_iter().collect();
67        let mut heap = BinaryHeap::with_capacity(streams.len());
68        for (i, s) in streams.iter_mut().enumerate() {
69            if let Some(e) = s.next() {
70                heap.push(Reverse(HeapItem {
71                    key: e.key,
72                    source: i,
73                    value: e.value,
74                }));
75            }
76        }
77        Self { streams, heap }
78    }
79
80    fn advance(&mut self, source: usize) {
81        if let Some(e) = self.streams[source].next() {
82            self.heap.push(Reverse(HeapItem {
83                key: e.key,
84                source,
85                value: e.value,
86            }));
87        }
88    }
89}
90
91impl<K, V, I> Iterator for DedupMergeIterator<K, V, I>
92where
93    K: Ord,
94    I: Iterator<Item = DedupEntry<K, V>>,
95{
96    type Item = DedupEntry<K, V>;
97
98    fn next(&mut self) -> Option<DedupEntry<K, V>> {
99        let Reverse(HeapItem {
100            key: winning_key,
101            source,
102            value: winning_value,
103        }) = self.heap.pop()?;
104        self.advance(source);
105
106        // Drop every other entry with the same key.
107        while let Some(Reverse(item)) = self.heap.peek() {
108            if item.key == winning_key {
109                let Reverse(item) = self.heap.pop().unwrap();
110                self.advance(item.source);
111            } else {
112                break;
113            }
114        }
115        Some(DedupEntry {
116            key: winning_key,
117            value: winning_value,
118        })
119    }
120}
121
122#[cfg(test)]
123#[path = "dedup_tests.rs"]
124mod tests;