Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
151 changes: 150 additions & 1 deletion quickwit/quickwit-control-plane/src/indexing_scheduler/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,13 +58,26 @@ pub(crate) const APPLY_INDEXING_PLAN_TIMEOUT: Duration = if cfg!(any(test, featu
Duration::from_secs(2)
};

/// Number of consecutive control-loop cycles in which the running plan keeps differing from the
/// last applied plan (each triggering a reapply of the identical plan) before logging escalates
/// from `info` to `warn`.
///
/// At `MIN_DURATION_BETWEEN_SCHEDULING` of 30s this corresponds to ~2.5 minutes of non-convergence,
/// long enough to exclude the normal transient case where an indexer has not yet reported
/// its freshly applied tasks, while surfacing a persistently stuck loop.
const REAPPLY_LOOP_WARN_THRESHOLD: usize = 5;

#[derive(Debug, Clone, Default, Serialize)]
pub struct IndexingSchedulerState {
pub num_applied_physical_indexing_plan: usize,
pub num_schedule_indexing_plan: usize,
pub last_applied_physical_plan: Option<PhysicalIndexingPlan>,
#[serde(skip)]
pub last_applied_plan_timestamp: Option<Instant>,
/// Consecutive control-loop cycles that reapplied the identical plan
/// without the cluster converging. Reset to 0 on convergence or reschedule;
/// once it reaches `REAPPLY_LOOP_WARN_THRESHOLD`, logging escalates from `info` to `warn`.
pub num_consecutive_reapplies: usize,
}

/// The [`IndexingScheduler`] is responsible for listing indexing tasks and assigning them to
Expand Down Expand Up @@ -384,12 +397,31 @@ impl IndexingScheduler {
last_applied_plan.indexing_tasks_per_indexer(),
);
if !indexing_plans_diff.has_same_nodes() {
// `rebuild_plan` installs a new plan, which resets the reapply counter in
// `apply_physical_indexing_plan`.
info!(plans_diff=?indexing_plans_diff, "running plan and last applied plan node IDs differ: schedule an indexing plan");
self.rebuild_plan(model);
} else if !indexing_plans_diff.has_same_tasks() {
// Some nodes may have not received their tasks, apply it again.
info!(plans_diff=?indexing_plans_diff, "running tasks and last applied tasks differ: reapply last plan");
self.state.num_consecutive_reapplies += 1;
// A single reapply is the normal transient case (an indexer has not reported its
// freshly applied tasks yet). Escalation to `warn` happens only once the loop persists,
// since a plan that never converges indicates a stuck cluster rather than propagation
// lag.
if self.state.num_consecutive_reapplies >= REAPPLY_LOOP_WARN_THRESHOLD {
warn!(
num_consecutive_reapplies = self.state.num_consecutive_reapplies,
plans_diff=?indexing_plans_diff,
"running tasks still differ from last applied plan after repeated reapplies: \
cluster is not converging"
);
} else {
info!(plans_diff=?indexing_plans_diff, "running tasks and last applied tasks differ: reapply last plan");
}
self.apply_physical_indexing_plan(last_applied_plan.clone(), None);
} else {
// Running plan matches the last applied plan: the cluster has converged.
self.state.num_consecutive_reapplies = 0;
}
}

