hermes_core/segment/merger/
mod.rs1mod dense;
4#[cfg(feature = "diagnostics")]
5mod diagnostics;
6mod fast_fields;
7mod postings;
8mod sparse;
9mod store;
10
11use std::sync::Arc;
12
13use rustc_hash::FxHashMap;
14
15use super::reader::SegmentReader;
16use super::types::{FieldStats, SegmentFiles, SegmentId, SegmentMeta};
17use super::{OffsetWriter, format_bytes};
18use crate::Result;
19use crate::directories::{Directory, DirectoryWriter};
20use crate::dsl::Schema;
21
22fn doc_offsets(segments: &[SegmentReader]) -> Result<Vec<u32>> {
26 let mut offsets = Vec::with_capacity(segments.len());
27 let mut acc = 0u32;
28 for seg in segments {
29 offsets.push(acc);
30 acc = acc.checked_add(seg.num_docs()).ok_or_else(|| {
31 crate::Error::Internal(format!(
32 "Total document count across segments exceeds u32::MAX ({})",
33 u32::MAX
34 ))
35 })?;
36 }
37 Ok(offsets)
38}
39
40#[derive(Debug, Clone, Default)]
42pub struct MergeStats {
43 pub terms_processed: usize,
45 pub term_dict_bytes: usize,
47 pub postings_bytes: usize,
49 pub store_bytes: usize,
51 pub vectors_bytes: usize,
53 pub sparse_bytes: usize,
55 pub fast_bytes: usize,
57}
58
59impl std::fmt::Display for MergeStats {
60 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
61 write!(
62 f,
63 "terms={}, term_dict={}, postings={}, store={}, vectors={}, sparse={}, fast={}",
64 self.terms_processed,
65 format_bytes(self.term_dict_bytes),
66 format_bytes(self.postings_bytes),
67 format_bytes(self.store_bytes),
68 format_bytes(self.vectors_bytes),
69 format_bytes(self.sparse_bytes),
70 format_bytes(self.fast_bytes),
71 )
72 }
73}
74
75pub use super::types::TrainedVectorStructures;
77
78pub(crate) fn block_in_place_if_multithread<R>(f: impl FnOnce() -> R) -> R {
82 if tokio::runtime::Handle::try_current()
83 .map(|h| h.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
84 .unwrap_or(false)
85 {
86 tokio::task::block_in_place(f)
87 } else {
88 f()
89 }
90}
91
92pub struct SegmentMerger {
94 schema: Arc<Schema>,
95 reorder_bmp: bool,
99 background_pool: Option<Arc<rayon::ThreadPool>>,
103}
104
105impl SegmentMerger {
106 pub fn new(schema: Arc<Schema>) -> Self {
107 Self {
108 schema,
109 reorder_bmp: false,
110 background_pool: None,
111 }
112 }
113
114 pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
116 self.reorder_bmp = reorder;
117 self
118 }
119
120 pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
122 self.background_pool = pool;
123 self
124 }
125
126 pub async fn merge<D: Directory + DirectoryWriter>(
136 &self,
137 dir: &D,
138 segments: &[SegmentReader],
139 new_segment_id: SegmentId,
140 trained: Option<&TrainedVectorStructures>,
141 ) -> Result<(SegmentMeta, MergeStats)> {
142 let mut stats = MergeStats::default();
143 let files = SegmentFiles::new(new_segment_id.0);
144
145 let merge_start = std::time::Instant::now();
158
159 let postings_fut = async {
161 let mut postings_writer =
162 OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
163 let mut positions_writer =
164 OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
165 let mut term_dict_writer =
166 OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
167
168 let terms_processed = self
169 .merge_postings(
170 segments,
171 &mut term_dict_writer,
172 &mut postings_writer,
173 &mut positions_writer,
174 )
175 .await?;
176
177 let postings_bytes = postings_writer.offset() as usize;
178 let term_dict_bytes = term_dict_writer.offset() as usize;
179 let positions_bytes = positions_writer.offset();
180
181 postings_writer.finish()?;
182 term_dict_writer.finish()?;
183 if positions_bytes > 0 {
184 positions_writer.finish()?;
185 } else {
186 drop(positions_writer);
187 let _ = dir.delete(&files.positions).await;
188 }
189 log::info!(
190 "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
191 terms_processed,
192 format_bytes(term_dict_bytes),
193 format_bytes(postings_bytes),
194 format_bytes(positions_bytes as usize),
195 );
196 Ok::<(usize, usize, usize), crate::Error>((
197 terms_processed,
198 term_dict_bytes,
199 postings_bytes,
200 ))
201 };
202
203 let store_fut = async {
204 let mut store_writer =
205 OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
206 let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
207 let bytes = store_writer.offset() as usize;
208 store_writer.finish()?;
209 Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
210 };
211
212 let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
213
214 let (postings_result, store_result, fast_bytes) =
215 tokio::try_join!(postings_fut, store_fut, fast_fut)?;
216
217 log::info!(
218 "[merge] stage 1 done in {:.1}s (postings + store + fast)",
219 merge_start.elapsed().as_secs_f64()
220 );
221
222 let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
226
227 let dense_fut = async {
228 self.merge_dense_vectors(dir, segments, &files, trained)
229 .await
230 };
231
232 let (sparse_bytes, vectors_bytes) = tokio::try_join!(sparse_fut, dense_fut)?;
233 let (store_bytes, store_num_docs) = store_result;
234 stats.terms_processed = postings_result.0;
235 stats.term_dict_bytes = postings_result.1;
236 stats.postings_bytes = postings_result.2;
237 stats.store_bytes = store_bytes;
238 stats.vectors_bytes = vectors_bytes;
239 stats.sparse_bytes = sparse_bytes;
240 stats.fast_bytes = fast_bytes;
241 log::info!(
242 "[merge] all phases done in {:.1}s: {}",
243 merge_start.elapsed().as_secs_f64(),
244 stats
245 );
246
247 let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
249 for segment in segments {
250 for (&field_id, field_stats) in &segment.meta().field_stats {
251 let entry = merged_field_stats.entry(field_id).or_default();
252 entry.total_tokens += field_stats.total_tokens;
253 entry.doc_count += field_stats.doc_count;
254 }
255 }
256
257 let total_docs: u32 = segments
258 .iter()
259 .try_fold(0u32, |acc, s| acc.checked_add(s.num_docs()))
260 .ok_or_else(|| {
261 crate::Error::Internal(format!(
262 "Total document count exceeds u32::MAX ({})",
263 u32::MAX
264 ))
265 })?;
266
267 if store_num_docs != total_docs {
271 log::error!(
272 "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
273 Per-segment: {:?}",
274 store_num_docs,
275 total_docs,
276 segments
277 .iter()
278 .map(|s| (
279 format!("{:016x}", s.meta().id),
280 s.num_docs(),
281 s.store().num_docs()
282 ))
283 .collect::<Vec<_>>()
284 );
285 return Err(crate::Error::Io(std::io::Error::new(
286 std::io::ErrorKind::InvalidData,
287 format!(
288 "Store/meta doc count mismatch: store={}, meta={}",
289 store_num_docs, total_docs
290 ),
291 )));
292 }
293
294 let meta = SegmentMeta {
295 id: new_segment_id.0,
296 num_docs: total_docs,
297 field_stats: merged_field_stats,
298 };
299
300 dir.write(&files.meta, &meta.serialize()?).await?;
301
302 let label = if trained.is_some() {
303 "ANN merge"
304 } else {
305 "Merge"
306 };
307 log::info!("{} complete: {} docs, {}", label, total_docs, stats);
308
309 Ok((meta, stats))
310 }
311}
312
313pub async fn delete_segment<D: Directory + DirectoryWriter>(
315 dir: &D,
316 segment_id: SegmentId,
317) -> Result<()> {
318 let files = SegmentFiles::new(segment_id.0);
319 let _ = tokio::join!(
320 dir.delete(&files.term_dict),
321 dir.delete(&files.postings),
322 dir.delete(&files.store),
323 dir.delete(&files.meta),
324 dir.delete(&files.vectors),
325 dir.delete(&files.sparse),
326 dir.delete(&files.positions),
327 dir.delete(&files.fast),
328 );
329 Ok(())
330}