From ab1dba9d87fbc2e1dd5668d1631c93120f59bdda Mon Sep 17 00:00:00 2001 From: "nadav.govari" Date: Tue, 25 Aug 2026 13:55:27 -0400 Subject: [PATCH] Make compaction planner configurable --- config/quickwit.yaml | 8 ++ docs/configuration/node-config.md | 20 +++++ .../src/planner/compaction_planner.rs | 80 +++++++++++-------- quickwit/quickwit-config/src/lib.rs | 8 +- .../quickwit-config/src/node_config/mod.rs | 62 ++++++++++++++ .../src/node_config/serialize.rs | 15 +++- quickwit/quickwit-serve/src/lib.rs | 12 ++- 7 files changed, 165 insertions(+), 40 deletions(-) diff --git a/config/quickwit.yaml b/config/quickwit.yaml index 58bdd4942ff..27ec3d1e90d 100644 --- a/config/quickwit.yaml +++ b/config/quickwit.yaml @@ -190,6 +190,14 @@ indexer: # compactor: # decommission_timeout: 300s # +# -------------------------- Compaction planner settings ---------------------------- +# https://quickwit.io/docs/configuration/node-config#compaction-planner-configuration +# +# compaction_planner: +# scan_page_size: 5000 +# scan_and_plan_interval: 5s +# max_excluded_split_ids: 50000 +# # -------------------------------- Searcher settings -------------------------------- # https://quickwit.io/docs/configuration/node-config#searcher-configuration # diff --git a/docs/configuration/node-config.md b/docs/configuration/node-config.md index cbe3f8b0022..c761781714e 100644 --- a/docs/configuration/node-config.md +++ b/docs/configuration/node-config.md @@ -297,6 +297,26 @@ compactor: decommission_timeout: 300s ``` +## Compaction planner configuration + +This section contains the configuration options for the compaction planner running on janitor +nodes when standalone compactors are enabled. + +| Property | Description | Default value | +| --- | --- | --- | +| `scan_page_size` | Maximum number of splits fetched from the metastore per scan. | `5000` | +| `scan_and_plan_interval` | Interval between compaction planner scan-and-plan cycles. | `5s` | +| `max_excluded_split_ids` | Maximum number of split IDs excluded from a metastore scan because they are already tracked. | `50000` | + +Example: + +```yaml +compaction_planner: + scan_page_size: 10000 + scan_and_plan_interval: 2s + max_excluded_split_ids: 100000 +``` + ## Searcher configuration This section contains the configuration options for a Searcher. diff --git a/quickwit/quickwit-compaction/src/planner/compaction_planner.rs b/quickwit/quickwit-compaction/src/planner/compaction_planner.rs index 0617dd0de97..736d9ea0084 100644 --- a/quickwit/quickwit-compaction/src/planner/compaction_planner.rs +++ b/quickwit/quickwit-compaction/src/planner/compaction_planner.rs @@ -40,29 +40,17 @@ use super::compaction_state::CompactionState; use super::index_config_metastore::{IndexConfigMetastore, IndexEntry}; use crate::planner::metrics::{METASTORE_ERRORS, NEW_SPLITS_SCANNED, OPERATION, SOURCE_UID}; -/// Cap on splits fetched per tick. Every tick, the planner re-scans the immature published set, -/// sorted by `maturity_timestamp` ASC so the most-urgent splits are processed first when a backlog -/// exists. Splits beyond this cap aren't lost -- they bubble into range as the front of the queue -/// is merged off. -const SCAN_PAGE_SIZE: usize = 5_000; - -/// Cap on the size of the `excluded_split_ids` list we send to the metastore. -/// It's a sanity max rather than some invariant. -const MAX_EXCLUDED_SPLIT_IDS: usize = 50_000; - #[derive(Debug)] pub struct CompactionPlanner { state: CompactionState, index_config_metastore: IndexConfigMetastore, metastore: MetastoreServiceClient, cluster: Cluster, + scan_page_size: usize, + scan_and_plan_interval: Duration, + max_excluded_split_ids: usize, } -const SCAN_AND_PLAN_INTERVAL: Duration = Duration::from_secs(5); -/// On initialization, we want to wait for two intervals to allow any in-progress workers to report -/// their progress, preventing us from frivolously rescheduling work. -const INITIAL_SCAN_AND_PLAN_INTERVAL: Duration = SCAN_AND_PLAN_INTERVAL.saturating_mul(2); - #[derive(Debug)] struct ScanAndPlan; @@ -80,11 +68,14 @@ impl Actor for CompactionPlanner { fn observable_state(&self) -> Self::ObservableState {} async fn initialize(&mut self, ctx: &ActorContext) -> Result<(), ActorExitStatus> { + // On initialization, wait for two intervals to allow any in-progress workers to report + // their progress, preventing us from frivolously rescheduling work. + let initial_scan_and_plan_interval = self.scan_and_plan_interval.saturating_mul(2); info!( "initializing compaction planner, waiting for indexers to hand off compaction in {}", - INITIAL_SCAN_AND_PLAN_INTERVAL.pretty_display() + initial_scan_and_plan_interval.pretty_display() ); - ctx.schedule_self_msg(INITIAL_SCAN_AND_PLAN_INTERVAL, AwaitIndexersMigrated); + ctx.schedule_self_msg(initial_scan_and_plan_interval, AwaitIndexersMigrated); Ok(()) } } @@ -102,7 +93,7 @@ impl Handler for CompactionPlanner { error!(%error, "failed to scan metastore and/or plan merges"); } self.state.check_heartbeat_timeouts(); - ctx.schedule_self_msg(SCAN_AND_PLAN_INTERVAL, ScanAndPlan); + ctx.schedule_self_msg(self.scan_and_plan_interval, ScanAndPlan); Ok(()) } } @@ -127,7 +118,7 @@ impl Handler for CompactionPlanner { "waiting for indexers to report standalone compactors enabled before planning \ merges" ); - ctx.schedule_self_msg(SCAN_AND_PLAN_INTERVAL, AwaitIndexersMigrated); + ctx.schedule_self_msg(self.scan_and_plan_interval, AwaitIndexersMigrated); } Ok(()) } @@ -152,12 +143,21 @@ impl Handler for CompactionPlanner { } impl CompactionPlanner { - pub fn new(metastore: MetastoreServiceClient, cluster: Cluster) -> Self { + pub fn new( + metastore: MetastoreServiceClient, + cluster: Cluster, + scan_page_size: usize, + scan_and_plan_interval: Duration, + max_excluded_split_ids: usize, + ) -> Self { CompactionPlanner { state: CompactionState::default(), index_config_metastore: IndexConfigMetastore::new(metastore.clone()), metastore, cluster, + scan_page_size, + scan_and_plan_interval, + max_excluded_split_ids, } } @@ -194,12 +194,12 @@ impl CompactionPlanner { } async fn scan_metastore(&self) -> Result> { - let excluded_split_ids = self.state.tracked_split_ids(MAX_EXCLUDED_SPLIT_IDS); + let excluded_split_ids = self.state.tracked_split_ids(self.max_excluded_split_ids); let query = ListSplitsQuery::for_all_indexes() .with_split_state(SplitState::Published) .retain_immature(OffsetDateTime::now_utc()) .sort_by_maturity_timestamp() - .with_limit(SCAN_PAGE_SIZE) + .with_limit(self.scan_page_size) .with_excluded_split_ids(excluded_split_ids); let request = ListSplitsRequest::try_from_list_splits_query(&query)?; let splits = self @@ -315,10 +315,10 @@ mod tests { use quickwit_cluster::{ChitchatTransport, create_cluster_for_test}; use quickwit_common::ServiceStream; use quickwit_common::test_utils::wait_until_predicate; - use quickwit_config::IndexingSettings; use quickwit_config::merge_policy_config::{ ConstWriteAmplificationMergePolicyConfig, MergePolicyConfig, }; + use quickwit_config::{CompactionPlannerConfig, IndexingSettings}; use quickwit_metastore::{ IndexMetadata, IndexMetadataResponseExt, ListSplitsRequestExt, ListSplitsResponseExt, SortBy, Split, SplitMaturity, SplitMetadata, SplitState, @@ -385,6 +385,17 @@ mod tests { .unwrap() } + fn planner_for_test(metastore: MetastoreServiceClient, cluster: Cluster) -> CompactionPlanner { + let config = CompactionPlannerConfig::default(); + CompactionPlanner::new( + metastore, + cluster, + config.scan_page_size(), + config.scan_and_plan_interval(), + config.max_excluded_split_ids(), + ) + } + #[tokio::test] async fn test_scan_metastore_query_shape_and_passthrough() { let index_uid = IndexUid::for_test("test-index", 0); @@ -400,7 +411,10 @@ mod tests { let query = req.deserialize_list_splits_query().unwrap(); assert_eq!(query.split_states, vec![SplitState::Published]); - assert_eq!(query.limit, Some(SCAN_PAGE_SIZE)); + assert_eq!( + query.limit, + Some(CompactionPlannerConfig::default().scan_page_size()) + ); assert_eq!(query.sort_by, SortBy::MaturityTimestamp); let Bound::Excluded(mature_at) = query.mature else { @@ -421,7 +435,7 @@ mod tests { Ok(ServiceStream::from(vec![Ok(response)])) }); - let planner = CompactionPlanner::new( + let planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -448,7 +462,7 @@ mod tests { Ok(ServiceStream::from(vec![Ok(response)])) }); - let mut planner = CompactionPlanner::new( + let mut planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -480,7 +494,7 @@ mod tests { mock.expect_index_metadata() .returning(move |_| Ok(response.clone())); - let mut planner = CompactionPlanner::new( + let mut planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -518,7 +532,7 @@ mod tests { }) }); - let mut planner = CompactionPlanner::new( + let mut planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -538,7 +552,7 @@ mod tests { }) }); - let mut planner = CompactionPlanner::new( + let mut planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -567,7 +581,7 @@ mod tests { mock.expect_index_metadata() .returning(move |_| Ok(index_metadata_response.clone())); - let mut planner = CompactionPlanner::new( + let mut planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -601,7 +615,7 @@ mod tests { mock.expect_index_metadata() .returning(move |_| Ok(index_metadata_response.clone())); - let mut planner = CompactionPlanner::new( + let mut planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -646,7 +660,7 @@ mod tests { )])) }); - let mut planner = CompactionPlanner::new( + let mut planner = planner_for_test( MetastoreServiceClient::from_mock(mock), test_cluster().await, ); @@ -763,7 +777,7 @@ mod tests { .await .unwrap(); let seeds = vec![janitor_cluster.gossip_listen_addr.to_string()]; - let planner = CompactionPlanner::new( + let planner = planner_for_test( MetastoreServiceClient::from_mock(MockMetastoreService::new()), janitor_cluster.clone(), ); diff --git a/quickwit/quickwit-config/src/lib.rs b/quickwit/quickwit-config/src/lib.rs index 904e65f8180..4e3b79c8ac6 100644 --- a/quickwit/quickwit-config/src/lib.rs +++ b/quickwit/quickwit-config/src/lib.rs @@ -77,10 +77,10 @@ pub use crate::metastore_config::{ MetastoreBackend, MetastoreConfig, MetastoreConfigs, PostgresMetastoreConfig, }; pub use crate::node_config::{ - CacheConfig, CachePolicy, CompactorConfig, DEFAULT_QW_CONFIG_PATH, GrpcConfig, HealthConfig, - IndexerConfig, IngestApiConfig, JaegerConfig, KeepAliveConfig, LambdaConfig, - LambdaDeployConfig, MAX_GOSSIP_PROTOCOL_VERSION, NodeConfig, RestConfig, SearcherConfig, - SplitCacheLimits, StorageTimeoutPolicy, TlsConfig, + CacheConfig, CachePolicy, CompactionPlannerConfig, CompactorConfig, DEFAULT_QW_CONFIG_PATH, + GrpcConfig, HealthConfig, IndexerConfig, IngestApiConfig, JaegerConfig, KeepAliveConfig, + LambdaConfig, LambdaDeployConfig, MAX_GOSSIP_PROTOCOL_VERSION, NodeConfig, RestConfig, + SearcherConfig, SplitCacheLimits, StorageTimeoutPolicy, TlsConfig, }; pub use crate::serde_utils::HumanDuration; use crate::source_config::serialize::{SourceConfigV0_7, SourceConfigV0_8, VersionedSourceConfig}; diff --git a/quickwit/quickwit-config/src/node_config/mod.rs b/quickwit/quickwit-config/src/node_config/mod.rs index 5e02be1af39..c1c68ad4a54 100644 --- a/quickwit/quickwit-config/src/node_config/mod.rs +++ b/quickwit/quickwit-config/src/node_config/mod.rs @@ -377,6 +377,43 @@ impl Default for CompactorConfig { } } +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields, default)] +pub struct CompactionPlannerConfig { + /// Maximum number of splits fetched from the metastore per scan. + scan_page_size: usize, + /// Interval between compaction planner scan-and-plan cycles. + scan_and_plan_interval: HumanDuration, + /// Maximum number of split IDs excluded from a metastore scan because they are already + /// tracked. + max_excluded_split_ids: usize, +} + +impl CompactionPlannerConfig { + pub fn scan_page_size(&self) -> usize { + self.scan_page_size + } + + pub fn scan_and_plan_interval(&self) -> Duration { + Duration::from(self.scan_and_plan_interval.clone()) + } + + pub fn max_excluded_split_ids(&self) -> usize { + self.max_excluded_split_ids + } +} + +impl Default for CompactionPlannerConfig { + fn default() -> Self { + Self { + scan_page_size: 5_000, + scan_and_plan_interval: HumanDuration::try_from("5s".to_string()) + .expect("`5s` should be a valid human duration"), + max_excluded_split_ids: 50_000, + } + } +} + #[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct SplitCacheLimits { @@ -943,6 +980,7 @@ pub struct NodeConfig { pub ingest_api_config: IngestApiConfig, pub jaeger_config: JaegerConfig, pub compactor_config: CompactorConfig, + pub compaction_planner_config: CompactionPlannerConfig, #[serde(skip_serializing)] pub enable_standalone_compactors: bool, #[serde(skip_serializing_if = "Option::is_none")] @@ -1197,6 +1235,30 @@ mod tests { assert_eq!(yaml_config.decommission_timeout(), Duration::from_mins(2)); } + #[test] + fn test_compaction_planner_config() { + let default_config: CompactionPlannerConfig = serde_yaml::from_str("").unwrap(); + assert_eq!(default_config, CompactionPlannerConfig::default()); + assert_eq!(default_config.scan_page_size(), 5_000); + assert_eq!( + default_config.scan_and_plan_interval(), + Duration::from_secs(5) + ); + assert_eq!(default_config.max_excluded_split_ids(), 50_000); + + let yaml_config: CompactionPlannerConfig = serde_yaml::from_str( + r#" + scan_page_size: 10000 + scan_and_plan_interval: 2s + max_excluded_split_ids: 100000 + "#, + ) + .unwrap(); + assert_eq!(yaml_config.scan_page_size(), 10_000); + assert_eq!(yaml_config.scan_and_plan_interval(), Duration::from_secs(2)); + assert_eq!(yaml_config.max_excluded_split_ids(), 100_000); + } + #[track_caller] fn test_keepalive_config_serialization_aux( keep_alive_json: serde_json::Value, diff --git a/quickwit/quickwit-config/src/node_config/serialize.rs b/quickwit/quickwit-config/src/node_config/serialize.rs index b69c354ed58..26bdd4cbb05 100644 --- a/quickwit/quickwit-config/src/node_config/serialize.rs +++ b/quickwit/quickwit-config/src/node_config/serialize.rs @@ -37,8 +37,9 @@ use crate::service::QuickwitService; use crate::storage_config::StorageConfigs; use crate::templating::render_config; use crate::{ - CompactorConfig, ConfigFormat, IndexerConfig, IngestApiConfig, JaegerConfig, MetastoreConfigs, - NodeConfig, SearcherConfig, TlsConfig, validate_identifier, validate_node_id, + CompactionPlannerConfig, CompactorConfig, ConfigFormat, IndexerConfig, IngestApiConfig, + JaegerConfig, MetastoreConfigs, NodeConfig, SearcherConfig, TlsConfig, validate_identifier, + validate_node_id, }; pub const DEFAULT_CLUSTER_ID: &str = "quickwit-default-cluster"; @@ -240,6 +241,9 @@ struct NodeConfigBuilder { #[serde(rename = "compactor")] #[serde(default)] compactor_config: CompactorConfig, + #[serde(rename = "compaction_planner")] + #[serde(default)] + compaction_planner_config: CompactionPlannerConfig, #[serde(rename = "docs_clustering")] #[serde(default)] docs_clustering_config: Option, @@ -388,6 +392,7 @@ impl NodeConfigBuilder { ingest_api_config: self.ingest_api_config, jaeger_config: self.jaeger_config, compactor_config: self.compactor_config, + compaction_planner_config: self.compaction_planner_config, enable_standalone_compactors, docs_clustering_config, }; @@ -534,6 +539,7 @@ impl Default for NodeConfigBuilder { ingest_api_config: IngestApiConfig::default(), jaeger_config: JaegerConfig::default(), compactor_config: CompactorConfig::default(), + compaction_planner_config: CompactionPlannerConfig::default(), docs_clustering_config: None, } } @@ -685,6 +691,7 @@ pub fn node_config_for_tests_from_ports( ingest_api_config: IngestApiConfig::default(), jaeger_config: JaegerConfig::default(), compactor_config: CompactorConfig::default(), + compaction_planner_config: CompactionPlannerConfig::default(), enable_standalone_compactors: false, docs_clustering_config: None, } @@ -1192,6 +1199,10 @@ mod tests { assert_eq!(config.searcher_config, SearcherConfig::default()); assert_eq!(config.ingest_api_config, IngestApiConfig::default()); assert_eq!(config.jaeger_config, JaegerConfig::default()); + assert_eq!( + config.compaction_planner_config, + CompactionPlannerConfig::default() + ); } #[tokio::test] diff --git a/quickwit/quickwit-serve/src/lib.rs b/quickwit/quickwit-serve/src/lib.rs index 1bd34f3f4ac..6297f199f53 100644 --- a/quickwit/quickwit-serve/src/lib.rs +++ b/quickwit/quickwit-serve/src/lib.rs @@ -293,7 +293,17 @@ async fn get_compaction_planner_client_if_needed( return Ok((None, None)); } if is_janitor { - let planner = CompactionPlanner::new(metastore_client.clone(), cluster.clone()); + let planner = CompactionPlanner::new( + metastore_client.clone(), + cluster.clone(), + node_config.compaction_planner_config.scan_page_size(), + node_config + .compaction_planner_config + .scan_and_plan_interval(), + node_config + .compaction_planner_config + .max_excluded_split_ids(), + ); let (mailbox, handle) = universe.spawn_builder().spawn(planner); info!("compaction planner actor started on janitor node"); return Ok((