Expand Down Expand Up @@ -425,6 +457,15 @@ impl IndexingScheduler {
) {
debug!(new_physical_plan=?new_physical_plan, "apply physical indexing plan");
APPLY_PLAN_TOTAL.inc();
// Installing a genuinely new plan (rebuild, rebalance, startup) resets the reapply counter,
// regardless of which path scheduled it. A reapply of the identical plan (from
// `control_running_plan`) leaves the counter untouched so it keeps growing while the
// cluster fails to converge.
let reapplies_same_plan =
self.state.last_applied_physical_plan.as_ref() == Some(&new_physical_plan);
if !reapplies_same_plan {
self.state.num_consecutive_reapplies = 0;
}
// Retiring and decommissioning indexers still receive the plan so they can gracefully shut
// down dropped pipelines; other states (initializing, decommissioned, failed) are skipped.
for indexer in self.indexer_pool.values().into_iter().filter(|indexer| {
Expand Down Expand Up @@ -1416,6 +1457,114 @@ mod tests {
);
}

// Builds an `IndexerNodeInfo` whose client accepts any number of `apply_indexing_plan` calls
// (including zero) and reports `indexing_tasks` as its running tasks.
fn accepting_indexer_node_info(
node_id: &str,
status: IngesterStatus,
indexing_tasks: Vec<IndexingTask>,
) -> IndexerNodeInfo {
let mut mock_indexer = MockIndexingService::new();
mock_indexer
.expect_apply_indexing_plan()
.times(0..)
.returning(|_| Ok(ApplyIndexingPlanResponse {}));
let client = IndexingServiceClient::from_mock(mock_indexer);
IndexerNodeInfo {
node_id: NodeId::from_str(node_id),
generation_id: 0,
client,
indexing_tasks,
indexing_capacity: CpuCapacity::from_cpu_millis(4_000),
ingester_status: status,
}
}

// When the running tasks keep differing from the last applied plan, the control loop reapplies
// the identical plan every cycle. Each consecutive reapply must be counted so that a cluster
// which never converges becomes observable, and the counter must reset to 0 as soon as the
// cluster converges.
#[tokio::test]
async fn test_control_running_plan_counts_consecutive_reapplies() {
let index_uid = IndexUid::from_str("index-1:11111111111111111111111111").unwrap();
let task = IndexingTask {
pipeline_uid: Some(PipelineUid::for_test(1u128)),
index_uid: Some(index_uid),
source_id: "source-1".to_string(),
shard_ids: Vec::new(),
params_fingerprint: 0,
};

// Same node in the running plan and the last applied plan, but the indexer runs no task
// while the plan expects one: this drives the "reapply last plan" branch every cycle.
let indexer_pool = IndexerPool::default();
let indexer = accepting_indexer_node_info("indexer-1", IngesterStatus::Ready, Vec::new());
indexer_pool.insert(indexer.node_id.clone(), indexer);

let mut scheduler = IndexingScheduler::new(
"test-cluster".to_string(),
NodeId::from_str("control-plane"),
indexer_pool.clone(),
);
let mut physical_plan = PhysicalIndexingPlan::with_indexer_ids(&["indexer-1".to_string()]);
physical_plan.add_indexing_task("indexer-1", task.clone());
scheduler.state.last_applied_physical_plan = Some(physical_plan);
// Leaving `last_applied_plan_timestamp` at `None` bypasses the min-interval guard so every
// call actually evaluates the diff.

let model = ControlPlaneModel::default();

// Reapplying a plan stamps `last_applied_plan_timestamp`, which would otherwise make the
// next call short-circuit on the min-interval guard; clearing it each cycle simulates that
// enough time has elapsed and isolates the counter logic under test.
let num_cycles = REAPPLY_LOOP_WARN_THRESHOLD + 1;
for expected in 1..=num_cycles {
scheduler.state.last_applied_plan_timestamp = None;
scheduler.control_running_plan(&model);
assert_eq!(scheduler.state.num_consecutive_reapplies, expected);
}

// The indexer now runs exactly the planned task: the cluster has converged (`Pool` is a
// shared handle, so re-inserting updates the scheduler's view) and the counter resets.
let converged_indexer =
accepting_indexer_node_info("indexer-1", IngesterStatus::Ready, vec![task]);
indexer_pool.insert(converged_indexer.node_id.clone(), converged_indexer);
scheduler.state.last_applied_plan_timestamp = None;
scheduler.control_running_plan(&model);
assert_eq!(scheduler.state.num_consecutive_reapplies, 0);
}

// The reapply counter must be reset whenever a genuinely new plan is installed, no matter which
// path scheduled it (control loop, `RebuildPlan` handler, rebalance). Installing a different
// plan resets it; reapplying the identical plan does not. An empty pool keeps
// `apply_physical_indexing_plan` from issuing RPCs.
#[tokio::test]
async fn test_apply_physical_indexing_plan_resets_reapplies_for_new_plan_only() {
let indexer_pool = IndexerPool::default();
let mut scheduler = IndexingScheduler::new(
"test-cluster".to_string(),
NodeId::from_str("control-plane"),
indexer_pool,
);

let plan_a = PhysicalIndexingPlan::with_indexer_ids(&["indexer-1".to_string()]);
scheduler.apply_physical_indexing_plan(plan_a.clone(), None);
// Simulate a loop that has been stuck reapplying `plan_a` for a while.
scheduler.state.num_consecutive_reapplies = 7;

// Reapplying the identical plan must not reset the counter.
scheduler.apply_physical_indexing_plan(plan_a, None);
assert_eq!(scheduler.state.num_consecutive_reapplies, 7);

// Installing a different plan (as rebuild/rebalance would) resets the counter.
let plan_b = PhysicalIndexingPlan::with_indexer_ids(&[
"indexer-1".to_string(),
"indexer-2".to_string(),
]);
scheduler.apply_physical_indexing_plan(plan_b, None);
assert_eq!(scheduler.state.num_consecutive_reapplies, 0);
}

fn kafka_source_params_for_test() -> SourceParams {
SourceParams::Kafka(KafkaSourceParams {
topic: "topic".to_string(),
Expand Down