Skip to content

Commit 9a5f824

Browse files
committed
storage/upsert: select the stash representation by dyncfg
enable_upsert_chunked_stash (off in production, on and randomized in CI) selects between two instantiations of the upsert-v2 operator loop: * Chunked (flag on): ChunkBatcher stash, chunk-spine feedback arrangement, bulk-probe drain. Spills committed chunk bodies through the process buffer pool. * Paged (flag off, the default): paged columnar merge batcher stash, ValRowSpine feedback arrangement, cursor drain. Spills cold chains through the storage-owned column pager. The operator loop itself is shared: build_upsert_operator is generic over an UpsertStashArm, whose two impls carry the flavor-specific pieces (stash batcher, feedback spine, flush, drain). UpsertStashFlavor names the arms and resolves from the config set once, at operator construction, mirroring compute's ArrangementBatcher. Both flavors' spill paths stay gated by enable_upsert_paged_spill, so storage_state applies that flag to both the chunk gate and the restored upsert_stash_pager. The unit-test harness runs every scenario under both flavors and asserts they produce identical output.
1 parent 25e43d1 commit 9a5f824

8 files changed

Lines changed: 812 additions & 213 deletions

File tree

‎misc/python/materialize/mzcompose/__init__.py‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -308,6 +308,14 @@ def get_variable_system_parameters(
308308
"true",
309309
["true", "false"],
310310
),
311+
# On by default so CI exercises the chunked stash flavor, which is
312+
# off in production while it earns trust. Only meaningful when
313+
# enable_upsert_v2 is true.
314+
VariableSystemParameter(
315+
"enable_upsert_chunked_stash",
316+
"true",
317+
["true", "false"],
318+
),
311319
VariableSystemParameter(
312320
"enable_upsert_v2",
313321
"false",

‎misc/python/materialize/parallel_workload/action.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3055,6 +3055,7 @@ def __init__(
30553055
"0.02",
30563056
]
30573057
self.flags_with_values["enable_upsert_paged_spill"] = BOOLEAN_FLAG_VALUES
3058+
self.flags_with_values["enable_upsert_chunked_stash"] = BOOLEAN_FLAG_VALUES
30583059
self.flags_with_values["column_chunk_compress_min_depth"] = [
30593060
"0", # compress every spilled body
30603061
"1", # the default: fresh chunks store uncompressed

‎src/storage-types/src/dyncfgs.rs‎

Lines changed: 43 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -392,21 +392,52 @@ pub const STORAGE_UPSERT_MAX_SNAPSHOT_BATCH_BUFFERING: Config<Option<usize>> = C
392392
ParameterScope::Replica,
393393
);
394394

395-
/// Storage's leg of the process-wide chunk spill gate
396-
/// (`mz_timely_util::columnar::chunk`). The gate is the OR of a compute leg
397-
/// (`enable_column_paged_batcher_spill`) and this storage leg: chunks spill
398-
/// while either is set, so this flag cannot veto spilling that the compute
399-
/// flag has enabled. Spilled chunks draw on the one shared pool budget.
395+
/// Allow the upsert-v2 stash to spill out of RSS. Off by default; while off,
396+
/// the stash keeps everything resident.
400397
///
401-
/// Off by default. Enabling it also installs the process buffer pool (via
402-
/// compute's config handler, which reads this flag from the aggregate dyncfg
403-
/// set), so storage-only spilling needs no compute-side gate.
398+
/// The spill mechanism depends on the stash flavor
399+
/// ([`ENABLE_UPSERT_CHUNKED_STASH`]):
400+
///
401+
/// * Chunked: sets storage's leg of the process-wide chunk spill gate
402+
/// (`mz_timely_util::columnar::chunk`). The gate is the OR of a compute
403+
/// leg (`enable_column_paged_batcher_spill`) and this storage leg: chunks
404+
/// spill while either is set, so this flag cannot veto spilling that the
405+
/// compute flag has enabled. Spilled chunks draw on the one shared pool
406+
/// budget, and the gate is consulted at every chunk commit, so flips
407+
/// apply to running dataflows.
408+
/// * Paged: gates the storage-owned column pager the stash and feedback
409+
/// arrangement route their chains through, independently of compute's
410+
/// `enable_column_paged_batcher_spill`. Captured at operator
411+
/// construction, so flips apply to dataflows created after the change.
412+
///
413+
/// Enabling it also installs the process buffer pool (via compute's config
414+
/// handler, which reads this flag from the aggregate dyncfg set), so
415+
/// storage-only spilling needs no compute-side gate.
404416
pub const ENABLE_UPSERT_PAGED_SPILL: Config<bool> = Config::new(
405417
"enable_upsert_paged_spill",
406418
false,
407-
"Allow upsert-v2 chunks to spill to the shared buffer pool. Sets the storage leg of the \
408-
process-wide spill gate, which is the OR of this flag and the compute \
409-
`enable_column_paged_batcher_spill`.",
419+
"Allow the upsert-v2 stash to spill out of RSS, through the buffer pool (chunked stash \
420+
flavor) or the column pager (paged stash flavor).",
421+
ParameterScope::Replica,
422+
);
423+
424+
/// Use the chunked stash flavor for the upsert-v2 operator: differential's
425+
/// chunk merge batcher for the source stash and a spine of chunk batches for
426+
/// the feedback arrangement, with a bulk-probe drain. When `false` (the
427+
/// default), the paged flavor is used: the paged columnar merge batcher and a
428+
/// `ValRowSpine`, with a cursor-based drain. See
429+
/// `mz_storage::upsert_continual_feedback_v2::UpsertStashFlavor` for the
430+
/// comparison.
431+
///
432+
/// Read at operator construction time; flips take effect on dataflows created
433+
/// after the change. Only meaningful when [`ENABLE_UPSERT_V2`] is `true`.
434+
/// Spilling in either flavor is gated by [`ENABLE_UPSERT_PAGED_SPILL`].
435+
pub const ENABLE_UPSERT_CHUNKED_STASH: Config<bool> = Config::new(
436+
"enable_upsert_chunked_stash",
437+
false,
438+
"Use the chunk batcher and chunk spine for the upsert-v2 stash and feedback arrangement, \
439+
instead of the paged columnar merge batcher and ValRowSpine. Only meaningful when \
440+
enable_upsert_v2 is true.",
410441
ParameterScope::Replica,
411442
);
412443

