hermes_core/segment/merger/
mod.rs1mod dense;
4mod fast_fields;
5mod postings;
6mod sparse;
7mod store;
8
9use std::sync::Arc;
10
11use rustc_hash::FxHashMap;
12
13use super::reader::SegmentReader;
14use super::types::{FieldStats, SegmentFiles, SegmentId, SegmentMeta};
15use super::{OffsetWriter, format_bytes};
16use crate::Result;
17use crate::directories::{Directory, DirectoryWriter};
18use crate::dsl::Schema;
19
20fn doc_offsets(segments: &[SegmentReader]) -> Result<Vec<u32>> {
24 let mut offsets = Vec::with_capacity(segments.len());
25 let mut acc = 0u32;
26 for seg in segments {
27 offsets.push(acc);
28 acc = acc.checked_add(seg.num_docs()).ok_or_else(|| {
29 crate::Error::Internal(format!(
30 "Total document count across segments exceeds u32::MAX ({})",
31 u32::MAX
32 ))
33 })?;
34 }
35 Ok(offsets)
36}
37
38#[derive(Debug, Clone, Default)]
40pub struct MergeStats {
41 pub terms_processed: usize,
43 pub term_dict_bytes: usize,
45 pub postings_bytes: usize,
47 pub store_bytes: usize,
49 pub vectors_bytes: usize,
51 pub sparse_bytes: usize,
53 pub bp_converged: bool,
59 pub fast_bytes: usize,
61}
62
63impl std::fmt::Display for MergeStats {
64 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
65 write!(
66 f,
67 "terms={}, term_dict={}, postings={}, store={}, vectors={}, sparse={}, fast={}",
68 self.terms_processed,
69 format_bytes(self.term_dict_bytes),
70 format_bytes(self.postings_bytes),
71 format_bytes(self.store_bytes),
72 format_bytes(self.vectors_bytes),
73 format_bytes(self.sparse_bytes),
74 format_bytes(self.fast_bytes),
75 )
76 }
77}
78
79pub use super::types::TrainedVectorStructures;
81
82pub(crate) fn block_in_place_if_multithread<R>(f: impl FnOnce() -> R) -> R {
86 if tokio::runtime::Handle::try_current()
87 .map(|h| h.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
88 .unwrap_or(false)
89 {
90 tokio::task::block_in_place(f)
91 } else {
92 f()
93 }
94}
95
96pub struct SegmentMerger {
98 schema: Arc<Schema>,
99 reorder_bmp: bool,
103 background_pool: Option<Arc<rayon::ThreadPool>>,
107 granularity: crate::segment::reorder::BpGranularity,
111 bp_budget: crate::segment::BpBudget,
116 bp_memory_budget: usize,
118 reorder_permits: Option<Arc<tokio::sync::Semaphore>>,
121}
122
123impl SegmentMerger {
124 pub fn new(schema: Arc<Schema>) -> Self {
125 Self {
126 schema,
127 reorder_bmp: false,
128 background_pool: None,
129 granularity: crate::segment::reorder::BpGranularity::Auto,
130 bp_budget: crate::segment::BpBudget::full(),
131 bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
132 reorder_permits: None,
133 }
134 }
135
136 pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
138 self.reorder_bmp = reorder;
139 self
140 }
141
142 pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
144 self.background_pool = pool;
145 self
146 }
147
148 pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
150 self.granularity = granularity;
151 self
152 }
153
154 pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
156 self.bp_budget = budget;
157 self
158 }
159
160 pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
162 self.bp_memory_budget = bytes;
163 self
164 }
165
166 pub fn with_reorder_permits(mut self, permits: Arc<tokio::sync::Semaphore>) -> Self {
168 self.reorder_permits = Some(permits);
169 self
170 }
171
172 pub async fn merge<D: Directory + DirectoryWriter>(
182 &self,
183 dir: &D,
184 segments: &[SegmentReader],
185 new_segment_id: SegmentId,
186 trained: Option<&TrainedVectorStructures>,
187 ) -> Result<(SegmentMeta, MergeStats)> {
188 let total_docs: u32 = segments
192 .iter()
193 .try_fold(0u32, |acc, segment| acc.checked_add(segment.num_docs()))
194 .ok_or_else(|| {
195 crate::Error::Internal(format!(
196 "Total document count exceeds u32::MAX ({})",
197 u32::MAX
198 ))
199 })?;
200
201 let mut stats = MergeStats::default();
202 let files = SegmentFiles::new(new_segment_id.0);
203
204 let merge_start = std::time::Instant::now();
218
219 let postings_fut = async {
221 let mut postings_writer =
222 OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
223 let mut positions_writer =
224 OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
225 let mut term_dict_writer =
226 OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
227
228 let terms_processed = self
229 .merge_postings(
230 segments,
231 &mut term_dict_writer,
232 &mut postings_writer,
233 &mut positions_writer,
234 )
235 .await?;
236
237 let postings_bytes = postings_writer.offset() as usize;
238 let term_dict_bytes = term_dict_writer.offset() as usize;
239 let positions_bytes = positions_writer.offset();
240
241 postings_writer.finish()?;
242 term_dict_writer.finish()?;
243 if positions_bytes > 0 {
244 positions_writer.finish()?;
245 } else {
246 drop(positions_writer);
247 let _ = dir.delete(&files.positions).await;
248 }
249 log::info!(
250 "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
251 terms_processed,
252 format_bytes(term_dict_bytes),
253 format_bytes(postings_bytes),
254 format_bytes(positions_bytes as usize),
255 );
256 Ok::<(usize, usize, usize), crate::Error>((
257 terms_processed,
258 term_dict_bytes,
259 postings_bytes,
260 ))
261 };
262
263 let store_fut = async {
264 let mut store_writer =
265 OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
266 let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
267 let bytes = store_writer.offset() as usize;
268 store_writer.finish()?;
269 Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
270 };
271
272 let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
273
274 let (postings_result, store_result, fast_bytes) =
275 tokio::try_join!(postings_fut, store_fut, fast_fut)?;
276
277 log::info!(
278 "[merge] stage 1 done in {:.1}s (postings + store + fast)",
279 merge_start.elapsed().as_secs_f64()
280 );
281
282 let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
286
287 let dense_fut = async {
288 self.merge_dense_vectors(dir, segments, &files, trained)
289 .await
290 };
291
292 let ((sparse_bytes, bp_converged), vectors_bytes) = if self.reorder_bmp {
296 let sparse = sparse_fut.await?;
297 let dense = dense_fut.await?;
298 (sparse, dense)
299 } else {
300 tokio::try_join!(sparse_fut, dense_fut)?
301 };
302 let (store_bytes, store_num_docs) = store_result;
303 stats.terms_processed = postings_result.0;
304 stats.term_dict_bytes = postings_result.1;
305 stats.postings_bytes = postings_result.2;
306 stats.store_bytes = store_bytes;
307 stats.vectors_bytes = vectors_bytes;
308 stats.sparse_bytes = sparse_bytes;
309 stats.bp_converged = bp_converged;
310 stats.fast_bytes = fast_bytes;
311 log::info!(
312 "[merge] all phases done in {:.1}s: {}",
313 merge_start.elapsed().as_secs_f64(),
314 stats
315 );
316
317 let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
319 for segment in segments {
320 for (&field_id, field_stats) in &segment.meta().field_stats {
321 let entry = merged_field_stats.entry(field_id).or_default();
322 entry.total_tokens = entry
323 .total_tokens
324 .checked_add(field_stats.total_tokens)
325 .ok_or_else(|| {
326 crate::Error::Corruption(format!(
327 "field {} total-token count overflow while merging",
328 field_id
329 ))
330 })?;
331 entry.doc_count = entry
332 .doc_count
333 .checked_add(field_stats.doc_count)
334 .ok_or_else(|| {
335 crate::Error::Corruption(format!(
336 "field {} document count overflow while merging",
337 field_id
338 ))
339 })?;
340 }
341 }
342
343 if store_num_docs != total_docs {
347 log::error!(
348 "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
349 Per-segment: {:?}",
350 store_num_docs,
351 total_docs,
352 segments
353 .iter()
354 .map(|s| (
355 format!("{:016x}", s.meta().id),
356 s.num_docs(),
357 s.store().num_docs()
358 ))
359 .collect::<Vec<_>>()
360 );
361 return Err(crate::Error::Io(std::io::Error::new(
362 std::io::ErrorKind::InvalidData,
363 format!(
364 "Store/meta doc count mismatch: store={}, meta={}",
365 store_num_docs, total_docs
366 ),
367 )));
368 }
369
370 let meta = SegmentMeta {
371 id: new_segment_id.0,
372 num_docs: total_docs,
373 field_stats: merged_field_stats,
374 };
375
376 dir.write(&files.meta, &meta.serialize()?).await?;
377
378 let label = if trained.is_some() {
379 "ANN merge"
380 } else {
381 "Merge"
382 };
383 log::info!("{} complete: {} docs, {}", label, total_docs, stats);
384
385 Ok((meta, stats))
386 }
387}
388
389pub async fn delete_segment<D: Directory + DirectoryWriter>(
391 dir: &D,
392 segment_id: SegmentId,
393) -> Result<()> {
394 let files = SegmentFiles::new(segment_id.0);
395 let paths = files.lifecycle_paths();
396 let results = futures::future::join_all(paths.iter().map(|path| dir.delete(path))).await;
397
398 for result in results {
402 if let Err(error) = result
403 && error.kind() != std::io::ErrorKind::NotFound
404 {
405 return Err(crate::Error::Io(error));
406 }
407 }
408 Ok(())
409}