Skip to content

Commit 90dc1b3

Browse files
storage: budget MySQL snapshot partition probing by table size
Bound the prefix partitioner to a per-table probe query budget scaled to the estimated row count, 2500 requests per billion rows by default via the mysql_source_snapshot_partition_requests_per_billion_rows dyncfg, with a floor of 256 so small tables can afford their handful of splits. Probes run sequentially, so this caps the wall-clock that partitioning can add to snapshot start in proportion to the snapshot work it optimizes. An exhausted budget stops splitting early and leaves coarser buckets, never incorrect ones. Adds a mock test asserting the partitioner never exceeds its budget and still returns a valid ordered boundary list when truncated.
1 parent 3cf0f01 commit 90dc1b3

4 files changed

Lines changed: 99 additions & 22 deletions

File tree

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -701,6 +701,9 @@ def get_default_system_parameters(
701701
# The estimated path is covered explicitly in mysql-cdc/statistics.td and
702702
# by parallel-workload.
703703
"mysql_source_snapshot_exact_count_max_rows",
704+
# Not varied here because the 256-request floor dominates for test-sized
705+
# tables. parallel-workload flips it.
706+
"mysql_source_snapshot_partition_requests_per_billion_rows",
704707
"postgres_fetch_slot_resume_lsn_interval",
705708
"pg_schema_validation_interval",
706709
"pg_source_validate_timeline",

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

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3046,6 +3046,14 @@ def __init__(
30463046
self.flags_with_values["mysql_source_snapshot_parallelism"] = (
30473047
BOOLEAN_FLAG_VALUES
30483048
)
3049+
# 0 leaves only the 256-request floor, the default scales with table
3050+
# size.
3051+
self.flags_with_values[
3052+
"mysql_source_snapshot_partition_requests_per_billion_rows"
3053+
] = [
3054+
"0",
3055+
"2500",
3056+
]
30493057

30503058
# If you are adding a new config flag in Materialize, consider using it
30513059
# here instead of just marking it as uninteresting to silence the

‎src/mysql-util/src/partition.rs‎

Lines changed: 75 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -16,13 +16,18 @@ use crate::{KeyProber, MySqlError, QualifiedTableRef};
1616
/// into `num_workers` roughly even partitions.
1717
/// This should be run in a repeatable read transaction against a primary key varchar/char column
1818
/// with the `utf8mb4_bin` collation.
19+
///
20+
/// At most `max_requests` probes are issued against the server. When the
21+
/// budget runs out, remaining prefixes stay unsplit, which skews partition
22+
/// sizes but never correctness.
1923
pub async fn partition_table(
2024
conn: &mut mysql_async::Conn,
2125
table: QualifiedTableRef<'_>,
2226
pk_col: &str,
2327
num_workers: usize,
2428
estimated_row_count: u64,
2529
min_rows_per_worker: u64,
30+
max_requests: u64,
2631
) -> Result<Vec<String>, MySqlError> {
2732
let (schema_name, table_name) = (table.schema_name, table.table_name);
2833
let mut db = KeyProber::new(conn, table, pk_col);
@@ -31,6 +36,7 @@ pub async fn partition_table(
3136
num_workers,
3237
estimated_row_count,
3338
min_rows_per_worker,
39+
max_requests,
3440
)
3541
.await?;
3642
tracing::trace!(
@@ -60,6 +66,7 @@ async fn partition<D: PrimaryKeyProber>(
6066
workers: usize,
6167
estimated_row_count: u64,
6268
min_rows_per_worker: u64,
69+
max_requests: u64,
6370
) -> Result<Vec<String>, MySqlError> {
6471
if workers <= 1 {
6572
return Ok(Vec::new());
@@ -75,15 +82,24 @@ async fn partition<D: PrimaryKeyProber>(
7582
.max(f64::cast_lossy(min_rows_per_worker))
7683
.max(1.0);
7784

78-
compute_boundaries(db, workers, estimated_row_count, target_max_rows_per_prefix).await
85+
compute_boundaries(
86+
db,
87+
workers,
88+
estimated_row_count,
89+
target_max_rows_per_prefix,
90+
max_requests,
91+
)
92+
.await
7993
}
8094

8195
async fn compute_boundaries<D: PrimaryKeyProber>(
8296
db: &mut D,
8397
workers: usize,
8498
estimated_row_count: u64,
8599
target_rows_per_prefix: f64,
100+
max_requests: u64,
86101
) -> Result<Vec<String>, MySqlError> {
102+
let mut budget = max_requests;
87103
// BFS of prefixes, splitting until estimates fall under the target.
88104
let mut ordered_prefixes = vec![Prefix {
89105
prefix: String::new(),
@@ -96,11 +112,14 @@ async fn compute_boundaries<D: PrimaryKeyProber>(
96112
let mut next_ordered_prefixes: Vec<Prefix> = vec![];
97113
let mut split_any = false;
98114
for prefix in ordered_prefixes {
99-
if f64::cast_lossy(prefix.estimated_rows) > target_rows_per_prefix {
115+
// Entering a split costs one probe for the first prefix plus two
116+
// for its first walk step, a smaller budget keeps the prefix as a
117+
// leaf.
118+
if budget >= 3 && f64::cast_lossy(prefix.estimated_rows) > target_rows_per_prefix {
100119
split_any = true;
101120
// Partitioning children can drop some rows from the parent prefix range. This
102121
// is acceptable given the approximate nature of the algorithm.
103-
let children = children_prefixes(db, &prefix).await?;
122+
let children = children_prefixes(db, &prefix, &mut budget).await?;
104123
next_ordered_prefixes.extend(children);
105124
} else {
106125
next_ordered_prefixes.push(prefix);
@@ -111,6 +130,11 @@ async fn compute_boundaries<D: PrimaryKeyProber>(
111130
break;
112131
}
113132
}
133+
tracing::debug!(
134+
prefixes = ordered_prefixes.len(),
135+
requests_spent = max_requests - budget,
136+
"split key space into prefixes"
137+
);
114138

115139
// Recompute the total after partitioning the table to get more even splits because the actual row count and the
116140
// granularly estimated row count can diverge from the original top level estimate.
@@ -147,13 +171,19 @@ async fn compute_boundaries<D: PrimaryKeyProber>(
147171
///
148172
/// Note: This will drop the key "a" on the floor, along with any keys
149173
/// sorting below their own prefix (below-space characters at this depth).
174+
///
175+
/// `budget` is decremented once per probe. The caller must provide at least
176+
/// 3, one for the first prefix and two for a walk step. The walk closes out
177+
/// with a tail child once it cannot afford another step.
150178
async fn children_prefixes<D: PrimaryKeyProber>(
151179
db: &mut D,
152180
parent: &Prefix,
181+
budget: &mut u64,
153182
) -> Result<Vec<Prefix>, MySqlError> {
154183
let depth = parent.depth + 1;
155184
let mut children = Vec::new();
156185

186+
*budget -= 1;
157187
let Some(mut cur) = db
158188
.prefix_of_first_key_in_range(&parent.prefix, parent.end.as_deref(), depth)
159189
.await?
@@ -162,10 +192,25 @@ async fn children_prefixes<D: PrimaryKeyProber>(
162192
};
163193

164194
loop {
195+
// A walk step costs two probes. When they are unaffordable, close out
196+
// with a tail child so the parent's key space stays covered,
197+
// estimated as the parent's unconsumed mass.
198+
if *budget < 2 {
199+
let consumed: u64 = children.iter().map(|c| c.estimated_rows).sum();
200+
children.push(Prefix {
201+
prefix: cur,
202+
end: parent.end.clone(),
203+
estimated_rows: parent.estimated_rows.saturating_sub(consumed).max(1),
204+
depth,
205+
});
206+
return Ok(children);
207+
}
208+
*budget -= 1;
165209
let next = db
166210
.prefix_of_first_row_not_matching_prefix(&cur, parent.end.as_deref(), depth)
167211
.await?;
168212
let end = next.clone().or_else(|| parent.end.clone());
213+
*budget -= 1;
169214
let estimated_rows = db.estimate_range_rows(&cur, end.as_deref()).await?;
170215
children.push(Prefix {
171216
prefix: cur,
@@ -247,9 +292,15 @@ mod tests {
247292
/// cover them.
248293
struct MockDb {
249294
keys: Vec<String>,
295+
/// Probes served, for asserting on the request budget.
296+
requests: u64,
250297
}
251298

252299
impl MockDb {
300+
fn new(keys: Vec<String>) -> Self {
301+
MockDb { keys, requests: 0 }
302+
}
303+
253304
fn bounds(&self, start: &str, end: Option<&str>) -> (usize, usize) {
254305
// The lower bound is exclusive, a key equal to `start` is skipped.
255306
let lo = self.keys.partition_point(|k| k.as_str() <= start);
@@ -267,6 +318,7 @@ mod tests {
267318
start: &str,
268319
end: Option<&str>,
269320
) -> Result<u64, MySqlError> {
321+
self.requests += 1;
270322
let (lo, hi) = self.bounds(start, end);
271323
Ok(u64::cast_from(hi - lo))
272324
}
@@ -277,6 +329,7 @@ mod tests {
277329
end: Option<&str>,
278330
len: usize,
279331
) -> Result<Option<String>, MySqlError> {
332+
self.requests += 1;
280333
let (lo, hi) = self.bounds(start, end);
281334
if lo >= hi {
282335
return Ok(None);
@@ -290,6 +343,7 @@ mod tests {
290343
end: Option<&str>,
291344
len: usize,
292345
) -> Result<Option<String>, MySqlError> {
346+
self.requests += 1;
293347
let (_, hi) = self.bounds("", end);
294348
// Find the last key matching `cur`, byte prefixes stand in for
295349
// the collation's LIKE matching.
@@ -310,9 +364,9 @@ mod tests {
310364

311365
#[mz_ore::test(tokio::test)]
312366
async fn single_worker_gets_no_boundaries() -> Result<(), MySqlError> {
313-
let mut db = MockDb { keys: keys(1000) };
367+
let mut db = MockDb::new(keys(1000));
314368
let count = u64::cast_from(db.keys.len());
315-
let boundaries = partition(&mut db, 1, count, MIN_ROWS_PER_WORKER).await?;
369+
let boundaries = partition(&mut db, 1, count, MIN_ROWS_PER_WORKER, u64::MAX).await?;
316370
assert!(boundaries.is_empty());
317371
Ok(())
318372
}
@@ -321,28 +375,26 @@ mod tests {
321375
async fn small_table_gets_no_boundaries() -> Result<(), MySqlError> {
322376
// All keys share one depth-1 prefix and fit under `min_rows_per_worker`,
323377
// so the single open-ended range yields no boundary.
324-
let mut db = MockDb { keys: keys(10_000) };
378+
let mut db = MockDb::new(keys(10_000));
325379
let count = u64::cast_from(db.keys.len());
326-
let boundaries = partition(&mut db, 4, count, MIN_ROWS_PER_WORKER).await?;
380+
let boundaries = partition(&mut db, 4, count, MIN_ROWS_PER_WORKER, u64::MAX).await?;
327381
assert!(boundaries.is_empty());
328382
Ok(())
329383
}
330384

331385
#[mz_ore::test(tokio::test)]
332386
async fn empty_table_gets_no_boundaries() -> Result<(), MySqlError> {
333-
let mut db = MockDb { keys: vec![] };
334-
let boundaries = partition(&mut db, 4, 0, MIN_ROWS_PER_WORKER).await?;
387+
let mut db = MockDb::new(vec![]);
388+
let boundaries = partition(&mut db, 4, 0, MIN_ROWS_PER_WORKER, u64::MAX).await?;
335389
assert!(boundaries.is_empty());
336390
Ok(())
337391
}
338392

339393
#[mz_ore::test(tokio::test)]
340394
async fn splits_evenly_across_workers() -> Result<(), MySqlError> {
341-
let mut db = MockDb {
342-
keys: keys(200_000),
343-
};
395+
let mut db = MockDb::new(keys(200_000));
344396
let count = u64::cast_from(db.keys.len());
345-
let boundaries = partition(&mut db, 4, count, MIN_ROWS_PER_WORKER).await?;
397+
let boundaries = partition(&mut db, 4, count, MIN_ROWS_PER_WORKER, u64::MAX).await?;
346398
assert_eq!(boundaries.len(), 3);
347399
// Boundaries must be sorted and split the keys into ~50k chunks.
348400
let mut prev = 0;
@@ -361,9 +413,9 @@ mod tests {
361413

362414
#[mz_ore::test(tokio::test)]
363415
async fn low_min_rows_per_worker_splits_small_tables() -> Result<(), MySqlError> {
364-
let mut db = MockDb { keys: keys(1000) };
416+
let mut db = MockDb::new(keys(1000));
365417
let count = u64::cast_from(db.keys.len());
366-
let boundaries = partition(&mut db, 4, count, 10).await?;
418+
let boundaries = partition(&mut db, 4, count, 10, u64::MAX).await?;
367419
assert_eq!(boundaries.len(), 3);
368420
let mut prev = 0;
369421
for b in &boundaries {
@@ -386,9 +438,9 @@ mod tests {
386438
// stalling on the all-encompassing "U" prefix.
387439
let mut all_keys = vec!["U".to_string()];
388440
all_keys.extend((0..1000).map(|i| format!("U{i:06}")));
389-
let mut db = MockDb { keys: all_keys };
441+
let mut db = MockDb::new(all_keys);
390442
let count = u64::cast_from(db.keys.len());
391-
let boundaries = partition(&mut db, 4, count, 10).await?;
443+
let boundaries = partition(&mut db, 4, count, 10, u64::MAX).await?;
392444
assert_eq!(boundaries.len(), 3);
393445
for b in &boundaries {
394446
assert!(
@@ -403,8 +455,8 @@ mod tests {
403455
async fn fractional_target_still_terminates() -> Result<(), MySqlError> {
404456
// count / (workers * 4) is fractional and the minimum is zero, so
405457
// the target floors at one row instead of splitting forever.
406-
let mut db = MockDb { keys: keys(3) };
407-
let boundaries = partition(&mut db, 4, 3, 0).await?;
458+
let mut db = MockDb::new(keys(3));
459+
let boundaries = partition(&mut db, 4, 3, 0, u64::MAX).await?;
408460
assert_eq!(boundaries, vec!["000001", "000002"]);
409461
Ok(())
410462
}
@@ -437,12 +489,13 @@ mod tests {
437489
let total = u64::cast_from(all_keys.len());
438490

439491
// A minimum above the table size yields no boundaries at all.
440-
let bounds = partition_table(&mut conn, table.clone(), "id", 4, total, 50_000).await?;
492+
let bounds =
493+
partition_table(&mut conn, table.clone(), "id", 4, total, 50_000, u64::MAX).await?;
441494
assert!(bounds.is_empty(), "{bounds:?}");
442495

443496
// A low minimum splits inside the 'a' extensions rather than stopping
444497
// at the exact key.
445-
let bounds = partition_table(&mut conn, table, "id", 4, total, 10).await?;
498+
let bounds = partition_table(&mut conn, table, "id", 4, total, 10, u64::MAX).await?;
446499
assert_eq!(bounds.len(), 3, "{bounds:?}");
447500

448501
// MySQL agrees the boundaries are strictly increasing.
@@ -488,7 +541,7 @@ mod tests {
488541
let total = u64::cast_from(all_keys.len());
489542

490543
// Partition for 4 workers with a minimum split size around 250.
491-
let bounds = partition_table(&mut conn, table, "id", 4, total, 250).await?;
544+
let bounds = partition_table(&mut conn, table, "id", 4, total, 250, u64::MAX).await?;
492545
assert_eq!(bounds.len(), 3);
493546
let counts = partition_counts(&mut conn, DB, &bounds, total).await?;
494547
// ~8k keys are visible, so each count gets at least 2k under perfect

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,18 @@ pub static MYSQL_SOURCE_SNAPSHOT_PARALLELISM: Config<bool> = Config::new(
215215
"Whether to split MySQL snapshot reads across workers by primary-key ranges.",
216216
);
217217

218+
/// Probe query budget for the MySQL snapshot prefix partitioner, scaled to
219+
/// the table's estimated size so probing effort stays proportional to the
220+
/// snapshot work it optimizes. A small floor applies so modest tables can
221+
/// still afford their handful of splits. When a table's budget runs out,
222+
/// splitting stops early and buckets come out coarser, never incorrect.
223+
pub static MYSQL_SOURCE_SNAPSHOT_PARTITION_REQUESTS_PER_BILLION_ROWS: Config<usize> = Config::new(
224+
"mysql_source_snapshot_partition_requests_per_billion_rows",
225+
2_500,
226+
"Cap on MySQL snapshot PK-prefix partitioning probe queries per table, per billion \
227+
estimated rows; when exhausted, splitting stops early with coarser buckets.",
228+
);
229+
218230
/// If the optimizer estimates the table has fewer rows than this, compute the exact row count
219231
/// with `COUNT(*)`. Otherwise, report the `information_schema` estimate directly.
220232
pub static MYSQL_SOURCE_SNAPSHOT_EXACT_COUNT_MAX_ROWS: Config<usize> = Config::new(
@@ -438,6 +450,7 @@ pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet {
438450
.add(&MYSQL_REPLICATION_HEARTBEAT_INTERVAL)
439451
.add(&MYSQL_SOURCE_SNAPSHOT_EXACT_COUNT_MAX_ROWS)
440452
.add(&MYSQL_SOURCE_SNAPSHOT_PARALLELISM)
453+
.add(&MYSQL_SOURCE_SNAPSHOT_PARTITION_REQUESTS_PER_BILLION_ROWS)
441454
.add(&ORE_OVERFLOWING_BEHAVIOR)
442455
.add(&PG_FETCH_SLOT_RESUME_LSN_INTERVAL)
443456
.add(&PG_SCHEMA_VALIDATION_INTERVAL)

0 commit comments

Comments
 (0)