hermes_core/segment/merger/
mod.rs1mod dense;
4mod fast_fields;
5mod postings;
6mod sparse;
7mod store;
8
9pub(crate) use dense::AnnWriteMode;
10
11use std::sync::Arc;
12
13use rustc_hash::FxHashMap;
14
15use super::OffsetWriter;
16use super::reader::SegmentReader;
17use super::types::{FieldStats, SegmentFiles, SegmentId, SegmentMeta};
18use crate::Result;
19use crate::directories::{Directory, DirectoryWriter};
20use crate::dsl::Schema;
21use crate::index::{ReorderConcurrencyGate, ReorderPriority};
22
23fn doc_offsets(segments: &[SegmentReader]) -> Result<Vec<u32>> {
27 let mut offsets = Vec::with_capacity(segments.len());
28 let mut acc = 0u32;
29 for seg in segments {
30 offsets.push(acc);
31 acc = acc.checked_add(seg.num_docs()).ok_or_else(|| {
32 crate::Error::Internal(format!(
33 "Total document count across segments exceeds u32::MAX ({})",
34 u32::MAX
35 ))
36 })?;
37 }
38 Ok(offsets)
39}
40
41#[derive(Debug, Clone, Default)]
43pub struct MergeStats {
44 pub terms_processed: usize,
46 pub term_dict_bytes: usize,
48 pub postings_bytes: usize,
50 pub store_bytes: usize,
52 pub vectors_bytes: usize,
54 pub sparse_bytes: usize,
56 pub bp_converged: bool,
62 pub fast_bytes: usize,
64}
65
66impl std::fmt::Display for MergeStats {
67 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
68 write!(
69 f,
70 "terms={}, term_dict={}, postings={}, store={}, dense_vectors={}, sparse_vectors={}, fast_fields={}",
71 self.terms_processed,
72 crate::format_bytes(self.term_dict_bytes as u64),
73 crate::format_bytes(self.postings_bytes as u64),
74 crate::format_bytes(self.store_bytes as u64),
75 crate::format_bytes(self.vectors_bytes as u64),
76 crate::format_bytes(self.sparse_bytes as u64),
77 crate::format_bytes(self.fast_bytes as u64),
78 )
79 }
80}
81
82pub use super::types::TrainedVectorStructures;
84
85pub(crate) fn block_in_place_if_multithread<R>(f: impl FnOnce() -> R) -> R {
89 if tokio::runtime::Handle::try_current()
90 .map(|h| h.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
91 .unwrap_or(false)
92 {
93 tokio::task::block_in_place(f)
94 } else {
95 f()
96 }
97}
98
99pub struct SegmentMerger {
101 schema: Arc<Schema>,
102 reorder_bmp: bool,
106 background_pool: Option<Arc<rayon::ThreadPool>>,
110 granularity: crate::segment::reorder::BpGranularity,
114 bp_budget: crate::segment::BpBudget,
119 bp_memory_budget: usize,
121 reorder_permits: Option<Arc<ReorderConcurrencyGate>>,
124 reorder_priority: ReorderPriority,
127}
128
129impl SegmentMerger {
130 pub fn new(schema: Arc<Schema>) -> Self {
131 Self {
132 schema,
133 reorder_bmp: false,
134 background_pool: None,
135 granularity: crate::segment::reorder::BpGranularity::Auto,
136 bp_budget: crate::segment::BpBudget::full(),
137 bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
138 reorder_permits: None,
139 reorder_priority: ReorderPriority::Background,
140 }
141 }
142
143 pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
145 self.reorder_bmp = reorder;
146 self
147 }
148
149 pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
151 self.background_pool = pool;
152 self
153 }
154
155 pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
157 self.granularity = granularity;
158 self
159 }
160
161 pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
163 self.bp_budget = budget;
164 self
165 }
166
167 pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
169 self.bp_memory_budget = bytes;
170 self
171 }
172
173 pub fn with_reorder_permits(mut self, permits: Arc<ReorderConcurrencyGate>) -> Self {
175 self.reorder_permits = Some(permits);
176 self
177 }
178
179 pub(crate) fn with_reorder_priority(mut self, priority: ReorderPriority) -> Self {
180 self.reorder_priority = priority;
181 self
182 }
183
184 pub async fn merge<D: Directory + DirectoryWriter>(
194 &self,
195 dir: &D,
196 segments: &[SegmentReader],
197 new_segment_id: SegmentId,
198 trained: Option<&TrainedVectorStructures>,
199 ) -> Result<(SegmentMeta, MergeStats)> {
200 let total_docs: u32 = segments
204 .iter()
205 .try_fold(0u32, |acc, segment| acc.checked_add(segment.num_docs()))
206 .ok_or_else(|| {
207 crate::Error::Internal(format!(
208 "Total document count exceeds u32::MAX ({})",
209 u32::MAX
210 ))
211 })?;
212
213 let mut stats = MergeStats::default();
214 let files = SegmentFiles::new(new_segment_id.0);
215
216 let merge_start = std::time::Instant::now();
230
231 let postings_fut = async {
233 let mut postings_writer =
234 OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
235 let mut positions_writer =
236 OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
237 let mut term_dict_writer =
238 OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
239
240 let terms_processed = self
241 .merge_postings(
242 segments,
243 &mut term_dict_writer,
244 &mut postings_writer,
245 &mut positions_writer,
246 )
247 .await?;
248
249 let postings_bytes = postings_writer.offset() as usize;
250 let term_dict_bytes = term_dict_writer.offset() as usize;
251 let positions_bytes = positions_writer.offset();
252
253 postings_writer.finish()?;
254 term_dict_writer.finish()?;
255 if positions_bytes > 0 {
256 positions_writer.finish()?;
257 } else {
258 drop(positions_writer);
259 let _ = dir.delete(&files.positions).await;
260 }
261 log::info!(
262 "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
263 terms_processed,
264 crate::format_bytes(term_dict_bytes as u64),
265 crate::format_bytes(postings_bytes as u64),
266 crate::format_bytes(positions_bytes),
267 );
268 Ok::<(usize, usize, usize), crate::Error>((
269 terms_processed,
270 term_dict_bytes,
271 postings_bytes,
272 ))
273 };
274
275 let store_fut = async {
276 let mut store_writer =
277 OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
278 let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
279 let bytes = store_writer.offset() as usize;
280 store_writer.finish()?;
281 Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
282 };
283
284 let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
285
286 let (postings_result, store_result, fast_bytes) =
287 tokio::try_join!(postings_fut, store_fut, fast_fut)?;
288
289 log::info!(
290 "[merge] stage 1 done in {:.1}s (postings + store + fast)",
291 merge_start.elapsed().as_secs_f64()
292 );
293
294 let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
298
299 let dense_fut = async {
300 self.merge_dense_vectors(dir, segments, &files, trained, AnnWriteMode::Copy)
301 .await
302 };
303
304 let ((sparse_bytes, bp_converged), vectors_bytes) = if self.reorder_bmp {
308 let sparse = sparse_fut.await?;
309 let dense = dense_fut.await?;
310 (sparse, dense)
311 } else {
312 tokio::try_join!(sparse_fut, dense_fut)?
313 };
314 let (store_bytes, store_num_docs) = store_result;
315 stats.terms_processed = postings_result.0;
316 stats.term_dict_bytes = postings_result.1;
317 stats.postings_bytes = postings_result.2;
318 stats.store_bytes = store_bytes;
319 stats.vectors_bytes = vectors_bytes;
320 stats.sparse_bytes = sparse_bytes;
321 stats.bp_converged = bp_converged;
322 stats.fast_bytes = fast_bytes;
323 log::info!(
324 "[merge] all phases done in {:.1}s: {}",
325 merge_start.elapsed().as_secs_f64(),
326 stats
327 );
328
329 let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
331 for segment in segments {
332 for (&field_id, field_stats) in &segment.meta().field_stats {
333 let entry = merged_field_stats.entry(field_id).or_default();
334 entry.total_tokens = entry
335 .total_tokens
336 .checked_add(field_stats.total_tokens)
337 .ok_or_else(|| {
338 crate::Error::Corruption(format!(
339 "field {} total-token count overflow while merging",
340 field_id
341 ))
342 })?;
343 entry.doc_count = entry
344 .doc_count
345 .checked_add(field_stats.doc_count)
346 .ok_or_else(|| {
347 crate::Error::Corruption(format!(
348 "field {} document count overflow while merging",
349 field_id
350 ))
351 })?;
352 }
353 }
354
355 if store_num_docs != total_docs {
359 log::error!(
360 "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
361 Per-segment: {:?}",
362 store_num_docs,
363 total_docs,
364 segments
365 .iter()
366 .map(|s| (
367 format!("{:016x}", s.meta().id),
368 s.num_docs(),
369 s.store().num_docs()
370 ))
371 .collect::<Vec<_>>()
372 );
373 return Err(crate::Error::Io(std::io::Error::new(
374 std::io::ErrorKind::InvalidData,
375 format!(
376 "Store/meta doc count mismatch: store={}, meta={}",
377 store_num_docs, total_docs
378 ),
379 )));
380 }
381
382 let meta = SegmentMeta {
383 id: new_segment_id.0,
384 num_docs: total_docs,
385 field_stats: merged_field_stats,
386 };
387
388 dir.write_durable(&files.meta, &meta.serialize()?).await?;
392
393 let label = if trained.is_some() {
394 "ANN merge"
395 } else {
396 "Merge"
397 };
398 log::info!("{} complete: {} docs, {}", label, total_docs, stats);
399
400 Ok((meta, stats))
401 }
402}
403
404pub async fn delete_segment<D: Directory + DirectoryWriter>(
406 dir: &D,
407 segment_id: SegmentId,
408) -> Result<()> {
409 let files = SegmentFiles::new(segment_id.0);
410 let paths = files.lifecycle_paths();
411 let results = futures::future::join_all(paths.iter().map(|path| dir.delete(path))).await;
412
413 for result in results {
417 if let Err(error) = result
418 && error.kind() != std::io::ErrorKind::NotFound
419 {
420 return Err(crate::Error::Io(error));
421 }
422 }
423 Ok(())
424}