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 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}
121
122impl SegmentMerger {
123 pub fn new(schema: Arc<Schema>) -> Self {
124 Self {
125 schema,
126 reorder_bmp: false,
127 background_pool: None,
128 granularity: crate::segment::reorder::BpGranularity::Auto,
129 bp_budget: crate::segment::BpBudget::full(),
130 bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
131 }
132 }
133
134 pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
136 self.reorder_bmp = reorder;
137 self
138 }
139
140 pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
142 self.background_pool = pool;
143 self
144 }
145
146 pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
148 self.granularity = granularity;
149 self
150 }
151
152 pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
154 self.bp_budget = budget;
155 self
156 }
157
158 pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
160 self.bp_memory_budget = bytes;
161 self
162 }
163
164 pub async fn merge<D: Directory + DirectoryWriter>(
174 &self,
175 dir: &D,
176 segments: &[SegmentReader],
177 new_segment_id: SegmentId,
178 trained: Option<&TrainedVectorStructures>,
179 ) -> Result<(SegmentMeta, MergeStats)> {
180 let mut stats = MergeStats::default();
181 let files = SegmentFiles::new(new_segment_id.0);
182
183 let merge_start = std::time::Instant::now();
196
197 let postings_fut = async {
199 let mut postings_writer =
200 OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
201 let mut positions_writer =
202 OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
203 let mut term_dict_writer =
204 OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
205
206 let terms_processed = self
207 .merge_postings(
208 segments,
209 &mut term_dict_writer,
210 &mut postings_writer,
211 &mut positions_writer,
212 )
213 .await?;
214
215 let postings_bytes = postings_writer.offset() as usize;
216 let term_dict_bytes = term_dict_writer.offset() as usize;
217 let positions_bytes = positions_writer.offset();
218
219 postings_writer.finish()?;
220 term_dict_writer.finish()?;
221 if positions_bytes > 0 {
222 positions_writer.finish()?;
223 } else {
224 drop(positions_writer);
225 let _ = dir.delete(&files.positions).await;
226 }
227 log::info!(
228 "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
229 terms_processed,
230 format_bytes(term_dict_bytes),
231 format_bytes(postings_bytes),
232 format_bytes(positions_bytes as usize),
233 );
234 Ok::<(usize, usize, usize), crate::Error>((
235 terms_processed,
236 term_dict_bytes,
237 postings_bytes,
238 ))
239 };
240
241 let store_fut = async {
242 let mut store_writer =
243 OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
244 let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
245 let bytes = store_writer.offset() as usize;
246 store_writer.finish()?;
247 Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
248 };
249
250 let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
251
252 let (postings_result, store_result, fast_bytes) =
253 tokio::try_join!(postings_fut, store_fut, fast_fut)?;
254
255 log::info!(
256 "[merge] stage 1 done in {:.1}s (postings + store + fast)",
257 merge_start.elapsed().as_secs_f64()
258 );
259
260 let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
264
265 let dense_fut = async {
266 self.merge_dense_vectors(dir, segments, &files, trained)
267 .await
268 };
269
270 let ((sparse_bytes, bp_converged), vectors_bytes) =
271 tokio::try_join!(sparse_fut, dense_fut)?;
272 let (store_bytes, store_num_docs) = store_result;
273 stats.terms_processed = postings_result.0;
274 stats.term_dict_bytes = postings_result.1;
275 stats.postings_bytes = postings_result.2;
276 stats.store_bytes = store_bytes;
277 stats.vectors_bytes = vectors_bytes;
278 stats.sparse_bytes = sparse_bytes;
279 stats.bp_converged = bp_converged;
280 stats.fast_bytes = fast_bytes;
281 log::info!(
282 "[merge] all phases done in {:.1}s: {}",
283 merge_start.elapsed().as_secs_f64(),
284 stats
285 );
286
287 let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
289 for segment in segments {
290 for (&field_id, field_stats) in &segment.meta().field_stats {
291 let entry = merged_field_stats.entry(field_id).or_default();
292 entry.total_tokens += field_stats.total_tokens;
293 entry.doc_count += field_stats.doc_count;
294 }
295 }
296
297 let total_docs: u32 = segments
298 .iter()
299 .try_fold(0u32, |acc, s| acc.checked_add(s.num_docs()))
300 .ok_or_else(|| {
301 crate::Error::Internal(format!(
302 "Total document count exceeds u32::MAX ({})",
303 u32::MAX
304 ))
305 })?;
306
307 if store_num_docs != total_docs {
311 log::error!(
312 "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
313 Per-segment: {:?}",
314 store_num_docs,
315 total_docs,
316 segments
317 .iter()
318 .map(|s| (
319 format!("{:016x}", s.meta().id),
320 s.num_docs(),
321 s.store().num_docs()
322 ))
323 .collect::<Vec<_>>()
324 );
325 return Err(crate::Error::Io(std::io::Error::new(
326 std::io::ErrorKind::InvalidData,
327 format!(
328 "Store/meta doc count mismatch: store={}, meta={}",
329 store_num_docs, total_docs
330 ),
331 )));
332 }
333
334 let meta = SegmentMeta {
335 id: new_segment_id.0,
336 num_docs: total_docs,
337 field_stats: merged_field_stats,
338 };
339
340 dir.write(&files.meta, &meta.serialize()?).await?;
341
342 let label = if trained.is_some() {
343 "ANN merge"
344 } else {
345 "Merge"
346 };
347 log::info!("{} complete: {} docs, {}", label, total_docs, stats);
348
349 Ok((meta, stats))
350 }
351}
352
353pub async fn delete_segment<D: Directory + DirectoryWriter>(
355 dir: &D,
356 segment_id: SegmentId,
357) -> Result<()> {
358 let files = SegmentFiles::new(segment_id.0);
359 let _ = tokio::join!(
360 dir.delete(&files.term_dict),
361 dir.delete(&files.postings),
362 dir.delete(&files.store),
363 dir.delete(&files.meta),
364 dir.delete(&files.vectors),
365 dir.delete(&files.sparse),
366 dir.delete(&files.positions),
367 dir.delete(&files.fast),
368 );
369 Ok(())
370}