@@ -4,8 +4,9 @@ use std::io::Write;
44use std:: ops:: Deref ;
55use std:: path:: PathBuf ;
66use std:: sync:: atomic:: { AtomicBool , Ordering } ;
7- use std:: sync:: { Arc , RwLock } ;
7+ use std:: sync:: Arc ;
88
9+ use parking_lot:: RwLock ;
910use rayon:: { ThreadPool , ThreadPoolBuilder } ;
1011
1112use super :: segment_manager:: SegmentManager ;
@@ -316,6 +317,12 @@ pub fn merge_filtered_segments<T: Into<Box<dyn Directory>>>(
316317 Ok ( merged_index)
317318}
318319
320+ struct Pools {
321+ pool : ThreadPool ,
322+ merge_thread_pool : ThreadPool ,
323+ merge_errors : Arc < RwLock < Vec < TantivyError > > > ,
324+ }
325+
319326pub ( crate ) struct InnerSegmentUpdater {
320327 // we keep a copy of the current active IndexMeta to
321328 // avoid loading the file every time we need it in the
@@ -324,10 +331,7 @@ pub(crate) struct InnerSegmentUpdater {
324331 // This should be up to date as all update happen through
325332 // the unique active `SegmentUpdater`.
326333 active_index_meta : RwLock < Arc < IndexMeta > > ,
327- pool : ThreadPool ,
328- merge_thread_pool : ThreadPool ,
329- merge_errors : Arc < RwLock < Vec < TantivyError > > > ,
330-
334+ pools : Option < Pools > ,
331335 index : Index ,
332336 segment_manager : SegmentManager ,
333337 merge_policy : RwLock < Arc < dyn MergePolicy > > ,
@@ -347,40 +351,56 @@ impl SegmentUpdater {
347351 ) -> crate :: Result < SegmentUpdater > {
348352 let segments = index. searchable_segment_metas ( ) ?;
349353 let segment_manager = SegmentManager :: from_segments ( segments, delete_cursor) ;
350- let mut builder = ThreadPoolBuilder :: new ( )
351- . thread_name ( |_| "segment_updater" . to_string ( ) )
352- . num_threads ( 1 ) ;
353-
354- if let Some ( panic_handler) = panic_handler. as_ref ( ) {
355- let panic_handler = panic_handler. clone ( ) ;
356- builder = builder. panic_handler ( move |any| {
357- panic_handler ( any) ;
358- } ) ;
359- }
360354
361- let pool = builder. build ( ) . map_err ( |_| {
362- crate :: TantivyError :: SystemError ( "Failed to spawn segment updater thread" . to_string ( ) )
363- } ) ?;
364- let mut builder = ThreadPoolBuilder :: new ( )
365- . thread_name ( |i| format ! ( "merge_thread_{i}" ) )
366- . num_threads ( num_merge_threads) ;
367- if let Some ( panic_handler) = panic_handler {
368- let panic_handler = panic_handler. clone ( ) ;
369- builder = builder. panic_handler ( move |any| {
370- panic_handler ( any) ;
371- } ) ;
372- }
373-
374- let merge_thread_pool = builder. build ( ) . map_err ( |_| {
375- crate :: TantivyError :: SystemError ( "Failed to spawn segment merging thread" . to_string ( ) )
376- } ) ?;
377355 let index_meta = index. load_metas ( ) ?;
378356 Ok ( SegmentUpdater {
379357 inner : Arc :: new ( InnerSegmentUpdater {
380358 active_index_meta : RwLock :: new ( Arc :: new ( index_meta) ) ,
381- pool,
382- merge_thread_pool,
383- merge_errors : Default :: default ( ) ,
359+ pools : ( num_merge_threads > 0 ) . then ( || {
360+ let mut builder = ThreadPoolBuilder :: new ( )
361+ . thread_name ( |_| "segment_updater" . to_string ( ) )
362+ . num_threads ( 1 ) ;
363+
364+ if let Some ( panic_handler) = panic_handler. as_ref ( ) {
365+ let panic_handler = panic_handler. clone ( ) ;
366+ builder = builder. panic_handler ( move |any| {
367+ panic_handler ( any) ;
368+ } ) ;
369+ }
370+
371+ let pool = builder
372+ . build ( )
373+ . map_err ( |_| {
374+ crate :: TantivyError :: SystemError (
375+ "Failed to spawn segment updater thread" . to_string ( ) ,
376+ )
377+ } )
378+ . unwrap ( ) ;
379+
380+ let mut builder = ThreadPoolBuilder :: new ( )
381+ . thread_name ( |i| format ! ( "merge_thread_{i}" ) )
382+ . num_threads ( num_merge_threads) ;
383+ if let Some ( panic_handler) = panic_handler {
384+ let panic_handler = panic_handler. clone ( ) ;
385+ builder = builder. panic_handler ( move |any| {
386+ panic_handler ( any) ;
387+ } ) ;
388+ }
389+ let merge_thread_pool = builder
390+ . build ( )
391+ . map_err ( |_| {
392+ crate :: TantivyError :: SystemError (
393+ "Failed to spawn segment merging thread" . to_string ( ) ,
394+ )
395+ } )
396+ . unwrap ( ) ;
397+
398+ Pools {
399+ pool,
400+ merge_thread_pool,
401+ merge_errors : Default :: default ( ) ,
402+ }
403+ } ) ,
384404 index,
385405 segment_manager,
386406 merge_policy : RwLock :: new ( Arc :: new ( DefaultMergePolicy :: default ( ) ) ) ,
@@ -393,12 +413,12 @@ impl SegmentUpdater {
393413 }
394414
395415 pub fn get_merge_policy ( & self ) -> Arc < dyn MergePolicy > {
396- self . merge_policy . read ( ) . unwrap ( ) . clone ( )
416+ self . merge_policy . read ( ) . clone ( )
397417 }
398418
399419 pub fn set_merge_policy ( & self , merge_policy : Box < dyn MergePolicy > ) {
400420 let arc_merge_policy = Arc :: from ( merge_policy) ;
401- * self . merge_policy . write ( ) . unwrap ( ) = arc_merge_policy;
421+ * self . merge_policy . write ( ) = arc_merge_policy;
402422 }
403423
404424 fn schedule_task < T : ' static + Send , F : FnOnce ( ) -> crate :: Result < T > + ' static + Send > (
@@ -411,10 +431,14 @@ impl SegmentUpdater {
411431 let ( scheduled_result, sender) = FutureResult :: create (
412432 "A segment_updater future did not succeed. This should never happen." ,
413433 ) ;
414- self . pool . spawn ( || {
415- let task_result = task ( ) ;
416- let _ = sender. send ( task_result) ;
417- } ) ;
434+ self . pools
435+ . as_ref ( )
436+ . expect ( "thread pools should have been configured" )
437+ . pool
438+ . spawn ( || {
439+ let task_result = task ( ) ;
440+ let _ = sender. send ( task_result) ;
441+ } ) ;
418442 scheduled_result
419443 }
420444
@@ -537,11 +561,11 @@ impl SegmentUpdater {
537561 }
538562
539563 fn store_meta ( & self , index_meta : & IndexMeta ) {
540- * self . active_index_meta . write ( ) . unwrap ( ) = Arc :: new ( index_meta. clone ( ) ) ;
564+ * self . active_index_meta . write ( ) = Arc :: new ( index_meta. clone ( ) ) ;
541565 }
542566
543567 fn load_meta ( & self ) -> Arc < IndexMeta > {
544- self . active_index_meta . read ( ) . unwrap ( ) . clone ( )
568+ self . active_index_meta . read ( ) . clone ( )
545569 }
546570
547571 pub ( crate ) fn make_merge_operation (
@@ -605,38 +629,48 @@ impl SegmentUpdater {
605629 FutureResult :: create ( "Merge operation failed." ) ;
606630
607631 let cancel = self . cancel . box_clone ( ) ;
608- let merge_errors = self . merge_errors . clone ( ) ;
609- self . merge_thread_pool . spawn ( move || {
610- // The fact that `merge_operation` is moved here is important.
611- // Its lifetime is used to track how many merging thread are currently running,
612- // as well as which segment is currently in merge and therefore should not be
613- // candidate for another merge.
614- match merge (
615- & segment_updater. index ,
616- segment_entries,
617- merge_operation. target_opstamp ( ) ,
618- cancel,
619- false ,
620- ) {
621- Ok ( after_merge_segment_entry) => {
622- let res = segment_updater. end_merge ( merge_operation, after_merge_segment_entry) ;
623- let _send_result = merging_future_send. send ( res) ;
624- }
625- Err ( merge_error) => {
626- warn ! (
627- "Merge of {:?} was cancelled: {:?}" ,
628- merge_operation. segment_ids( ) . to_vec( ) ,
629- merge_error
630- ) ;
631- if cfg ! ( test) {
632- panic ! ( "{merge_error:?}" ) ;
632+ let merge_errors = self
633+ . pools
634+ . as_ref ( )
635+ . expect ( "thread pools should have been configured" )
636+ . merge_errors
637+ . clone ( ) ;
638+ self . pools
639+ . as_ref ( )
640+ . expect ( "thread pools should have been configured" )
641+ . merge_thread_pool
642+ . spawn ( move || {
643+ // The fact that `merge_operation` is moved here is important.
644+ // Its lifetime is used to track how many merging thread are currently running,
645+ // as well as which segment is currently in merge and therefore should not be
646+ // candidate for another merge.
647+ match merge (
648+ & segment_updater. index ,
649+ segment_entries,
650+ merge_operation. target_opstamp ( ) ,
651+ cancel,
652+ false ,
653+ ) {
654+ Ok ( after_merge_segment_entry) => {
655+ let res =
656+ segment_updater. end_merge ( merge_operation, after_merge_segment_entry) ;
657+ let _send_result = merging_future_send. send ( res) ;
633658 }
659+ Err ( merge_error) => {
660+ warn ! (
661+ "Merge of {:?} was cancelled: {:?}" ,
662+ merge_operation. segment_ids( ) . to_vec( ) ,
663+ merge_error
664+ ) ;
665+ if cfg ! ( test) {
666+ panic ! ( "{merge_error:?}" ) ;
667+ }
634668
635- merge_errors. write ( ) . unwrap ( ) . push ( merge_error. clone ( ) ) ;
636- let _send_result = merging_future_send. send ( Err ( merge_error) ) ;
669+ merge_errors. write ( ) . push ( merge_error. clone ( ) ) ;
670+ let _send_result = merging_future_send. send ( Err ( merge_error) ) ;
671+ }
637672 }
638- }
639- } ) ;
673+ } ) ;
640674
641675 scheduled_result
642676 }
@@ -679,7 +713,11 @@ impl SegmentUpdater {
679713 }
680714
681715 pub ( crate ) fn get_merge_errors ( & self ) -> Vec < TantivyError > {
682- self . merge_errors . read ( ) . unwrap ( ) . clone ( )
716+ if let Some ( pools) = self . pools . as_ref ( ) {
717+ pools. merge_errors . read ( ) . clone ( )
718+ } else {
719+ Vec :: new ( )
720+ }
683721 }
684722
685723 fn consider_merge_options ( & self ) {
0 commit comments