@@ -533,6 +564,7 @@ pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet {
533564
.add(&ENABLE_UPSERT_V2)
534565
.add(&SUSPENDABLE_SOURCES)
535566
.add(&ENABLE_UPSERT_PAGED_SPILL)
567+
.add(&ENABLE_UPSERT_CHUNKED_STASH)
536568
.add(&WALLCLOCK_GLOBAL_LAG_HISTOGRAM_RETENTION_INTERVAL)
537569
.add(&WALLCLOCK_LAG_HISTORY_RETENTION_INTERVAL)
538570
.add(&crate::sources::sql_server::CDC_CLEANUP_CHANGE_TABLE)

‎src/storage/src/render/sources.rs‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -339,6 +339,13 @@ where
339339
if dyncfgs::ENABLE_UPSERT_V2
340340
.get(storage_state.storage_configuration.config_set())
341341
{
342+
// Resolved here, at operator construction, so the
343+
// dataflow keeps one stash flavor for its whole
344+
// life even if the flag flips underneath it.
345+
let stash_flavor =
346+
crate::upsert_continual_feedback_v2::UpsertStashFlavor::from_config(
347+
storage_state.storage_configuration.config_set(),
348+
);
342349
crate::upsert::upsert_v2(
343350
upsert_input.enter(scope),
344351
upsert_envelope.clone(),
@@ -347,6 +354,7 @@ where
347354
previous_token,
348355
export_config,
349356
backpressure_metrics,
357+
stash_flavor,
350358
)
351359
} else {
352360
crate::upsert::upsert(

‎src/storage/src/storage_state.rs‎

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -846,12 +846,14 @@ impl<'w> Worker<'w> {
846846
STORAGE_SERVER_MAINTENANCE_INTERVAL
847847
.get(self.storage_state.storage_configuration.config_set());
848848

849-
// Set storage's leg of the process-wide chunk spill gate.
850-
// The buffer pool and its budget are the shared ones
851-
// configured by compute's `apply_worker_config` (compute and
852-
// storage run in the same process). The gate ORs this leg
853-
// with compute's, so chunks spill while either subsystem's
854-
// flag is set.
849+
// Apply storage's upsert spill flag to both stash flavors'
850+
// mechanisms: the storage leg of the process-wide chunk
851+
// spill gate (chunked flavor) and the storage-owned column
852+
// pager (paged flavor). The buffer pool, the pager pool, and
853+
// their budgets are the shared ones configured by compute's
854+
// `apply_worker_config` (compute and storage run in the same
855+
// process). The chunk gate ORs storage's leg with compute's,
856+
// so chunks spill while either subsystem's flag is set.
855857
//
856858
// The flag is replica-scoped: the storage controller merges
857859
// per-replica overrides into the `UpdateConfiguration`
@@ -867,6 +869,7 @@ impl<'w> Worker<'w> {
867869
enabled, "upsert stash spill: applying gate",
868870
);
869871
crate::upsert::upsert_stash_spill::set_enabled(enabled);
872+
crate::upsert::upsert_stash_pager::set_enabled(enabled);
870873
}
871874
}
872875
InternalStorageCommand::StatisticsUpdate { sources, sinks } => self

‎src/storage/src/upsert.rs‎

Lines changed: 43 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -296,9 +296,10 @@ macro_rules! upsert_source_time_unit {
296296
}
297297
upsert_source_time_unit!(GtidPartition, Lsn);
298298

299-
/// Storage's leg of the process-wide chunk spill gate.
299+
/// Storage's leg of the process-wide chunk spill gate, used by the chunked
300+
/// upsert-v2 stash flavor.
300301
///
301-
/// The upsert-v2 source stash and feedback arrangement spill through the
302+
/// In that flavor the source stash and feedback arrangement spill through the
302303
/// process buffer pool ([`mz_timely_util::columnar::chunk`]): committed chunk
303304
/// bodies land in the pool once compute's config handler has installed and
304305
/// budgeted it (storage and compute run in the same `clusterd` process).
@@ -315,6 +316,43 @@ pub mod upsert_stash_spill {
315316
}
316317
}
317318

319+
/// Pager for the paged upsert-v2 stash flavor.
320+
///
321+
/// This draws from the same process-wide [`TieredPolicy`] budget pool as the
322+
/// compute column-paged batcher — there is one budget and one underlying
323+
/// `mz_ore::pager` — but whether the stash *uses* it is gated by storage's own
324+
/// `enable_upsert_paged_spill` flag, independently of compute's
325+
/// `enable_column_paged_batcher_spill`. The shared pool's budget / backend /
326+
/// codec are configured by compute's `apply_tiered_config` (storage and compute
327+
/// run in the same `clusterd` process).
328+
///
329+
/// Flipping the flag takes effect on dataflows created after the change: the
330+
/// paged flavor captures the pager once at operator construction.
331+
///
332+
/// [`TieredPolicy`]: mz_timely_util::column_pager::policy::TieredPolicy
333+
pub mod upsert_stash_pager {
334+
use std::sync::{LazyLock, RwLock};
335+
336+
use mz_timely_util::column_pager::{ColumnPager, shared_pager};
337+
338+
/// Active pager handed to upsert source-stash batchers. Defaults to
339+
/// disabled (every chunk resident) until [`set_enabled`] turns it on.
340+
static PAGER: LazyLock<RwLock<ColumnPager>> =
341+
LazyLock::new(|| RwLock::new(ColumnPager::disabled()));
342+
343+
/// Enable or disable the stash's use of the shared column pager. When
344+
/// enabled, the stash spills through the shared budget pool; when disabled
345+
/// it keeps every chunk resident.
346+
pub fn set_enabled(enabled: bool) {
347+
*PAGER.write().expect("upsert stash pager poisoned") = shared_pager(enabled);
348+
}
349+
350+
/// The current upsert-stash pager. Cheap: clones the inner `Arc`.
351+
pub fn pager() -> ColumnPager {
352+
PAGER.read().expect("upsert stash pager poisoned").clone()
353+
}
354+
}
355+
318356
impl Debug for UpsertKey {
319357
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
320358
write!(f, "0x")?;
@@ -604,6 +642,7 @@ pub(crate) fn upsert_v2<'scope, T, FromTime>(
604642
previous_token: Option<Vec<PressOnDropButton>>,
605643
source_config: crate::source::SourceExportCreationConfig,
606644
backpressure_metrics: Option<BackpressureMetrics>,
645+
stash_flavor: upsert_continual_feedback_v2::UpsertStashFlavor,
607646
) -> (
608647
VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
609648
StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
@@ -630,10 +669,12 @@ where
630669
tracing::info!(
631670
worker_id = %source_config.worker_id,
632671
source_id = %source_config.id,
672+
?stash_flavor,
633673
"rendering upsert source (btreemap backend)"
634674
);
635675

636676
upsert_continual_feedback_v2::upsert_inner(
677+
stash_flavor,
637678
thin_input,
638679
upsert_envelope.key_indices,
639680
resume_upper,

0 commit comments

Comments
 (0)