Skip to content

Commit ff0b623

Browse files
committed
cluster-controller: batch ON REFRESH oracle reads
The controller requested refresh-window inputs one cluster at a time. Each request issued and awaited a timestamp oracle read, so a slow oracle made every scheduled cluster add another round trip before unrelated clusters could reconcile. Gather each scheduled cluster's catalog and storage inputs in a separate coordinator turn, then attach one shared oracle timestamp to the completed phase batch. Skip clusters whose inputs are unavailable instead of treating them as outside their refresh window. Add seam tests for one batch per phase and unavailable input handling. Add an mzcompose regression that compares probe-cluster convergence with zero and eight scheduled clusters under injected oracle latency. Fixes SQL-569.
1 parent 735e5e9 commit ff0b623

5 files changed

Lines changed: 424 additions & 96 deletions

File tree

‎src/adapter/src/coord/cluster_controller.rs‎

Lines changed: 66 additions & 50 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,10 @@
1414
//! the controller as a **separate task** and implements the ctx by marshaling
1515
//! each pull/apply to the Coordinator over the internal command channel, because
1616
//! the catalog and the live compute/storage signals are reachable only from the
17-
//! coordinator loop. The two whole-tick reads are batched; the per-cluster live
18-
//! signals are pulled on demand, so a tick's round-trips scale with the number of
19-
//! managed clusters that need a live signal, not with a constant.
17+
//! coordinator loop. Whole-tick reads are batched. Refresh-window catalog inputs
18+
//! are pulled one cluster at a time and completed with one shared oracle read.
19+
//! The remaining per-cluster live signals are pulled on demand, so steady
20+
//! clusters do not pay for signals they do not use.
2021
//!
2122
//! Everything here is gated by [`ENABLE_CLUSTER_CONTROLLER`] (default on). With
2223
//! the gate off the task does not tick, so the legacy scheduling and graceful
@@ -26,7 +27,7 @@
2627
//! replicas are materialized by `reconcile_builtin_cluster_replicas` at catalog
2728
//! open, which derives the same target from the same config.)
2829
29-
use std::collections::BTreeSet;
30+
use std::collections::{BTreeMap, BTreeSet};
3031
use std::sync::Arc;
3132
use std::time::Duration;
3233

