mito2/memtable/
version.rs1use std::sync::Arc;
18use std::time::Duration;
19
20use common_time::Timestamp;
21use smallvec::SmallVec;
22use store_api::metadata::RegionMetadataRef;
23use store_api::storage::SequenceNumber;
24
25use crate::error::Result;
26use crate::memtable::time_partition::TimePartitionsRef;
27use crate::memtable::{MemtableId, MemtableRef};
28
29pub(crate) type SmallMemtableVec = SmallVec<[MemtableRef; 2]>;
30
31#[derive(Debug, Clone)]
33pub(crate) struct MemtableVersion {
34 pub(crate) mutable: TimePartitionsRef,
36 immutables: SmallMemtableVec,
42}
43
44pub(crate) type MemtableVersionRef = Arc<MemtableVersion>;
45
46impl MemtableVersion {
47 pub(crate) fn new(mutable: TimePartitionsRef) -> MemtableVersion {
49 MemtableVersion {
50 mutable,
51 immutables: SmallVec::new(),
52 }
53 }
54
55 pub(crate) fn immutables(&self) -> &[MemtableRef] {
57 &self.immutables
58 }
59
60 pub(crate) fn list_memtables(&self) -> Vec<MemtableRef> {
62 let mut mems = Vec::with_capacity(self.immutables.len() + self.mutable.num_partitions());
63 self.mutable.list_memtables(&mut mems);
64 mems.extend_from_slice(&self.immutables);
65 mems
66 }
67
68 pub(crate) fn min_sequence(&self) -> Option<SequenceNumber> {
71 self.list_memtables()
72 .iter()
73 .filter(|mem| !mem.is_empty())
74 .map(|mem| mem.min_sequence())
75 .min()
76 }
77
78 pub(crate) fn freeze_mutable(
85 &self,
86 metadata: &RegionMetadataRef,
87 time_window: Option<Duration>,
88 ) -> Result<Option<MemtableVersion>> {
89 if self.mutable.is_empty() {
90 if Some(self.mutable.part_duration()) == time_window {
92 return Ok(None);
94 }
95
96 let mutable = self.mutable.new_with_part_duration(time_window, None);
98 common_telemetry::debug!(
99 "Freeze empty memtable, update partition duration from {:?} to {:?}",
100 self.mutable.part_duration(),
101 time_window
102 );
103 return Ok(Some(MemtableVersion {
104 mutable: Arc::new(mutable),
105 immutables: self.immutables.clone(),
106 }));
107 }
108
109 self.mutable.freeze()?;
112 if Some(self.mutable.part_duration()) != time_window {
114 common_telemetry::debug!(
115 "Fork memtable, update partition duration from {:?}, to {:?}",
116 self.mutable.part_duration(),
117 time_window
118 );
119 }
120 let mutable = Arc::new(self.mutable.fork(metadata, time_window));
121
122 let mut immutables =
123 SmallVec::with_capacity(self.immutables.len() + self.mutable.num_partitions());
124 immutables.extend(self.immutables.iter().cloned());
125 self.mutable.list_memtables_to_small_vec(&mut immutables);
127
128 Ok(Some(MemtableVersion {
129 mutable,
130 immutables,
131 }))
132 }
133
134 pub(crate) fn remove_memtables(&mut self, ids: &[MemtableId]) {
136 self.immutables = self
137 .immutables
138 .iter()
139 .filter(|mem| !ids.contains(&mem.id()))
140 .cloned()
141 .collect();
142 }
143
144 pub(crate) fn mutable_usage(&self) -> usize {
146 self.mutable.memory_usage()
147 }
148
149 pub(crate) fn immutables_usage(&self) -> usize {
151 self.immutables
152 .iter()
153 .map(|mem| mem.stats().estimated_bytes)
154 .sum()
155 }
156
157 pub(crate) fn num_rows(&self) -> u64 {
159 self.immutables
160 .iter()
161 .map(|mem| mem.stats().num_rows as u64)
162 .sum::<u64>()
163 + self.mutable.num_rows()
164 }
165
166 pub(crate) fn time_range(&self) -> Option<(Timestamp, Timestamp)> {
168 let mut mutables = Vec::new();
169 self.mutable.list_memtables(&mut mutables);
170 self.immutables
171 .iter()
172 .chain(mutables.iter())
173 .filter_map(|mem| mem.stats().time_range())
174 .reduce(|(min_a, max_a), (min_b, max_b)| (min_a.min(min_b), max_a.max(max_b)))
175 }
176
177 pub(crate) fn is_empty(&self) -> bool {
182 self.mutable.is_empty() && self.immutables.is_empty()
183 }
184}