Skip to content

Commit 7054e1d

Browse files
authored
Rework merge task for both compactor and indexer merge flows (#6464)
1 parent ce46988 commit 7054e1d

10 files changed

Lines changed: 117 additions & 113 deletions

File tree

quickwit/quickwit-compaction/src/compaction_pipeline.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ use quickwit_doc_mapper::DocMapper;
2525
use quickwit_indexing::actors::{
2626
MergeExecutor, MergeSplitDownloader, Packager, Publisher, Uploader, UploaderType,
2727
};
28-
use quickwit_indexing::merge_policy::MergeOperation;
28+
use quickwit_indexing::merge_policy::{MergeOperation, MergeSource};
2929
use quickwit_indexing::{IndexingSplitStore, SplitsUpdateMailbox};
3030
use quickwit_metrics::{counter, gauge, histogram, label_values};
3131
use quickwit_proto::indexing::MergePipelineId;
@@ -351,7 +351,7 @@ impl CompactionPipeline {
351351
self.pipeline_start = Some(now);
352352
// Kick off the pipeline.
353353
merge_split_downloader_mailbox
354-
.try_send_message(self.merge_operation.clone())
354+
.try_send_message(MergeSource::Operation(self.merge_operation.clone()))
355355
.map_err(|err| {
356356
anyhow::anyhow!("failed to send merge operation to downloader: {err:?}")
357357
})?;

quickwit/quickwit-indexing/failpoints/mod.rs

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ use quickwit_common::rand::append_random_suffix;
4040
use quickwit_common::split_file;
4141
use quickwit_common::temp_dir::TempDirectory;
4242
use quickwit_indexing::actors::MergeExecutor;
43-
use quickwit_indexing::merge_policy::{MergeOperation, MergeTask};
43+
use quickwit_indexing::merge_policy::{MergeOperation, MergeSource, MergeTask};
4444
use quickwit_indexing::models::MergeScratch;
4545
use quickwit_indexing::{TestSandbox, get_tantivy_directory_from_split_bundle};
4646
use quickwit_metastore::{
@@ -285,10 +285,9 @@ async fn test_merge_executor_controlled_directory_kill_switch() -> anyhow::Resul
285285
tantivy_dirs.push(get_tantivy_directory_from_split_bundle(&dest_filepath).unwrap());
286286
}
287287
let merge_operation = MergeOperation::new_merge_operation(split_metadatas);
288-
let merge_task = MergeTask::from_merge_operation_for_test(merge_operation.clone());
288+
let merge_task = MergeTask::from_merge_operation_for_test(merge_operation);
289289
let merge_scratch = MergeScratch {
290-
merge_operation,
291-
merge_task: Some(merge_task),
290+
merge_source: MergeSource::Task(merge_task),
292291
merge_scratch_directory,
293292
downloaded_splits_directory,
294293
tantivy_dirs,

quickwit/quickwit-indexing/src/actors/indexing_service.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1398,9 +1398,9 @@ mod tests {
13981398
// change whenever IndexingSettings fields are added/removed. Recompute
13991399
// by temporarily adding a test that prints
14001400
// `indexing_pipeline_params_fingerprint(&index_config, &source_config)`.
1401-
const PARAMS_FINGERPRINT_INGEST_API: u64 = 1637744865450232394;
1402-
const PARAMS_FINGERPRINT_SOURCE_1: u64 = 1705211905504908791;
1403-
const PARAMS_FINGERPRINT_SOURCE_2: u64 = 8706667372658059428;
1401+
const PARAMS_FINGERPRINT_INGEST_API: u64 = 7973087274884969148;
1402+
const PARAMS_FINGERPRINT_SOURCE_1: u64 = 9420938500552890840;
1403+
const PARAMS_FINGERPRINT_SOURCE_2: u64 = 16199199787360162635;
14041404

14051405
quickwit_common::setup_logging_for_tests();
14061406
let transport = ChannelTransport::default();

quickwit/quickwit-indexing/src/actors/merge_executor.rs

Lines changed: 16 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@ use tracing::{debug, error, info, instrument, warn};
4747

4848
use crate::actors::Packager;
4949
use crate::controlled_directory::ControlledDirectory;
50-
use crate::merge_policy::MergeOperationType;
50+
use crate::merge_policy::{MergeOperationType, MergeSource};
5151
use crate::models::{IndexedSplit, IndexedSplitBatch, MergeScratch, PublishLock, SplitAttrs};
5252

5353
#[derive(Clone)]
@@ -85,20 +85,20 @@ impl Actor for MergeExecutor {
8585
impl Handler<MergeScratch> for MergeExecutor {
8686
type Reply = ();
8787

88-
#[instrument(level = "info", name = "merge_executor", parent = merge_scratch.merge_operation.merge_parent_span.id(), skip_all)]
88+
#[instrument(level = "info", name = "merge_executor", parent = merge_scratch.merge_source.as_operation().merge_parent_span.id(), skip_all)]
8989
async fn handle(
9090
&mut self,
9191
merge_scratch: MergeScratch,
9292
ctx: &ActorContext<Self>,
9393
) -> Result<(), ActorExitStatus> {
9494
let start = Instant::now();
9595
let MergeScratch {
96-
merge_operation,
97-
merge_task,
96+
merge_source,
9897
tantivy_dirs,
9998
merge_scratch_directory,
10099
..
101100
} = merge_scratch;
101+
let merge_operation = merge_source.as_operation();
102102
// On nodes running the split compaction architecture, merge pipelines are ephemeral, and we
103103
// need to make sure there aren't too many CPU-bound operations occurring concurrently.
104104
let _cpu_permit = match &self.merge_execution_semaphore {
@@ -164,15 +164,20 @@ impl Handler<MergeScratch> for MergeExecutor {
164164
operation_type = %merge_operation.operation_type,
165165
"merge-operation-success"
166166
);
167+
let batch_parent_span = merge_operation.merge_parent_span.clone();
168+
let merge_task_opt = match merge_source {
169+
MergeSource::Task(task) => Some(task),
170+
MergeSource::Operation(_) => None,
171+
};
167172
ctx.send_message(
168173
&self.merge_packager_mailbox,
169174
IndexedSplitBatch {
170175
splits: vec![indexed_split],
171176
checkpoint_delta_opt: Default::default(),
172177
publish_lock: PublishLock::default(),
173178
publish_token_opt: None,
174-
batch_parent_span: merge_operation.merge_parent_span.clone(),
175-
merge_task_opt: merge_task,
179+
batch_parent_span,
180+
merge_task_opt,
176181
},
177182
)
178183
.await?;
@@ -615,7 +620,7 @@ mod tests {
615620
use tantivy::{Document, ReloadPolicy, TantivyDocument};
616621

617622
use super::*;
618-
use crate::merge_policy::{MergeOperation, MergeTask};
623+
use crate::merge_policy::{MergeOperation, MergeSource, MergeTask};
619624
use crate::{TestSandbox, get_tantivy_directory_from_split_bundle, new_split_id};
620625

621626
#[tokio::test]
@@ -664,10 +669,9 @@ mod tests {
664669
tantivy_dirs.push(get_tantivy_directory_from_split_bundle(&dest_filepath).unwrap())
665670
}
666671
let merge_operation = MergeOperation::new_merge_operation(split_metas);
667-
let merge_task = MergeTask::from_merge_operation_for_test(merge_operation.clone());
672+
let merge_task = MergeTask::from_merge_operation_for_test(merge_operation);
668673
let merge_scratch = MergeScratch {
669-
merge_operation,
670-
merge_task: Some(merge_task),
674+
merge_source: MergeSource::Task(merge_task),
671675
tantivy_dirs,
672676
merge_scratch_directory,
673677
downloaded_splits_directory,
@@ -810,10 +814,9 @@ mod tests {
810814
.await?;
811815
let tantivy_dir = get_tantivy_directory_from_split_bundle(&dest_filepath).unwrap();
812816
let merge_operation = MergeOperation::new_delete_and_merge_operation(new_split_metadata);
813-
let merge_task = MergeTask::from_merge_operation_for_test(merge_operation.clone());
817+
let merge_task = MergeTask::from_merge_operation_for_test(merge_operation);
814818
let merge_scratch = MergeScratch {
815-
merge_operation,
816-
merge_task: Some(merge_task),
819+
merge_source: MergeSource::Task(merge_task),
817820
tantivy_dirs: vec![tantivy_dir],
818821
merge_scratch_directory,
819822
downloaded_splits_directory,

quickwit/quickwit-indexing/src/actors/merge_planner.rs

Lines changed: 18 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -376,7 +376,7 @@ mod tests {
376376

377377
use crate::actors::MergePlanner;
378378
use crate::merge_policy::{
379-
MergePolicy, MergeTask, StableLogMergePolicy, merge_policy_from_settings,
379+
MergePolicy, MergeSource, StableLogMergePolicy, merge_policy_from_settings,
380380
};
381381
use crate::models::NewSplits;
382382

@@ -481,36 +481,40 @@ mod tests {
481481
};
482482
merge_planner_mailbox.send_message(message).await?;
483483
merge_planner_handle.process_pending_and_observe().await;
484-
let operations = merge_split_downloader_inbox.drain_for_test_typed::<MergeTask>();
484+
let operations = merge_split_downloader_inbox.drain_for_test_typed::<MergeSource>();
485485
assert_eq!(operations.len(), 3);
486-
let mut merge_operations = operations
487-
.into_iter()
488-
.sorted_by_key(|op| (op.splits[0].partition_id, op.splits[0].doc_mapping_uid));
486+
let mut merge_operations = operations.into_iter().sorted_by_key(|op| {
487+
let op = op.as_operation();
488+
(op.splits[0].partition_id, op.splits[0].doc_mapping_uid)
489+
});
489490

490491
let first_merge_operation = merge_operations.next().unwrap();
491-
assert_eq!(first_merge_operation.splits.len(), 4);
492+
assert_eq!(first_merge_operation.as_operation().splits.len(), 4);
492493
assert!(
493494
first_merge_operation
495+
.as_operation()
494496
.splits
495497
.iter()
496498
.all(|split| split.partition_id == 1
497499
&& split.doc_mapping_uid == doc_mapping_uid1)
498500
);
499501

500502
let second_merge_operation = merge_operations.next().unwrap();
501-
assert_eq!(second_merge_operation.splits.len(), 3);
503+
assert_eq!(second_merge_operation.as_operation().splits.len(), 3);
502504
assert!(
503505
second_merge_operation
506+
.as_operation()
504507
.splits
505508
.iter()
506509
.all(|split| split.partition_id == 1
507510
&& split.doc_mapping_uid == doc_mapping_uid2)
508511
);
509512

510513
let third_merge_operation = merge_operations.next().unwrap();
511-
assert_eq!(third_merge_operation.splits.len(), 3);
514+
assert_eq!(third_merge_operation.as_operation().splits.len(), 3);
512515
assert!(
513516
third_merge_operation
517+
.as_operation()
514518
.splits
515519
.iter()
516520
.all(|split| split.partition_id == 2
@@ -580,7 +584,7 @@ mod tests {
580584
// We wait for the first merge ops. If we sent the Quit message right away, it would have
581585
// been queue before first `PlanMerge` message.
582586
let merge_task_res = merge_split_downloader_inbox
583-
.recv_typed_message::<MergeTask>()
587+
.recv_typed_message::<MergeSource>()
584588
.await;
585589
assert!(merge_task_res.is_ok());
586590

@@ -594,15 +598,15 @@ mod tests {
594598

595599
let _ = merge_planner_handle.process_pending_and_observe().await;
596600

597-
let merge_ops = merge_split_downloader_inbox.drain_for_test_typed::<MergeTask>();
601+
let merge_ops = merge_split_downloader_inbox.drain_for_test_typed::<MergeSource>();
598602

599603
assert!(merge_ops.is_empty());
600604

601605
merge_planner_mailbox.send_message(Command::Quit).await?;
602606

603607
let (exit_status, _last_state) = merge_planner_handle.join().await;
604608
assert!(matches!(exit_status, ActorExitStatus::Quit));
605-
let merge_ops = merge_split_downloader_inbox.drain_for_test_typed::<MergeTask>();
609+
let merge_ops = merge_split_downloader_inbox.drain_for_test_typed::<MergeSource>();
606610
assert!(merge_ops.is_empty());
607611
universe.assert_quit().await;
608612
Ok(())
@@ -672,7 +676,7 @@ mod tests {
672676
merge_planner_mailbox.send_message(Command::Quit).await?;
673677
let (exit_status, _last_state) = merge_planner_handle.join().await;
674678
assert!(matches!(exit_status, ActorExitStatus::Quit));
675-
let merge_tasks = merge_split_downloader_inbox.drain_for_test_typed::<MergeTask>();
679+
let merge_tasks = merge_split_downloader_inbox.drain_for_test_typed::<MergeSource>();
676680

677681
assert!(merge_tasks.is_empty());
678682
universe.assert_quit().await;
@@ -750,7 +754,7 @@ mod tests {
750754

751755
// Instead, we wait for the first merge ops.
752756
let merge_task_res = merge_split_downloader_inbox
753-
.recv_typed_message::<MergeTask>()
757+
.recv_typed_message::<MergeSource>()
754758
.await;
755759
assert!(merge_task_res.is_ok());
756760

@@ -759,7 +763,7 @@ mod tests {
759763
let (exit_status, _last_state) = merge_planner_handle.join().await;
760764

761765
assert!(matches!(exit_status, ActorExitStatus::Quit));
762-
let merge_tasks = merge_split_downloader_inbox.drain_for_test_typed::<MergeTask>();
766+
let merge_tasks = merge_split_downloader_inbox.drain_for_test_typed::<MergeSource>();
763767
assert!(merge_tasks.is_empty());
764768

765769
universe.assert_quit().await;

quickwit/quickwit-indexing/src/actors/merge_scheduler_service.rs

Lines changed: 19 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@ use tracing::error;
2929
use super::MergeSplitDownloader;
3030
#[cfg(feature = "metrics")]
3131
use super::parquet_pipeline::{ParquetMergeSplitDownloader, ParquetMergeTask};
32-
use crate::merge_policy::{MergeOperation, MergeTask, compute_merge_score};
32+
use crate::merge_policy::{MergeOperation, MergeSource, MergeTask, compute_merge_score};
3333
use crate::metrics::{ONGOING_MERGE_OPERATIONS, PENDING_MERGE_BYTES, PENDING_MERGE_OPERATIONS};
3434

3535
pub struct MergePermit {
@@ -229,7 +229,7 @@ impl MergeSchedulerService {
229229
self.pending_merge_bytes -= merge_task.merge_operation.total_num_bytes();
230230
PENDING_MERGE_OPERATIONS.set(self.pending_merge_queue.len() as f64);
231231
PENDING_MERGE_BYTES.set(self.pending_merge_bytes as f64);
232-
match split_downloader_mailbox.try_send_message(merge_task) {
232+
match split_downloader_mailbox.try_send_message(MergeSource::Task(merge_task)) {
233233
Ok(_) => {}
234234
Err(quickwit_actors::TrySendError::Full(_)) => {
235235
// The split downloader mailbox has an unbounded queue capacity,
@@ -452,7 +452,7 @@ mod tests {
452452
use tokio::time::timeout;
453453

454454
use super::*;
455-
use crate::merge_policy::{MergeOperation, MergeTask};
455+
use crate::merge_policy::{MergeOperation, MergeSource};
456456

457457
fn build_merge_operation(num_splits: usize, num_bytes_per_split: u64) -> MergeOperation {
458458
let splits: Vec<SplitMetadata> = std::iter::repeat_with(|| SplitMetadata {
@@ -530,58 +530,58 @@ mod tests {
530530
.unwrap();
531531
}
532532
{
533-
let merge_task: MergeTask = merge_split_downloader_inbox
534-
.recv_typed_message::<MergeTask>()
533+
let merge_source = merge_split_downloader_inbox
534+
.recv_typed_message::<MergeSource>()
535535
.await
536536
.unwrap();
537537
assert_eq!(
538-
merge_task.merge_operation.splits[0].footer_offsets.end,
538+
merge_source.as_operation().splits[0].footer_offsets.end,
539539
4_000_000
540540
);
541-
let merge_task2: MergeTask = merge_split_downloader_inbox
542-
.recv_typed_message::<MergeTask>()
541+
let merge_source2 = merge_split_downloader_inbox
542+
.recv_typed_message::<MergeSource>()
543543
.await
544544
.unwrap();
545545
assert_eq!(
546-
merge_task2.merge_operation.splits[0].footer_offsets.end,
546+
merge_source2.as_operation().splits[0].footer_offsets.end,
547547
3_000_000
548548
);
549549
assert!(
550550
timeout(
551551
Duration::from_millis(200),
552-
merge_split_downloader_inbox.recv_typed_message::<MergeTask>()
552+
merge_split_downloader_inbox.recv_typed_message::<MergeSource>()
553553
)
554554
.await
555555
.is_err()
556556
);
557557
}
558558
{
559-
let merge_task: MergeTask = merge_split_downloader_inbox
560-
.recv_typed_message::<MergeTask>()
559+
let merge_source = merge_split_downloader_inbox
560+
.recv_typed_message::<MergeSource>()
561561
.await
562562
.unwrap();
563563
assert_eq!(
564-
merge_task.merge_operation.splits[0].footer_offsets.end,
564+
merge_source.as_operation().splits[0].footer_offsets.end,
565565
1_000_000
566566
);
567567
}
568568
{
569-
let merge_task: MergeTask = merge_split_downloader_inbox
570-
.recv_typed_message::<MergeTask>()
569+
let merge_source = merge_split_downloader_inbox
570+
.recv_typed_message::<MergeSource>()
571571
.await
572572
.unwrap();
573573
assert_eq!(
574-
merge_task.merge_operation.splits[0].footer_offsets.end,
574+
merge_source.as_operation().splits[0].footer_offsets.end,
575575
2_000_000
576576
);
577577
}
578578
{
579-
let merge_task: MergeTask = merge_split_downloader_inbox
580-
.recv_typed_message::<MergeTask>()
579+
let merge_source = merge_split_downloader_inbox
580+
.recv_typed_message::<MergeSource>()
581581
.await
582582
.unwrap();
583583
assert_eq!(
584-
merge_task.merge_operation.splits[0].footer_offsets.end,
584+
merge_source.as_operation().splits[0].footer_offsets.end,
585585
5_000_000
586586
);
587587
}

0 commit comments

Comments
 (0)