@@ -36,7 +37,8 @@ use mz_cluster_controller::ClusterController;
3637
use mz_cluster_controller::ctx::{
3738
ApplyOutcome, AvailabilityZones, ClusterControllerCtx, ClusterState, CreateReason, Decision,
3839
ExpectedClusterState, ObservedReplica, OnTimeout, ReconfigurationRecord, ReconfigurationStatus,
39-
ReconfigurationTarget, RefreshMvInfo, RefreshWindowInputs, ReplicaShape, StateWrite,
40+
ReconfigurationTarget, RefreshMvInfo, RefreshWindowClusterInputs, RefreshWindowInputsBatch,
41+
ReplicaShape, StateWrite,
4042
};
4143
use mz_compute_types::config::ComputeReplicaConfig;
4244
use mz_controller::clusters::ClusterStatus;
@@ -54,8 +56,9 @@ use crate::error::AdapterError;
5456
/// [`ClusterControllerCtx`] call. Each variant carries a oneshot for the reply.
5557
///
5658
/// `ManagedClusterIds` and `ClusterStates` are the per-tick batched reads. The
57-
/// `ClusterStates` reply also carries `now`. `HydratedReplicas` is a
58-
/// per-cluster live signal a strategy pulls on demand.
59+
/// `ClusterStates` reply also carries `now`. Refresh-window catalog inputs are
60+
/// pulled one cluster at a time, followed by one shared oracle read.
61+
/// `HydratedReplicas` is a per-cluster live signal a strategy pulls on demand.
5962
#[derive(Debug)]
6063
pub enum ClusterControllerRequest {
6164
/// The ids of all *user* managed clusters the controller owns this tick.
@@ -81,13 +84,14 @@ pub enum ClusterControllerRequest {
8184
cluster_id: ClusterId,
8285
tx: oneshot::Sender<bool>,
8386
},
84-
/// The refresh-window live signals for one scheduled cluster (read ts,
85-
/// compaction estimate, bound REFRESH MVs). `None` for a cluster that is not
86-
/// scheduled `ON REFRESH`.
87-
RefreshWindowInputs {
87+
/// The catalog and storage refresh-window inputs for one scheduled cluster.
88+
/// `None` if the cluster no longer qualifies at pull time.
89+
RefreshWindowClusterInputs {
8890
cluster_id: ClusterId,
89-
tx: oneshot::Sender<Option<RefreshWindowInputs>>,
91+
tx: oneshot::Sender<Option<RefreshWindowClusterInputs>>,
9092
},
93+
/// One timestamp-oracle read for a completed refresh-window input batch.
94+
RefreshWindowReadTs { tx: oneshot::Sender<Timestamp> },
9195
/// Apply a tick's batch of decisions under their compare-and-append guards.
9296
Apply {
9397
decisions: Vec<Decision>,
@@ -184,11 +188,38 @@ impl ClusterControllerCtx for CoordCtx {
184188

185189
async fn refresh_window_inputs(
186190
&mut self,
187-
cluster_id: ClusterId,
188-
) -> Option<RefreshWindowInputs> {
189-
self.request(|tx| ClusterControllerRequest::RefreshWindowInputs { cluster_id, tx })
190-
.await
191-
.flatten()
191+
cluster_ids: &[ClusterId],
192+
) -> Option<RefreshWindowInputsBatch> {
193+
let mut cluster_inputs = BTreeMap::new();
194+
for (index, &cluster_id) in cluster_ids.iter().enumerate() {
195+
if index > 0 {
196+
// The coordinator prioritizes its internal command channel. Give
197+
// it a chance to service already-queued user commands instead of
198+
// keeping that channel continuously ready for the whole batch.
199+
tokio::task::yield_now().await;
200+
}
201+
let inputs = self
202+
.request(|tx| ClusterControllerRequest::RefreshWindowClusterInputs {
203+
cluster_id,
204+
tx,
205+
})
206+
.await
207+
.flatten();
208+
if let Some(inputs) = inputs {
209+
cluster_inputs.insert(cluster_id, inputs);
210+
}
211+
}
212+
if cluster_inputs.is_empty() {
213+
return None;
214+
}
215+
216+
let read_ts = self
217+
.request(|tx| ClusterControllerRequest::RefreshWindowReadTs { tx })
218+
.await?;
219+
Some(RefreshWindowInputsBatch {
220+
read_ts,
221+
cluster_inputs,
222+
})
192223
}
193224

194225
async fn apply(&mut self, decisions: Vec<Decision>) -> ApplyOutcome {
@@ -326,36 +357,18 @@ impl Coordinator {
326357
ClusterControllerRequest::HasHydratableObjects { cluster_id, tx } => {
327358
let _ = tx.send(self.cluster_has_hydratable_objects(cluster_id));
328359
}
329-
ClusterControllerRequest::RefreshWindowInputs { cluster_id, tx } => {
330-
// Gather the catalog- and storage-derived inputs on the loop,
331-
// then complete the reply from a spawned task: the oracle
332-
// read is a network round-trip (to the Postgres/CRDB-backed
333-
// timestamp oracle) and must never run on the serial
334-
// coordinator loop. The legacy `check_refresh_policy` makes
335-
// the same split.
336-
match self.refresh_window_catalog_inputs(cluster_id) {
337-
None => {
338-
let _ = tx.send(None);
339-
}
340-
Some((compaction_estimate, refresh_mvs)) => {
341-
let oracle = self.get_local_timestamp_oracle();
342-
// NOTE: this is one oracle read per scheduled cluster
343-
// per tick, and the controller awaits each pull before
344-
// the next, so the reads are sequential and the
345-
// batching oracle cannot coalesce them. Fine at the
346-
// tick cadence for realistic scheduled-cluster counts.
347-
// TODO: hoist to one read per tick if that stops
348-
// holding.
349-
spawn(|| "cluster_controller_refresh_window_read_ts", async move {
350-
let read_ts = oracle.read_ts().await;
351-
let _ = tx.send(Some(RefreshWindowInputs {
352-
read_ts,
353-
compaction_estimate,
354-
refresh_mvs,
355-
}));
356-
});
357-
}
358-
}
360+
ClusterControllerRequest::RefreshWindowClusterInputs { cluster_id, tx } => {
361+
let _ = tx.send(self.refresh_window_catalog_inputs(cluster_id));
362+
}
363+
ClusterControllerRequest::RefreshWindowReadTs { tx } => {
364+
// The oracle read is a network round trip to the Postgres or
365+
// CRDB-backed timestamp oracle. It must not run on the serial
366+
// coordinator loop.
367+
let oracle = self.get_local_timestamp_oracle();
368+
spawn(|| "cluster_controller_refresh_window_read_ts", async move {
369+
let read_ts = oracle.read_ts().await;
370+
let _ = tx.send(read_ts);
371+
});
359372
}
360373
ClusterControllerRequest::Apply { decisions, tx } => {
361374
let outcome = if active {
@@ -514,7 +527,7 @@ impl Coordinator {
514527
/// REFRESH`. These are the same signals the legacy `check_refresh_policy`
515528
/// reads.
516529
///
517-
/// The oracle read timestamp completing [`RefreshWindowInputs`] is
530+
/// The oracle read timestamp completing [`RefreshWindowInputsBatch`] is
518531
/// deliberately not fetched here: this runs on the coordinator loop, and
519532
/// the oracle read is a network round-trip the request handler performs on
520533
/// a spawned task instead.
@@ -525,7 +538,7 @@ impl Coordinator {
525538
fn refresh_window_catalog_inputs(
526539
&self,
527540
cluster_id: ClusterId,
528-
) -> Option<(Duration, Vec<RefreshMvInfo>)> {
541+
) -> Option<RefreshWindowClusterInputs> {
529542
use mz_catalog::memory::objects::CatalogItem;
530543

531544
let cluster = self.catalog().try_get_cluster(cluster_id)?;
@@ -568,7 +581,10 @@ impl Coordinator {
568581
.system_config()
569582
.cluster_refresh_mv_compaction_estimate();
570583

571-
Some((compaction_estimate, refresh_mvs))
584+
Some(RefreshWindowClusterInputs {
585+
compaction_estimate,
586+
refresh_mvs,
587+
})
572588
}
573589

574590
/// Apply one batch of decisions under their compare-and-append guards.

‎src/cluster-controller/src/ctx.rs‎

Lines changed: 43 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@
2424
//! controller drives what is fetched. Read methods are batched so a separate-task
2525
//! deployment can bound its round-trips to the Coordinator.
2626
27-
use std::collections::BTreeSet;
27+
use std::collections::{BTreeMap, BTreeSet};
2828
use std::time::Duration;
2929

3030
use async_trait::async_trait;
@@ -176,16 +176,36 @@ impl RefreshWindowDecision {
176176
}
177177
}
178178

179-
/// The live signals the on-refresh strategy reads to decide whether a scheduled
180-
/// cluster is inside a refresh window: the current read timestamp, the
181-
/// Persist-compaction time estimate, and the bound REFRESH MVs' frontiers and
182-
/// schedules.
179+
/// The catalog and storage inputs for one scheduled cluster's refresh window.
180+
#[derive(Clone, Debug, PartialEq, Eq)]
181+
pub struct RefreshWindowClusterInputs {
182+
/// How long after a refresh an MV is estimated to still need Persist
183+
/// compaction, which also keeps the cluster on.
184+
pub compaction_estimate: Duration,
185+
/// The REFRESH MVs bound to the cluster.
186+
pub refresh_mvs: Vec<RefreshMvInfo>,
187+
}
188+
189+
/// Refresh-window inputs gathered for one reconciliation phase.
183190
///
184-
/// Pulled on demand only for scheduled clusters. A MANUAL cluster carries `None`
185-
/// and is never probed.
191+
/// The top-level timestamp makes sharing one oracle read across every included
192+
/// cluster structural. A cluster absent from `cluster_inputs` has unavailable
193+
/// inputs and must not be reconciled during the phase.
194+
#[derive(Clone, Debug, PartialEq, Eq)]
195+
pub struct RefreshWindowInputsBatch {
196+
/// The local oracle read timestamp for every cluster in the batch.
197+
pub read_ts: Timestamp,
198+
/// The available catalog and storage inputs, keyed by cluster.
199+
pub cluster_inputs: BTreeMap<ClusterId, RefreshWindowClusterInputs>,
200+
}
201+
202+
/// The fulfilled live signal the on-refresh strategy uses for one cluster.
203+
///
204+
/// Pulled on demand only for scheduled clusters. A MANUAL cluster carries
205+
/// `None` and is never probed.
186206
#[derive(Clone, Debug, PartialEq, Eq)]
187207
pub struct RefreshWindowInputs {
188-
/// The local oracle read timestamp the window decision is taken against.
208+
/// The shared local oracle read timestamp the window decision uses.
189209
pub read_ts: Timestamp,
190210
/// How long after a refresh an MV is estimated to still need Persist
191211
/// compaction, which also keeps the cluster on.
@@ -417,21 +437,24 @@ pub trait ClusterControllerCtx: Send {
417437
/// burst winds down via its linger.
418438
async fn has_hydratable_objects(&mut self, cluster_id: ClusterId) -> bool;
419439

420-
/// The refresh-window live signals for one scheduled cluster: the read
421-
/// timestamp, the compaction estimate, and the bound REFRESH MVs' write
422-
/// frontiers and schedules. Returns `None` when the cluster is missing,
423-
/// unmanaged, or no longer scheduled `ON REFRESH` at pull time. The
424-
/// controller only asks about clusters it observed as scheduled, so `None`
425-
/// means a concurrent DDL moved the cluster mid-tick, and the schedule's
426-
/// membership in the compare-and-append witness rejects any decision
427-
/// derived from the stale observation.
440+
/// The refresh-window live signals for the given scheduled clusters.
441+
/// Returns one shared read timestamp plus the available per-cluster catalog
442+
/// and storage inputs. Omits a cluster when its inputs are unavailable,
443+
/// including when it is missing, unmanaged, or no longer scheduled `ON
444+
/// REFRESH` at pull time. Returns `None` when the batch fails or no requested
445+
/// cluster has valid inputs.
428446
///
429447
/// Pulled on demand the same way as [`Self::hydrated_replicas`]: the
430-
/// controller probes a cluster only when the on-refresh strategy needs the
448+
/// controller includes a cluster only when the on-refresh strategy needs the
431449
/// signal (i.e. the cluster is scheduled), so a steady MANUAL cluster never
432-
/// pays for it.
433-
async fn refresh_window_inputs(&mut self, cluster_id: ClusterId)
434-
-> Option<RefreshWindowInputs>;
450+
/// pays for it. Implementations fetch the shared read timestamp once after
451+
/// gathering the per-cluster inputs, so oracle latency does not scale with
452+
/// the cluster count. The controller skips any omitted cluster for the
453+
/// reconciliation phase.
454+
async fn refresh_window_inputs(
455+
&mut self,
456+
cluster_ids: &[ClusterId],
457+
) -> Option<RefreshWindowInputsBatch>;
435458

436459
/// Apply a tick's batch of decisions under their compare-and-append guards.
437460
/// Each decision carries the [`ExpectedClusterState`] it was derived from;

‎src/cluster-controller/src/lib.rs‎

Lines changed: 44 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ use mz_ore::soft_panic_or_log;
4646
use crate::ctx::{
4747
ApplyOutcome, ClusterControllerCtx, ClusterState, CreateReason, Decision, ObservedReplica,
4848
ReconfigurationAudit, ReconfigurationRecord, ReconfigurationStatus, ReconfigurationWrite,
49-
ReplicaShape, StateWrite,
49+
RefreshWindowInputs, ReplicaShape, StateWrite,
5050
};
5151
use crate::strategy::{
5252
BaselineStrategy, ConfigSignals, DesiredReplica, GracefulReconfigurationStrategy,
@@ -131,7 +131,10 @@ impl ClusterController {
131131
// that is probably about to go stale.
132132
let mut rejected = BTreeSet::new();
133133
for state in &states {
134-
let write = self.merge_state_writes(state, &signals[&state.cluster_id], &config, now);
134+
let Some(signals) = signals.get(&state.cluster_id) else {
135+
continue;
136+
};
137+
let write = self.merge_state_writes(state, signals, &config, now);
135138
if write.is_empty() {
136139
continue;
137140
}
@@ -166,8 +169,10 @@ impl ClusterController {
166169
if rejected.contains(&state.cluster_id) {
167170
continue;
168171
}
169-
let decisions =
170-
self.collect_replica_decisions(state, &signals[&state.cluster_id], &config, now);
172+
let Some(signals) = signals.get(&state.cluster_id) else {
173+
continue;
174+
};
175+
let decisions = self.collect_replica_decisions(state, signals, &config, now);
171176
if decisions.is_empty() {
172177
continue;
173178
}
@@ -325,17 +330,20 @@ impl ClusterController {
325330
///
326331
/// Each strategy names its needs as a pure function of the durable state
327332
/// and the tick's config signals ([`Strategy::signal_request`]), so the
328-
/// kernel stays ignorant of when a strategy engages. Signals are fetched per
329-
/// cluster and only where requested: a steady cluster is never probed,
330-
/// keeping the ctx seam pay-for-what-you-use. The returned map has an entry
331-
/// for every state.
333+
/// kernel stays ignorant of when a strategy engages. Signals are fetched
334+
/// only where requested: a steady cluster is never probed, keeping the ctx
335+
/// seam pay-for-what-you-use. Refresh-window inputs are fetched as one batch
336+
/// so every scheduled cluster shares one oracle read per phase. The returned
337+
/// map omits a state when one of its required inputs was unavailable, which
338+
/// causes the reconciliation phase to skip that cluster.
332339
async fn fetch_signals(
333340
&self,
334341
ctx: &mut dyn ClusterControllerCtx,
335342
states: &[ClusterState],
336343
config: &ConfigSignals,
337344
) -> BTreeMap<ClusterId, LiveSignals> {
338345
let mut signals = BTreeMap::new();
346+
let mut refresh_window_clusters = Vec::new();
339347
for state in states {
340348
let request = self
341349
.strategies
@@ -360,10 +368,37 @@ impl ClusterController {
360368
}
361369
}
362370
if request.refresh_window {
363-
live.refresh_window = ctx.refresh_window_inputs(state.cluster_id).await;
371+
refresh_window_clusters.push(state.cluster_id);
364372
}
365373
signals.insert(state.cluster_id, live);
366374
}
375+
if !refresh_window_clusters.is_empty() {
376+
match ctx.refresh_window_inputs(&refresh_window_clusters).await {
377+
Some(batch) => {
378+
let read_ts = batch.read_ts;
379+
let mut cluster_inputs = batch.cluster_inputs;
380+
for cluster_id in refresh_window_clusters {
381+
let Some(inputs) = cluster_inputs.remove(&cluster_id) else {
382+
signals.remove(&cluster_id);
383+
continue;
384+
};
385+
let live = signals
386+
.get_mut(&cluster_id)
387+
.expect("signal entry inserted for requested cluster");
388+
live.refresh_window = Some(RefreshWindowInputs {
389+
read_ts,
390+
compaction_estimate: inputs.compaction_estimate,
391+
refresh_mvs: inputs.refresh_mvs,
392+
});
393+
}
394+
}
395+
None => {
396+
for cluster_id in refresh_window_clusters {
397+
signals.remove(&cluster_id);
398+
}
399+
}
400+
}
401+
}
367402
signals
368403
}
369404

0 commit comments

Comments
 (0)