From 5e7cddc81d07d3a8a7e5fe19a420477e3e444afb Mon Sep 17 00:00:00 2001 From: "nadav.govari" Date: Mon, 5 Oct 2026 10:09:57 -0400 Subject: [PATCH 1/2] Reconcile shard scaling periodically with the v2 controller --- .../src/control_plane.rs | 77 ++- .../src/ingest/legacy_scaling_controller.rs | 254 ++------ .../quickwit-control-plane/src/ingest/mod.rs | 2 + .../src/ingest/scaling_controller.rs | 552 ++++++++++++++++++ .../quickwit-control-plane/src/model/mod.rs | 13 +- .../src/model/shard_table.rs | 120 ++++ 6 files changed, 770 insertions(+), 248 deletions(-) create mode 100644 quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs diff --git a/quickwit/quickwit-control-plane/src/control_plane.rs b/quickwit/quickwit-control-plane/src/control_plane.rs index 4dd0ff6c3e8..1230fcd7fb6 100644 --- a/quickwit/quickwit-control-plane/src/control_plane.rs +++ b/quickwit/quickwit-control-plane/src/control_plane.rs @@ -60,7 +60,7 @@ use crate::cooldown_map::{CooldownMap, CooldownStatus}; use crate::debouncer::Debouncer; use crate::indexing_scheduler::{IndexingScheduler, IndexingSchedulerState}; use crate::ingest::ingest_controller::{IngestControllerStats, RebalanceShardsCallback}; -use crate::ingest::{IngestController, LegacyScalingController}; +use crate::ingest::{IngestController, LegacyScalingController, ScalingController}; use crate::metrics::{METASTORE_ERROR_ABORTED, METASTORE_ERROR_MAYBE_EXECUTED, RESTART_TOTAL}; use crate::model::ControlPlaneModel; @@ -99,6 +99,7 @@ pub struct ControlPlane { indexing_scheduler: IndexingScheduler, ingest_controller: IngestController, legacy_scaling_controller: LegacyScalingController, + scaling_controller: ScalingController, metastore: MetastoreServiceClient, model: ControlPlaneModel, prune_shard_cooldown: CooldownMap<(IndexId, SourceId)>, @@ -138,6 +139,9 @@ impl ControlPlane { cluster_config.shard_throughput_limit, cluster_config.shard_scale_up_factor, ); + let scaling_controller = ScalingController::with_shard_throughput_limit( + cluster_config.shard_throughput_limit, + ); let readiness_tx = readiness_tx.clone(); let _ = readiness_tx.send(false); @@ -147,6 +151,7 @@ impl ControlPlane { indexing_scheduler, ingest_controller, legacy_scaling_controller, + scaling_controller, metastore: metastore.clone(), model: Default::default(), prune_shard_cooldown: CooldownMap::new(NonZeroUsize::new(1024).unwrap()), @@ -491,11 +496,26 @@ impl Handler for ControlPlane { _message: ControlPlaneLoop, ctx: &ActorContext, ) -> Result<(), ActorExitStatus> { - if let Err(metastore_error) = self - .ingest_controller - .rebalance_shards(&mut self.model, ctx.mailbox(), ctx.progress()) - .await - { + let reconcile_shards_result = if self.ingest_controller.all_indexers_migrated() { + self.scaling_controller + .reconcile_shards( + &mut self.ingest_controller, + &mut self.model, + ctx.mailbox(), + ctx.progress(), + ) + .await + } else { + self.legacy_scaling_controller + .reconcile_shards( + &mut self.ingest_controller, + &mut self.model, + ctx.mailbox(), + ctx.progress(), + ) + .await + }; + if let Err(metastore_error) = reconcile_shards_result { if let Err(actor_exit_status) = convert_metastore_error::<()>(metastore_error) { // See convert_metastore_error's spec. If it returns an error, it // means we do not know if all metastore tx were aborted or not. @@ -507,6 +527,7 @@ impl Handler for ControlPlane { return Err(actor_exit_status); } } + let _rebuild_plan_waiter = self.rebuild_plan_debounced(ctx); self.indexing_scheduler.control_running_plan(&self.model); ctx.schedule_self_msg(CONTROL_PLAN_LOOP_INTERVAL, ControlPlaneLoop); Ok(()) @@ -983,7 +1004,7 @@ impl DeferableReplyHandler for ControlPlane { &mut self, request: ReportIndexerStateRequest, reply: impl FnOnce(Self::Reply) + Send + Sync + 'static, - ctx: &ActorContext, + _ctx: &ActorContext, ) -> Result<(), ActorExitStatus> { reply(Ok(ReportIndexerStateResponse {})); @@ -994,23 +1015,15 @@ impl DeferableReplyHandler for ControlPlane { indexing_tasks_update.indexing_tasks, ); } - if let Some(shards_update) = request.shards_update - && let Err(metastore_error) = self - .legacy_scaling_controller - .handle_shards_update( - &mut self.ingest_controller, - &request.node_id, - request.generation_id, - shards_update, - &mut self.model, - ctx.progress(), - ) - .await - { - // Return () if there's no metastore error; return the error if there is one. - return convert_metastore_error::<()>(metastore_error).map(|_| ()); + if let Some(shards_update) = request.shards_update { + self.scaling_controller.handle_shards_update( + &self.ingest_controller, + &request.node_id, + request.generation_id, + shards_update, + &mut self.model, + ); } - let _rebuild_plan_waiter = self.rebuild_plan_debounced(ctx); Ok(()) } } @@ -1183,6 +1196,9 @@ mod tests { indexer_pool, ), ingest_controller: IngestController::new(metastore.clone(), ingester_pool), + scaling_controller: ScalingController::with_shard_throughput_limit( + cluster_config.shard_throughput_limit, + ), legacy_scaling_controller: LegacyScalingController::new( cluster_config.shard_throughput_limit, cluster_config.shard_scale_up_factor, @@ -1217,7 +1233,7 @@ mod tests { } #[tokio::test(start_paused = true)] - async fn test_report_acknowledged_before_io_and_processing_errors() { + async fn test_report_acknowledged_before_periodic_reconciliation() { for outcome in ["success", "aborted", "uncertain"] { let universe = Universe::new(); let (mailbox, _inbox) = universe.create_test_mailbox(); @@ -1289,17 +1305,22 @@ mod tests { 2, ); let (reply_tx, mut reply_rx) = tokio::sync::oneshot::channel(); - let result = { - let handling = control_plane.handle_message( + control_plane + .handle_message( report_for_test(), move |reply| { reply_tx.send(reply).unwrap(); }, &ctx, - ); + ) + .await + .unwrap(); + assert!(reply_rx.try_recv().unwrap().is_ok()); + assert_eq!(control_plane.model.all_shards().count(), 1); + let result = { + let handling = Handler::handle(&mut control_plane, ControlPlaneLoop, &ctx); tokio::pin!(handling); assert!(futures::poll!(handling.as_mut()).is_pending()); - assert!(reply_rx.try_recv().unwrap().is_ok()); handling.await }; assert_eq!(result.is_err(), outcome == "uncertain"); diff --git a/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs b/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs index 5dc9c704f7d..c094e2fd1c4 100644 --- a/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs +++ b/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs @@ -16,9 +16,9 @@ use std::collections::HashMap; use std::num::NonZeroUsize; use bytesize::ByteSize; +use quickwit_actors::Mailbox; use quickwit_common::Progress; -use quickwit_ingest::{ShardInfos, SourceShardReport}; -use quickwit_proto::control_plane::ShardsUpdate; +use quickwit_ingest::ShardInfos; use quickwit_proto::ingest::Shard; use quickwit_proto::metastore::MetastoreResult; use quickwit_proto::types::{NodeId, SourceUid}; @@ -27,6 +27,7 @@ use rand::{Rng, rng}; use tracing::{error, info, warn}; use super::legacy_scaling_arbiter::LegacyScalingArbiter; +use crate::control_plane::ControlPlane; use crate::ingest::IngestController; use crate::model::{ControlPlaneModel, ScalingMode, ShardEntry, ShardStats}; @@ -44,51 +45,49 @@ impl LegacyScalingController { } } - pub(crate) async fn handle_shards_update( + pub(crate) async fn update_local_shards( &self, + // TODO: hold the ingest controller on the struct instead of passing it in. ingest_controller: &mut IngestController, - node_id: &str, - generation_id: u64, - shards_update: ShardsUpdate, + source_uid: SourceUid, + shard_infos: &ShardInfos, model: &mut ControlPlaneModel, progress: &Progress, ) -> MetastoreResult<()> { - if let Some(ingester) = ingest_controller.ingester_pool.get(node_id) - && generation_id != ingester.generation_id.as_u64() - { - return Ok(()); - } - for source_shard_infos in &shards_update.shard_infos_by_source { - let SourceShardReport { - source_uid, - shard_infos, - } = source_shard_infos.into(); - - if let Err(metastore_error) = self - .update_local_shards(ingest_controller, source_uid, &shard_infos, model, progress) - .await - { - if !metastore_error.is_transaction_certainly_aborted() { - return Err(metastore_error); - } - error!(error=?metastore_error, "failed to update source shards"); - } - } - Ok(()) + model.update_shards(&source_uid, shard_infos); + self.scale_source_shards(ingest_controller, source_uid, model, progress) + .await } - pub(crate) async fn update_local_shards( + pub(crate) async fn reconcile_shards( &self, // TODO: hold the ingest controller on the struct instead of passing it in. ingest_controller: &mut IngestController, - source_uid: SourceUid, - shard_infos: &ShardInfos, model: &mut ControlPlaneModel, + mailbox: &Mailbox, progress: &Progress, ) -> MetastoreResult<()> { - model.update_shards(&source_uid, shard_infos); - self.scale_source_shards(ingest_controller, source_uid, model, progress) - .await + let source_uids: Vec = model + .source_configs() + .map(|(source_uid, _source_config)| source_uid) + .collect(); + + for source_uid in source_uids { + let scale_source_shards_result = self + .scale_source_shards(ingest_controller, source_uid, model, progress) + .await; + let Err(metastore_error) = scale_source_shards_result else { + continue; + }; + if !metastore_error.is_transaction_certainly_aborted() { + return Err(metastore_error); + } + error!(error=?metastore_error, "failed to scale source shards"); + } + ingest_controller + .rebalance_shards(model, mailbox, progress) + .await?; + Ok(()) } async fn scale_source_shards( @@ -281,6 +280,7 @@ fn find_scale_down_candidate(source_uid: &SourceUid, model: &ControlPlaneModel) .max_by_key(|(_ingester_id, shard_entries)| (shard_entries.len(), rng.next_u32())) .map(|(_ingester_id, shard_entries)| shard_entries.choose(&mut rng).unwrap().shard.clone()) } + #[cfg(test)] mod tests { use std::collections::BTreeSet; @@ -290,9 +290,8 @@ mod tests { use quickwit_common::Progress; use quickwit_common::shared_consts::DEFAULT_SHARD_THROUGHPUT_LIMIT; use quickwit_config::{INGEST_V2_SOURCE_ID, SourceConfig}; - use quickwit_ingest::{IngesterPool, IngesterPoolEntry, ShardInfo, SourceShardReport}; + use quickwit_ingest::{IngesterPool, IngesterPoolEntry, ShardInfo}; use quickwit_metastore::IndexMetadata; - use quickwit_proto::control_plane::ShardsUpdate; use quickwit_proto::ingest::ingester::{ CloseShardsResponse, IngesterServiceClient, InitShardSubrequest, InitShardSuccess, InitShardsRequest, InitShardsResponse, MockIngesterService, @@ -1101,187 +1100,4 @@ mod tests { // We pick ingester 1 has it has more open shard assert_eq!(shard.ingester_id, "test-ingester-1"); } - - fn shard_reports_for_test(min_shards: usize) -> (ControlPlaneModel, ShardsUpdate) { - let mut model = ControlPlaneModel::default(); - let mut metadata = IndexMetadata::for_test("index", "ram:///index"); - metadata.index_config.ingest_settings.min_shards = - std::num::NonZeroUsize::new(min_shards).unwrap(); - let index_uid = metadata.index_uid.clone(); - model.add_index(metadata); - let mut update = ShardsUpdate::default(); - for source_id in ["source-a", "source-b"] { - model - .add_source( - &index_uid, - SourceConfig::for_test(source_id, quickwit_config::SourceParams::void()), - ) - .unwrap(); - model.insert_shards( - &index_uid, - &source_id.to_string(), - vec![Shard { - index_uid: Some(index_uid.clone()), - source_id: source_id.to_string(), - shard_id: Some(ShardId::from(1)), - ingester_id: "ingester".to_string(), - shard_state: ShardState::Open as i32, - ..Default::default() - }], - ); - update.shard_infos_by_source.push( - SourceShardReport { - source_uid: SourceUid { - index_uid: index_uid.clone(), - source_id: source_id.to_string(), - }, - shard_infos: BTreeSet::from([ShardInfo { - shard_id: ShardId::from(1), - shard_state: ShardState::Open, - short_term_ingestion_rate: ByteSize::b(123), - long_term_ingestion_rate: ByteSize::b(456), - }]), - } - .into(), - ); - } - (model, update) - } - - #[tokio::test] - async fn test_shards_update_sources_and_generation() { - let pool = IngesterPool::default(); - let ingester = IngesterPoolEntry::ready_with_client(IngesterServiceClient::mocked()); - let generation = ingester.generation_id.as_u64(); - pool.insert(NodeId::from_str("ingester"), ingester); - let mut controller = IngestController::new(MetastoreServiceClient::mocked(), pool); - let scaling_controller = LegacyScalingController::new(DEFAULT_SHARD_THROUGHPUT_LIMIT, 1.5); - let progress = Progress::default(); - let (mut model, update) = shard_reports_for_test(1); - for wrong_generation in [generation - 1, generation + 1] { - scaling_controller - .handle_shards_update( - &mut controller, - "ingester", - wrong_generation, - update.clone(), - &mut model, - &progress, - ) - .await - .unwrap(); - assert!( - model - .all_shards() - .all(|shard| shard.short_term_ingestion_rate == ByteSize::default()) - ); - } - scaling_controller - .handle_shards_update( - &mut controller, - "ingester", - generation, - update.clone(), - &mut model, - &progress, - ) - .await - .unwrap(); - assert_eq!(model.all_shards().count(), 2); - assert!( - model - .all_shards() - .all(|shard| shard.short_term_ingestion_rate == ByteSize::b(123) - && shard.long_term_ingestion_rate == ByteSize::b(456)) - ); - - let (mut model, _) = shard_reports_for_test(1); - scaling_controller - .handle_shards_update( - &mut controller, - "joining-ingester", - generation, - update, - &mut model, - &progress, - ) - .await - .unwrap(); - assert!( - model - .all_shards() - .all(|shard| shard.short_term_ingestion_rate == ByteSize::b(123)) - ); - } - - #[tokio::test] - async fn test_shards_update_continues_only_after_certainly_aborted_errors() { - for certainly_aborted in [true, false] { - let (mut model, update) = shard_reports_for_test(2); - let mut mock_metastore = MockMetastoreService::new(); - mock_metastore - .expect_open_shards() - .times(if certainly_aborted { 2 } else { 1 }) - .returning(move |_| { - if certainly_aborted { - Err(MetastoreError::InvalidArgument { - message: "aborted".to_string(), - }) - } else { - Err(MetastoreError::Connection { - message: "uncertain".to_string(), - }) - } - }); - let mut mock_ingester = MockIngesterService::new(); - mock_ingester - .expect_init_shards() - .times(if certainly_aborted { 2 } else { 1 }) - .returning(|request| { - Ok(InitShardsResponse { - successes: request - .subrequests - .into_iter() - .map(|request| InitShardSuccess { - subrequest_id: request.subrequest_id, - shard: request.shard, - }) - .collect(), - failures: Vec::new(), - }) - }); - let pool = IngesterPool::default(); - pool.insert( - NodeId::from_str("ingester"), - IngesterPoolEntry::ready_with_client(IngesterServiceClient::from_mock( - mock_ingester, - )), - ); - let mut controller = - IngestController::new(MetastoreServiceClient::from_mock(mock_metastore), pool); - let scaling_controller = - LegacyScalingController::new(DEFAULT_SHARD_THROUGHPUT_LIMIT, 1.5); - let result = scaling_controller - .handle_shards_update( - &mut controller, - "ingester", - 1, - update, - &mut model, - &Progress::default(), - ) - .await; - assert_eq!(result.is_ok(), certainly_aborted); - let second_source = model - .all_shards() - .find(|shard| shard.source_id == "source-b") - .unwrap(); - let expected = if certainly_aborted { - ByteSize::b(123) - } else { - ByteSize::default() - }; - assert_eq!(second_source.short_term_ingestion_rate, expected); - } - } } diff --git a/quickwit/quickwit-control-plane/src/ingest/mod.rs b/quickwit/quickwit-control-plane/src/ingest/mod.rs index 3aca5497826..65fed21634b 100644 --- a/quickwit/quickwit-control-plane/src/ingest/mod.rs +++ b/quickwit/quickwit-control-plane/src/ingest/mod.rs @@ -15,8 +15,10 @@ pub(crate) mod ingest_controller; mod legacy_scaling_arbiter; mod legacy_scaling_controller; +mod scaling_controller; mod wait_handle; pub use ingest_controller::IngestController; pub(crate) use legacy_scaling_controller::LegacyScalingController; +pub(crate) use scaling_controller::ScalingController; pub use wait_handle::WaitHandle; diff --git a/quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs b/quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs new file mode 100644 index 00000000000..84785782f64 --- /dev/null +++ b/quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs @@ -0,0 +1,552 @@ +// Copyright 2021-Present Datadog, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::collections::HashMap; +use std::num::NonZeroUsize; +use std::time::{Duration, Instant}; + +use bytesize::ByteSize; +use fnv::FnvHashSet; +use quickwit_actors::Mailbox; +use quickwit_common::Progress; +use quickwit_ingest::SourceShardReport; +use quickwit_proto::control_plane::ShardsUpdate; +use quickwit_proto::ingest::Shard; +use quickwit_proto::metastore::MetastoreResult; +use quickwit_proto::types::{NodeId, SourceUid}; +use tracing::{error, info, warn}; + +use crate::control_plane::ControlPlane; +use crate::ingest::IngestController; +use crate::model::{ControlPlaneModel, ShardThroughputStats}; + +const SCALE_DOWN_COOLDOWN: Duration = Duration::from_mins(5); + +#[derive(Debug, Clone, Copy, Eq, PartialEq)] +pub(crate) enum ScalingDecision { + ScaleUp { target_num_open_shards: usize }, + ScaleDown { target_num_open_shards: usize }, +} + +pub(crate) struct ScalingController { + target_shard_throughput: ByteSize, + scale_up_shard_throughput_threshold: ByteSize, + scale_down_shard_throughput_threshold: ByteSize, + last_shard_count_changes: HashMap, +} + +impl ScalingController { + pub fn with_shard_throughput_limit(shard_throughput_limit: ByteSize) -> ScalingController { + let shard_throughput_limit_bytes = shard_throughput_limit.as_u64(); + ScalingController { + target_shard_throughput: ByteSize::b(shard_throughput_limit_bytes * 8 / 10), + scale_up_shard_throughput_threshold: shard_throughput_limit, + scale_down_shard_throughput_threshold: ByteSize::b(shard_throughput_limit_bytes / 2), + last_shard_count_changes: HashMap::new(), + } + } + + fn is_scale_down_cooldown_expired(&self, source_uid: &SourceUid, now: Instant) -> bool { + let Some(last_shard_count_change) = self.last_shard_count_changes.get(source_uid) else { + return true; + }; + now.duration_since(*last_shard_count_change) >= SCALE_DOWN_COOLDOWN + } + + fn restart_scale_down_cooldown(&mut self, source_uid: &SourceUid, now: Instant) { + self.last_shard_count_changes + .insert(source_uid.clone(), now); + } + + pub(crate) fn handle_shards_update( + &self, + // TODO: hold the ingest controller on the struct instead of passing it in. + ingest_controller: &IngestController, + node_id: &str, + generation_id: u64, + shards_update: ShardsUpdate, + model: &mut ControlPlaneModel, + ) { + if let Some(ingester) = ingest_controller.ingester_pool.get(node_id) + && generation_id != ingester.generation_id.as_u64() + { + return; + } + for source_shard_infos in &shards_update.shard_infos_by_source { + let SourceShardReport { + source_uid, + shard_infos, + } = source_shard_infos.into(); + + model.update_shards(&source_uid, &shard_infos); + } + } + + pub(crate) async fn reconcile_shards( + &mut self, + // TODO: hold the ingest controller on the struct instead of passing it in. + ingest_controller: &mut IngestController, + model: &mut ControlPlaneModel, + mailbox: &Mailbox, + progress: &Progress, + ) -> MetastoreResult<()> { + let live_ingesters: FnvHashSet = + ingest_controller.ingester_pool.keys().into_iter().collect(); + let source_uids: Vec = model + .source_configs() + .map(|(source_uid, _source_config)| source_uid) + .collect(); + + for source_uid in source_uids { + let scale_source_shards_result = self + .scale_source_shards( + ingest_controller, + &source_uid, + &live_ingesters, + model, + progress, + ) + .await; + let Err(metastore_error) = scale_source_shards_result else { + continue; + }; + if !metastore_error.is_transaction_certainly_aborted() { + return Err(metastore_error); + } + error!(error=?metastore_error, "failed to scale source shards"); + } + ingest_controller + .rebalance_shards(model, mailbox, progress) + .await?; + Ok(()) + } + + async fn scale_source_shards( + &mut self, + // TODO: hold the ingest controller on the struct instead of passing it in. + ingest_controller: &mut IngestController, + source_uid: &SourceUid, + live_ingesters: &FnvHashSet, + model: &mut ControlPlaneModel, + progress: &Progress, + ) -> MetastoreResult<()> { + let Some(shard_throughput_stats) = model.shard_throughput_stats(source_uid, live_ingesters) + else { + return Ok(()); + }; + let min_shards = model + .index_metadata(&source_uid.index_uid) + .expect("index should exist") + .index_config + .ingest_settings + .min_shards; + let cooldown_expired = self.is_scale_down_cooldown_expired(source_uid, Instant::now()); + let scaling_decision_opt = + self.should_scale(shard_throughput_stats, min_shards, cooldown_expired); + let Some(scaling_decision) = scaling_decision_opt else { + return Ok(()); + }; + let num_open_shards = shard_throughput_stats.num_open_shards; + + match scaling_decision { + ScalingDecision::ScaleUp { + target_num_open_shards, + } => { + self.scale_up_shards( + ingest_controller, + source_uid, + num_open_shards, + target_num_open_shards, + model, + progress, + ) + .await + } + ScalingDecision::ScaleDown { + target_num_open_shards, + } => { + self.scale_down_shards( + ingest_controller, + source_uid, + num_open_shards, + target_num_open_shards, + live_ingesters, + model, + progress, + ) + .await; + Ok(()) + } + } + } + + async fn scale_up_shards( + &mut self, + // TODO: hold the ingest controller on the struct instead of passing it in. + ingest_controller: &mut IngestController, + source_uid: &SourceUid, + num_open_shards: usize, + target_num_open_shards: usize, + model: &mut ControlPlaneModel, + progress: &Progress, + ) -> MetastoreResult<()> { + let num_shards_to_open = target_num_open_shards - num_open_shards; + let num_opened_shards = ingest_controller + .open_shards_for_source(source_uid, num_shards_to_open, model, progress) + .await?; + + if num_opened_shards == 0 { + warn!( + index_uid=%source_uid.index_uid, + source_id=%source_uid.source_id, + "failed to scale up number of shards from {num_open_shards} to {target_num_open_shards}" + ); + return Ok(()); + } + self.restart_scale_down_cooldown(source_uid, Instant::now()); + info!( + index_uid=%source_uid.index_uid, + source_id=%source_uid.source_id, + "scaled up number of shards from {num_open_shards} by {num_opened_shards} (target {target_num_open_shards})" + ); + Ok(()) + } + + #[allow(clippy::too_many_arguments)] + async fn scale_down_shards( + &mut self, + // TODO: hold the ingest controller on the struct instead of passing it in. + ingest_controller: &IngestController, + source_uid: &SourceUid, + num_open_shards: usize, + target_num_open_shards: usize, + live_ingesters: &FnvHashSet, + model: &mut ControlPlaneModel, + progress: &Progress, + ) { + let num_shards_to_close = num_open_shards - target_num_open_shards; + let shards_to_close = + find_scale_down_candidates(source_uid, num_shards_to_close, live_ingesters, model); + let closed_shard_ids = ingest_controller + .close_source_shards(source_uid, shards_to_close, model, progress) + .await; + + if closed_shard_ids.is_empty() { + warn!( + index_uid=%source_uid.index_uid, + source_id=%source_uid.source_id, + "failed to scale down number of shards from {num_open_shards} to {target_num_open_shards}" + ); + return; + } + self.restart_scale_down_cooldown(source_uid, Instant::now()); + info!( + index_uid=%source_uid.index_uid, + source_id=%source_uid.source_id, + "scaled down number of shards from {num_open_shards} by {} (target {target_num_open_shards})", + closed_shard_ids.len() + ); + } + + pub fn should_scale( + &self, + shard_throughput_stats: ShardThroughputStats, + min_shards: NonZeroUsize, + cooldown_expired: bool, + ) -> Option { + let num_open_shards = shard_throughput_stats.num_open_shards; + if num_open_shards == 0 { + return None; + } + let ingestion_rate = shard_throughput_stats + .total_short_term_ingestion_rate + .max(shard_throughput_stats.total_long_term_ingestion_rate); + let target_num_open_shards = self.target_num_open_shards(ingestion_rate, min_shards); + + if num_open_shards < min_shards.get() { + let scale_up = ScalingDecision::ScaleUp { + target_num_open_shards, + }; + return Some(scale_up); + } + let scale_up_ingestion_rate = + self.scale_up_shard_throughput_threshold * num_open_shards as u64; + if ingestion_rate > scale_up_ingestion_rate { + let scale_up = ScalingDecision::ScaleUp { + target_num_open_shards, + }; + return Some(scale_up); + } + let scale_down_ingestion_rate = + self.scale_down_shard_throughput_threshold * num_open_shards as u64; + if ingestion_rate < scale_down_ingestion_rate + && cooldown_expired + && target_num_open_shards < num_open_shards + { + let scale_down = ScalingDecision::ScaleDown { + target_num_open_shards, + }; + return Some(scale_down); + } + None + } + + fn target_num_open_shards(&self, ingestion_rate: ByteSize, min_shards: NonZeroUsize) -> usize { + let num_shards_for_ingestion_rate = ingestion_rate + .as_u64() + .div_ceil(self.target_shard_throughput.as_u64()) + as usize; + num_shards_for_ingestion_rate.max(min_shards.get()) + } +} + +fn find_scale_down_candidates( + source_uid: &SourceUid, + num_shards_to_close: usize, + live_ingesters: &FnvHashSet, + model: &ControlPlaneModel, +) -> Vec { + let Some(source_shard_entries) = model.get_shards_for_source(source_uid) else { + return Vec::new(); + }; + let mut num_open_shards_by_ingester_id: HashMap<&str, usize> = HashMap::new(); + + for shard_entry in model.all_shards() { + if shard_entry.is_open() { + *num_open_shards_by_ingester_id + .entry(shard_entry.ingester_id.as_str()) + .or_default() += 1; + } + } + let mut source_open_shards_by_ingester_id: HashMap<&str, Vec<&Shard>> = HashMap::new(); + + for shard_entry in source_shard_entries.values() { + if shard_entry.is_open() && live_ingesters.contains(shard_entry.ingester_id.as_str()) { + source_open_shards_by_ingester_id + .entry(shard_entry.ingester_id.as_str()) + .or_default() + .push(&shard_entry.shard); + } + } + let mut shards_to_close: Vec = Vec::with_capacity(num_shards_to_close); + + for _ in 0..num_shards_to_close { + let most_loaded_ingester_id_opt = source_open_shards_by_ingester_id + .iter() + .filter(|(_ingester_id, shards)| !shards.is_empty()) + .map(|(ingester_id, _shards)| *ingester_id) + .max_by_key(|ingester_id| num_open_shards_by_ingester_id[ingester_id]); + + let Some(most_loaded_ingester_id) = most_loaded_ingester_id_opt else { + break; + }; + let shard = source_open_shards_by_ingester_id + .get_mut(most_loaded_ingester_id) + .and_then(|shards| shards.pop()) + .expect("ingester should have an open shard for the source"); + *num_open_shards_by_ingester_id + .get_mut(most_loaded_ingester_id) + .expect("ingester should have open shards") -= 1; + shards_to_close.push(shard.clone()); + } + shards_to_close +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeSet; + + use bytesize::ByteSize; + use quickwit_actors::Universe; + use quickwit_common::Progress; + use quickwit_common::shared_consts::DEFAULT_SHARD_THROUGHPUT_LIMIT; + use quickwit_config::SourceConfig; + use quickwit_ingest::{IngesterPool, IngesterPoolEntry, ShardInfo, SourceShardReport}; + use quickwit_metastore::IndexMetadata; + use quickwit_proto::control_plane::ShardsUpdate; + use quickwit_proto::ingest::ingester::{ + IngesterServiceClient, InitShardSuccess, InitShardsResponse, MockIngesterService, + }; + use quickwit_proto::ingest::{Shard, ShardState}; + use quickwit_proto::metastore::{MetastoreError, MetastoreServiceClient, MockMetastoreService}; + use quickwit_proto::types::{NodeId, ShardId, SourceUid}; + + use super::{IngestController, ScalingController}; + use crate::ingest::LegacyScalingController; + use crate::model::ControlPlaneModel; + + fn shard_reports_for_test(min_shards: usize) -> (ControlPlaneModel, ShardsUpdate) { + let mut model = ControlPlaneModel::default(); + let mut metadata = IndexMetadata::for_test("index", "ram:///index"); + metadata.index_config.ingest_settings.min_shards = + std::num::NonZeroUsize::new(min_shards).unwrap(); + let index_uid = metadata.index_uid.clone(); + model.add_index(metadata); + let mut update = ShardsUpdate::default(); + for source_id in ["source-a", "source-b"] { + model + .add_source( + &index_uid, + SourceConfig::for_test(source_id, quickwit_config::SourceParams::void()), + ) + .unwrap(); + model.insert_shards( + &index_uid, + &source_id.to_string(), + vec![Shard { + index_uid: Some(index_uid.clone()), + source_id: source_id.to_string(), + shard_id: Some(ShardId::from(1)), + ingester_id: "ingester".to_string(), + shard_state: ShardState::Open as i32, + ..Default::default() + }], + ); + update.shard_infos_by_source.push( + SourceShardReport { + source_uid: SourceUid { + index_uid: index_uid.clone(), + source_id: source_id.to_string(), + }, + shard_infos: BTreeSet::from([ShardInfo { + shard_id: ShardId::from(1), + shard_state: ShardState::Open, + short_term_ingestion_rate: ByteSize::b(123), + long_term_ingestion_rate: ByteSize::b(456), + }]), + } + .into(), + ); + } + (model, update) + } + + #[tokio::test] + async fn test_shards_update_sources_and_generation() { + let pool = IngesterPool::default(); + let ingester = IngesterPoolEntry::ready_with_client(IngesterServiceClient::mocked()); + let generation = ingester.generation_id.as_u64(); + pool.insert(NodeId::from_str("ingester"), ingester); + let controller = IngestController::new(MetastoreServiceClient::mocked(), pool); + let scaling_controller = + ScalingController::with_shard_throughput_limit(DEFAULT_SHARD_THROUGHPUT_LIMIT); + let (mut model, update) = shard_reports_for_test(1); + for wrong_generation in [generation - 1, generation + 1] { + scaling_controller.handle_shards_update( + &controller, + "ingester", + wrong_generation, + update.clone(), + &mut model, + ); + assert!( + model + .all_shards() + .all(|shard| shard.short_term_ingestion_rate == ByteSize::default()) + ); + } + scaling_controller.handle_shards_update( + &controller, + "ingester", + generation, + update.clone(), + &mut model, + ); + assert_eq!(model.all_shards().count(), 2); + assert!( + model + .all_shards() + .all(|shard| shard.short_term_ingestion_rate == ByteSize::b(123) + && shard.long_term_ingestion_rate == ByteSize::b(456)) + ); + + let (mut model, _) = shard_reports_for_test(1); + scaling_controller.handle_shards_update( + &controller, + "joining-ingester", + generation, + update, + &mut model, + ); + assert!( + model + .all_shards() + .all(|shard| shard.short_term_ingestion_rate == ByteSize::b(123)) + ); + } + + #[tokio::test] + async fn test_reconciliation_continues_only_after_certainly_aborted_errors() { + for certainly_aborted in [true, false] { + let (mut model, update) = shard_reports_for_test(2); + let universe = Universe::new(); + let (mailbox, _inbox) = universe.create_test_mailbox(); + let mut mock_metastore = MockMetastoreService::new(); + mock_metastore + .expect_open_shards() + .times(if certainly_aborted { 2 } else { 1 }) + .returning(move |_| { + if certainly_aborted { + Err(MetastoreError::InvalidArgument { + message: "aborted".to_string(), + }) + } else { + Err(MetastoreError::Connection { + message: "uncertain".to_string(), + }) + } + }); + let mut mock_ingester = MockIngesterService::new(); + mock_ingester + .expect_init_shards() + .times(if certainly_aborted { 2 } else { 1 }) + .returning(|request| { + Ok(InitShardsResponse { + successes: request + .subrequests + .into_iter() + .map(|request| InitShardSuccess { + subrequest_id: request.subrequest_id, + shard: request.shard, + }) + .collect(), + failures: Vec::new(), + }) + }); + let pool = IngesterPool::default(); + pool.insert( + NodeId::from_str("ingester"), + IngesterPoolEntry::ready_with_client(IngesterServiceClient::from_mock( + mock_ingester, + )), + ); + let mut controller = + IngestController::new(MetastoreServiceClient::from_mock(mock_metastore), pool); + let scaling_controller = + LegacyScalingController::new(DEFAULT_SHARD_THROUGHPUT_LIMIT, 1.5); + ScalingController::with_shard_throughput_limit(DEFAULT_SHARD_THROUGHPUT_LIMIT) + .handle_shards_update(&controller, "ingester", 1, update, &mut model); + let result = scaling_controller + .reconcile_shards(&mut controller, &mut model, &mailbox, &Progress::default()) + .await; + assert_eq!(result.is_ok(), certainly_aborted); + let second_source = model + .all_shards() + .find(|shard| shard.source_id == "source-b") + .unwrap(); + assert_eq!(second_source.short_term_ingestion_rate, ByteSize::b(123)); + universe.assert_quit().await; + } + } +} diff --git a/quickwit/quickwit-control-plane/src/model/mod.rs b/quickwit/quickwit-control-plane/src/model/mod.rs index 6972fa5572b..f8c8622f305 100644 --- a/quickwit/quickwit-control-plane/src/model/mod.rs +++ b/quickwit/quickwit-control-plane/src/model/mod.rs @@ -36,7 +36,9 @@ use quickwit_proto::metastore::{ MetastoreServiceClient, SourceType, ToggleSourceRequest, }; use quickwit_proto::types::{IndexId, IndexUid, NodeId, ShardId, SourceId, SourceUid}; -pub(super) use shard_table::{ScalingMode, ShardEntry, ShardLocations, ShardStats, ShardTable}; +pub(super) use shard_table::{ + ScalingMode, ShardEntry, ShardLocations, ShardStats, ShardTable, ShardThroughputStats, +}; use tracing::{debug, error, info, instrument, warn}; use crate::metrics::INDEXES_TOTAL; @@ -370,6 +372,15 @@ impl ControlPlaneModel { .find_open_shards(index_uid, source_id, unavailable_ingesters) } + pub fn shard_throughput_stats( + &self, + source_uid: &SourceUid, + live_ingesters: &FnvHashSet, + ) -> Option { + self.shard_table + .shard_throughput_stats(source_uid, live_ingesters) + } + pub fn legacy_shard_stats(&self, source_uid: &SourceUid) -> Option { self.shard_table.legacy_shard_stats(source_uid) } diff --git a/quickwit/quickwit-control-plane/src/model/shard_table.rs b/quickwit/quickwit-control-plane/src/model/shard_table.rs index 59c9849869b..8cab630862e 100644 --- a/quickwit/quickwit-control-plane/src/model/shard_table.rs +++ b/quickwit/quickwit-control-plane/src/model/shard_table.rs @@ -139,6 +139,28 @@ impl ShardTableEntry { avg_long_term_ingestion_rate, } } + + fn shard_throughput_stats(&self, live_ingesters: &FnvHashSet) -> ShardThroughputStats { + let mut num_open_shards = 0; + let mut total_short_term_ingestion_rate = ByteSize::default(); + let mut total_long_term_ingestion_rate = ByteSize::default(); + + for shard_entry in self.shard_entries.values() { + if !live_ingesters.contains(shard_entry.ingester_id.as_str()) { + continue; + } + if shard_entry.is_open() { + num_open_shards += 1; + } + total_short_term_ingestion_rate += shard_entry.short_term_ingestion_rate; + total_long_term_ingestion_rate += shard_entry.long_term_ingestion_rate; + } + ShardThroughputStats { + num_open_shards, + total_short_term_ingestion_rate, + total_long_term_ingestion_rate, + } + } } #[derive(Default)] @@ -451,6 +473,16 @@ impl ShardTable { Some(open_shards) } + pub fn shard_throughput_stats( + &self, + source_uid: &SourceUid, + live_ingesters: &FnvHashSet, + ) -> Option { + let table_entry = self.table_entries.get(source_uid)?; + let shard_throughput_stats = table_entry.shard_throughput_stats(live_ingesters); + Some(shard_throughput_stats) + } + pub fn legacy_shard_stats(&self, source_uid: &SourceUid) -> Option { let table_entry = self.table_entries.get(source_uid)?; let shard_stats = table_entry.shards_stats(); @@ -609,6 +641,13 @@ pub(crate) struct ShardStats { pub avg_long_term_ingestion_rate: ByteSize, } +#[derive(Clone, Copy, Default)] +pub(crate) struct ShardThroughputStats { + pub num_open_shards: usize, + pub total_short_term_ingestion_rate: ByteSize, + pub total_long_term_ingestion_rate: ByteSize, +} + #[cfg(test)] mod tests { use std::collections::BTreeSet; @@ -849,6 +888,87 @@ mod tests { assert_eq!(open_shards[0].shard, shard_04); } + #[test] + fn test_shard_table_shard_throughput_stats() { + let index_uid: IndexUid = IndexUid::for_test("test-index", 0); + let source_id = "test-source".to_string(); + let source_uid = SourceUid { + index_uid: index_uid.clone(), + source_id: source_id.clone(), + }; + let live_ingesters = FnvHashSet::from_iter([NodeId::from_str("test-ingester-0")]); + + let mut shard_table = ShardTable::default(); + assert!( + shard_table + .shard_throughput_stats(&source_uid, &live_ingesters) + .is_none() + ); + + shard_table.add_source(&index_uid, &source_id); + + let shard_01 = Shard { + index_uid: index_uid.clone().into(), + source_id: source_id.clone(), + shard_id: Some(ShardId::from(1)), + ingester_id: "test-ingester-0".to_string(), + shard_state: ShardState::Open as i32, + ..Default::default() + }; + let shard_02 = Shard { + index_uid: index_uid.clone().into(), + source_id: source_id.clone(), + shard_id: Some(ShardId::from(2)), + ingester_id: "test-ingester-0".to_string(), + shard_state: ShardState::Closed as i32, + ..Default::default() + }; + let shard_03 = Shard { + index_uid: index_uid.clone().into(), + source_id: source_id.clone(), + shard_id: Some(ShardId::from(3)), + ingester_id: "test-ingester-1".to_string(), + shard_state: ShardState::Open as i32, + ..Default::default() + }; + shard_table.insert_shards(&index_uid, &source_id, vec![shard_01, shard_02, shard_03]); + + let shard_infos = BTreeSet::from_iter([ + ShardInfo { + shard_id: ShardId::from(1), + shard_state: ShardState::Open, + short_term_ingestion_rate: ByteSize::mib(2), + long_term_ingestion_rate: ByteSize::mib(1), + }, + ShardInfo { + shard_id: ShardId::from(2), + shard_state: ShardState::Closed, + short_term_ingestion_rate: ByteSize::mib(0), + long_term_ingestion_rate: ByteSize::mib(3), + }, + ShardInfo { + shard_id: ShardId::from(3), + shard_state: ShardState::Open, + short_term_ingestion_rate: ByteSize::mib(4), + long_term_ingestion_rate: ByteSize::mib(4), + }, + ]); + shard_table.update_shards(&source_uid, &shard_infos); + + let shard_throughput_stats = shard_table + .shard_throughput_stats(&source_uid, &live_ingesters) + .unwrap(); + assert_eq!(shard_throughput_stats.num_open_shards, 1); + assert_eq!( + shard_throughput_stats.total_short_term_ingestion_rate, + ByteSize::mib(2) + ); + assert_eq!( + shard_throughput_stats.total_long_term_ingestion_rate, + ByteSize::mib(4) + ); + } + #[test] fn test_shard_table_update_shards() { let index_uid: IndexUid = IndexUid::for_test("test-index", 0); From 8b86ce55e690781bea003381ef4b98b2f0f57070 Mon Sep 17 00:00:00 2001 From: "nadav.govari" Date: Mon, 5 Oct 2026 11:05:26 -0400 Subject: [PATCH 2/2] Dont scale down during rebalance --- .../quickwit-control-plane/src/ingest/ingest_controller.rs | 4 ++++ .../src/ingest/legacy_scaling_controller.rs | 2 +- .../quickwit-control-plane/src/ingest/scaling_controller.rs | 3 +++ 3 files changed, 8 insertions(+), 1 deletion(-) diff --git a/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs b/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs index f4112d07be6..4d7909f6e80 100644 --- a/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs +++ b/quickwit/quickwit-control-plane/src/ingest/ingest_controller.rs @@ -958,6 +958,10 @@ impl IngestController { Ok(num_opened_shards) } + pub(crate) fn is_rebalancing(&self) -> bool { + self.rebalance_semaphore.available_permits() == 0 + } + pub(crate) async fn close_source_shards( &self, source_uid: &SourceUid, diff --git a/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs b/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs index c094e2fd1c4..df38d04d02a 100644 --- a/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs +++ b/quickwit/quickwit-control-plane/src/ingest/legacy_scaling_controller.rs @@ -219,7 +219,7 @@ impl LegacyScalingController { ) -> MetastoreResult<()> { // The scaling arbiter should not suggest scaling down if the number of shards is already // below the minimum, but we're just being defensive here. - if shard_stats.num_open_shards <= min_shards.get() { + if ingest_controller.is_rebalancing() || shard_stats.num_open_shards <= min_shards.get() { return Ok(()); } if !model diff --git a/quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs b/quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs index 84785782f64..143001954bc 100644 --- a/quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs +++ b/quickwit/quickwit-control-plane/src/ingest/scaling_controller.rs @@ -235,6 +235,9 @@ impl ScalingController { model: &mut ControlPlaneModel, progress: &Progress, ) { + if ingest_controller.is_rebalancing() { + return; + } let num_shards_to_close = num_open_shards - target_num_open_shards; let shards_to_close = find_scale_down_candidates(source_uid, num_shards_to_close, live_ingesters, model);