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::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 bp_converged: bool,
61 pub fast_bytes: usize,
63}
64
65impl std::fmt::Display for MergeStats {
66 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
67 write!(
68 f,
69 "terms={}, term_dict={}, postings={}, store={}, vectors={}, sparse={}, fast={}",
70 self.terms_processed,
71 format_bytes(self.term_dict_bytes),
72 format_bytes(self.postings_bytes),
73 format_bytes(self.store_bytes),
74 format_bytes(self.vectors_bytes),
75 format_bytes(self.sparse_bytes),
76 format_bytes(self.fast_bytes),
77 )
78 }
79}
80
81pub use super::types::TrainedVectorStructures;
83
84pub(crate) fn block_in_place_if_multithread<R>(f: impl FnOnce() -> R) -> R {
88 if tokio::runtime::Handle::try_current()
89 .map(|h| h.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
90 .unwrap_or(false)
91 {
92 tokio::task::block_in_place(f)
93 } else {
94 f()
95 }
96}
97
98pub struct SegmentMerger {
100 schema: Arc<Schema>,
101 reorder_bmp: bool,
105 background_pool: Option<Arc<rayon::ThreadPool>>,
109 granularity: crate::segment::reorder::BpGranularity,
113 bp_budget: crate::segment::BpBudget,
118 bp_memory_budget: usize,
120 reorder_permits: Option<Arc<tokio::sync::Semaphore>>,
123}
124
125impl SegmentMerger {
126 pub fn new(schema: Arc<Schema>) -> Self {
127 Self {
128 schema,
129 reorder_bmp: false,
130 background_pool: None,
131 granularity: crate::segment::reorder::BpGranularity::Auto,
132 bp_budget: crate::segment::BpBudget::full(),
133 bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
134 reorder_permits: None,
135 }
136 }
137
138 pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
140 self.reorder_bmp = reorder;
141 self
142 }
143
144 pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
146 self.background_pool = pool;
147 self
148 }
149
150 pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
152 self.granularity = granularity;
153 self
154 }
155
156 pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
158 self.bp_budget = budget;
159 self
160 }
161
162 pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
164 self.bp_memory_budget = bytes;
165 self
166 }
167
168 pub fn with_reorder_permits(mut self, permits: Arc<tokio::sync::Semaphore>) -> Self {
170 self.reorder_permits = Some(permits);
171 self
172 }
173
174 pub async fn merge<D: Directory + DirectoryWriter>(
184 &self,
185 dir: &D,
186 segments: &[SegmentReader],
187 new_segment_id: SegmentId,
188 trained: Option<&TrainedVectorStructures>,
189 ) -> Result<(SegmentMeta, MergeStats)> {
190 let total_docs: u32 = segments
194 .iter()
195 .try_fold(0u32, |acc, segment| acc.checked_add(segment.num_docs()))
196 .ok_or_else(|| {
197 crate::Error::Internal(format!(
198 "Total document count exceeds u32::MAX ({})",
199 u32::MAX
200 ))
201 })?;
202
203 let mut stats = MergeStats::default();
204 let files = SegmentFiles::new(new_segment_id.0);
205
206 let merge_start = std::time::Instant::now();
220
221 let postings_fut = async {
223 let mut postings_writer =
224 OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
225 let mut positions_writer =
226 OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
227 let mut term_dict_writer =
228 OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
229
230 let terms_processed = self
231 .merge_postings(
232 segments,
233 &mut term_dict_writer,
234 &mut postings_writer,
235 &mut positions_writer,
236 )
237 .await?;
238
239 let postings_bytes = postings_writer.offset() as usize;
240 let term_dict_bytes = term_dict_writer.offset() as usize;
241 let positions_bytes = positions_writer.offset();
242
243 postings_writer.finish()?;
244 term_dict_writer.finish()?;
245 if positions_bytes > 0 {
246 positions_writer.finish()?;
247 } else {
248 drop(positions_writer);
249 let _ = dir.delete(&files.positions).await;
250 }
251 log::info!(
252 "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
253 terms_processed,
254 format_bytes(term_dict_bytes),
255 format_bytes(postings_bytes),
256 format_bytes(positions_bytes as usize),
257 );
258 Ok::<(usize, usize, usize), crate::Error>((
259 terms_processed,
260 term_dict_bytes,
261 postings_bytes,
262 ))
263 };
264
265 let store_fut = async {
266 let mut store_writer =
267 OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
268 let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
269 let bytes = store_writer.offset() as usize;
270 store_writer.finish()?;
271 Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
272 };
273
274 let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
275
276 let (postings_result, store_result, fast_bytes) =
277 tokio::try_join!(postings_fut, store_fut, fast_fut)?;
278
279 log::info!(
280 "[merge] stage 1 done in {:.1}s (postings + store + fast)",
281 merge_start.elapsed().as_secs_f64()
282 );
283
284 let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
288
289 let dense_fut = async {
290 self.merge_dense_vectors(dir, segments, &files, trained, AnnWriteMode::Copy)
291 .await
292 };
293
294 let ((sparse_bytes, bp_converged), vectors_bytes) = if self.reorder_bmp {
298 let sparse = sparse_fut.await?;
299 let dense = dense_fut.await?;
300 (sparse, dense)
301 } else {
302 tokio::try_join!(sparse_fut, dense_fut)?
303 };
304 let (store_bytes, store_num_docs) = store_result;
305 stats.terms_processed = postings_result.0;
306 stats.term_dict_bytes = postings_result.1;
307 stats.postings_bytes = postings_result.2;
308 stats.store_bytes = store_bytes;
309 stats.vectors_bytes = vectors_bytes;
310 stats.sparse_bytes = sparse_bytes;
311 stats.bp_converged = bp_converged;
312 stats.fast_bytes = fast_bytes;
313 log::info!(
314 "[merge] all phases done in {:.1}s: {}",
315 merge_start.elapsed().as_secs_f64(),
316 stats
317 );
318
319 let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
321 for segment in segments {
322 for (&field_id, field_stats) in &segment.meta().field_stats {
323 let entry = merged_field_stats.entry(field_id).or_default();
324 entry.total_tokens = entry
325 .total_tokens
326 .checked_add(field_stats.total_tokens)
327 .ok_or_else(|| {
328 crate::Error::Corruption(format!(
329 "field {} total-token count overflow while merging",
330 field_id
331 ))
332 })?;
333 entry.doc_count = entry
334 .doc_count
335 .checked_add(field_stats.doc_count)
336 .ok_or_else(|| {
337 crate::Error::Corruption(format!(
338 "field {} document count overflow while merging",
339 field_id
340 ))
341 })?;
342 }
343 }
344
345 if store_num_docs != total_docs {
349 log::error!(
350 "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
351 Per-segment: {:?}",
352 store_num_docs,
353 total_docs,
354 segments
355 .iter()
356 .map(|s| (
357 format!("{:016x}", s.meta().id),
358 s.num_docs(),
359 s.store().num_docs()
360 ))
361 .collect::<Vec<_>>()
362 );
363 return Err(crate::Error::Io(std::io::Error::new(
364 std::io::ErrorKind::InvalidData,
365 format!(
366 "Store/meta doc count mismatch: store={}, meta={}",
367 store_num_docs, total_docs
368 ),
369 )));
370 }
371
372 let meta = SegmentMeta {
373 id: new_segment_id.0,
374 num_docs: total_docs,
375 field_stats: merged_field_stats,
376 };
377
378 dir.write_durable(&files.meta, &meta.serialize()?).await?;
382
383 let label = if trained.is_some() {
384 "ANN merge"
385 } else {
386 "Merge"
387 };
388 log::info!("{} complete: {} docs, {}", label, total_docs, stats);
389
390 Ok((meta, stats))
391 }
392}
393
394pub async fn delete_segment<D: Directory + DirectoryWriter>(
396 dir: &D,
397 segment_id: SegmentId,
398) -> Result<()> {
399 let files = SegmentFiles::new(segment_id.0);
400 let paths = files.lifecycle_paths();
401 let results = futures::future::join_all(paths.iter().map(|path| dir.delete(path))).await;
402
403 for result in results {
407 if let Err(error) = result
408 && error.kind() != std::io::ErrorKind::NotFound
409 {
410 return Err(crate::Error::Io(error));
411 }
412 }
413 Ok(())
414}