4646 VarChar ,
4747)
4848from materialize .data_ingest .query_error import QueryError
49- from materialize .data_ingest .row import Operation
5049from materialize .mzcompose import get_default_system_parameters
5150from materialize .mzcompose .composition import Composition
5251from materialize .mzcompose .services .materialized import (
@@ -243,6 +242,25 @@ def applicable(self, exe: Executor) -> bool:
243242 coverage check meaningful."""
244243 return True
245244
245+ def insert_batch_size (self , exe : Executor , table : Table ) -> int :
246+ """How many rows the next insert into `table` may add, 0 once the table
247+ sits at MAX_ROWS."""
248+ available = MAX_ROWS - table .num_rows
249+ if available < 1 :
250+ return 0
251+ # A per-worker share of the budget, not the whole remainder: `num_rows`
252+ # is read here but only bumped once the statement returns, so every
253+ # worker would otherwise size a full-budget batch off the same stale
254+ # counter and the table would overshoot MAX_ROWS by roughly a
255+ # worker-count factor, which is what kept the cap from binding at all.
256+ # A share bounds the overshoot at one extra round of concurrent batches.
257+ # Charging the budget up front instead would bind exactly, but a claim
258+ # is not given back when the statement's transaction rolls back, so the
259+ # insert actions would stall against tables that read as full while
260+ # holding far fewer rows.
261+ share = max (1 , MAX_ROWS // exe .db .num_threads )
262+ return self .rng .randint (1 , min (available , share ))
263+
246264 def create_system_connection (
247265 self , exe : Executor , num_attempts : int = 10
248266 ) -> Connection :
@@ -1066,10 +1084,13 @@ def run(self, exe: Executor) -> bool:
10661084 return False
10671085 table = self .rng .choice (tables )
10681086
1087+ num_rows = self .insert_batch_size (exe , table )
1088+ if not num_rows :
1089+ return False
1090+
10691091 column_names = ", " .join (column .name (True ) for column in table .columns )
10701092 column_values = []
1071- max_rows = min (100 , MAX_ROWS - table .num_rows )
1072- for i in range (self .rng .randrange (1 , max_rows + 1 )):
1093+ for i in range (num_rows ):
10731094 column_values .append (
10741095 ", " .join (column .value (self .rng , True ) for column in table .columns )
10751096 )
@@ -1081,7 +1102,7 @@ def run(self, exe: Executor) -> bool:
10811102 self .exe_prepared (query , f"insert{ self .stmt_id } " , exe )
10821103 else :
10831104 exe .execute (query , http = Http .RANDOM )
1084- table .num_rows += len ( column_values )
1105+ table .num_rows += num_rows
10851106 exe .insert_table = table .table_id
10861107 return True
10871108
@@ -1137,10 +1158,11 @@ def run(self, exe: Executor) -> bool:
11371158 f"({ expression (column .data_type , source .columns , self .rng , kind = ExprKind .WRITE )} )::{ column .data_type .name ()} "
11381159 for column in table .columns
11391160 )
1140- # `num_rows` can be stale, so clamp instead of trusting the filter
1141- # above. The LIMIT keeps a self-insert from doubling the table on every
1161+ # The LIMIT keeps a self-insert from doubling the table on every
11421162 # attempt.
1143- limit = self .rng .randint (1 , max (1 , min (100 , MAX_ROWS - table .num_rows )))
1163+ limit = self .insert_batch_size (exe , table )
1164+ if not limit :
1165+ return False
11441166 query = (
11451167 f"INSERT INTO { table } ({ column_names } ) SELECT { expressions } FROM { source } "
11461168 f" WHERE { expression (Boolean , source .columns , self .rng , kind = ExprKind .WRITE )} "
@@ -1199,13 +1221,16 @@ def run(self, exe: Executor) -> bool:
11991221 return False
12001222 table = self .rng .choice (tables )
12011223
1224+ num_rows = self .insert_batch_size (exe , table )
1225+ if not num_rows :
1226+ return False
1227+
12021228 values = []
1203- max_rows = min (100 , MAX_ROWS - table .num_rows )
1204- for i in range (self .rng .randrange (1 , max_rows + 1 )):
1229+ for i in range (num_rows ):
12051230 values .append ([column .value (self .rng , False ) for column in table .columns ])
12061231 query = f"COPY INTO { table } FROM STDIN"
12071232 exe .copy (query , values )
1208- table .num_rows += len ( values )
1233+ table .num_rows += num_rows
12091234 exe .insert_table = table .table_id
12101235 return True
12111236
@@ -1250,10 +1275,13 @@ def run(self, exe: Executor) -> bool:
12501275 return False
12511276 table = self .rng .choice (tables )
12521277
1278+ num_rows = self .insert_batch_size (exe , table )
1279+ if not num_rows :
1280+ return False
1281+
12531282 column_names = ", " .join (column .name (True ) for column in table .columns )
12541283 column_values = []
1255- max_rows = min (100 , MAX_ROWS - table .num_rows )
1256- for i in range (self .rng .randrange (1 , max_rows + 1 )):
1284+ for i in range (num_rows ):
12571285 column_values .append (
12581286 ", " .join (column .value (self .rng , True ) for column in table .columns )
12591287 )
@@ -1280,7 +1308,7 @@ def run(self, exe: Executor) -> bool:
12801308 self .exe_prepared (query , f"insert_returning{ self .stmt_id } " , exe )
12811309 else :
12821310 exe .execute (query , http = Http .RANDOM )
1283- table .num_rows += len ( column_values )
1311+ table .num_rows += num_rows
12841312 exe .insert_table = table .table_id
12851313 return True
12861314
@@ -1319,16 +1347,28 @@ def run(self, exe: Executor) -> bool:
13191347
13201348
13211349class SourceInsertAction (Action ):
1350+ """Feed one workload transaction to a random CDC source's upstream system.
1351+
1352+ Unlike the table DML actions this carries no MAX_ROWS gate, because a
1353+ `data_ingest` workload bounds its own upstream table: the upsert workloads
1354+ only ever write key 0, and the delete-at-end-of-day workloads insert
1355+ `Records.SOME` keys and then delete exactly those keys again, so the table
1356+ oscillates and every cycle nets to zero. A gate could only ever misfire.
1357+ `next()` has already consumed a transaction by the time a row count could be
1358+ checked, so refusing to run it would drop one half of an insert/delete pair
1359+ and leave the count drifting away from the upstream table for the rest of the
1360+ run. A workload that grows without bound would need a different mechanism,
1361+ not a row count, since `MySqlSource.prepopulate_rows` already puts up to
1362+ 30,000 rows upstream before the source is even created."""
1363+
13221364 def run (self , exe : Executor ) -> bool :
13231365 with exe .db .lock :
1324- sources = [
1325- source
1326- for source in exe .db .kafka_sources
1366+ sources = (
1367+ exe .db .kafka_sources
13271368 + exe .db .postgres_sources
13281369 + exe .db .mysql_sources
13291370 + exe .db .sql_server_sources
1330- if source .num_rows < MAX_ROWS
1331- ]
1371+ )
13321372 if not sources :
13331373 return False
13341374 source = self .rng .choice (sources )
@@ -1341,14 +1381,7 @@ def run(self, exe: Executor) -> bool:
13411381 ]:
13421382 return False
13431383
1344- transaction = next (source .generator )
1345- for row_list in transaction .row_lists :
1346- for row in row_list .rows :
1347- if row .operation == Operation .INSERT :
1348- source .num_rows += 1
1349- elif row .operation == Operation .DELETE :
1350- source .num_rows -= 1
1351- source .executor .run (transaction , logging_exe = exe )
1384+ source .executor .run (next (source .generator ), logging_exe = exe )
13521385 return True
13531386
13541387
@@ -1469,7 +1502,7 @@ def run(self, exe: Executor) -> bool:
14691502 for t in exe .db .tables
14701503 if t != table and (not t .temp or t in exe .temp_objects )
14711504 ]
1472- # TODO: Drop the RepeatRow gate once database-issues#9308 is fixed.
1505+ # TODO: Drop the RepeatRow gate once STG-36 is fixed.
14731506 # DELETE .. USING lowers to a semijoin whose DistinctBy can leave the
14741507 # target table with a net-negative row, and every later reader of that
14751508 # table then surfaces the corruption. Tolerating that class outside
@@ -2768,6 +2801,18 @@ def run(self, exe: Executor) -> bool:
27682801
27692802
27702803class FlipFlagsAction (Action ):
2804+ # Shortest gap between two flips, fleet-wide. Every worker's action list
2805+ # carries this action, so weights alone put it at ~30 flips/s, which drove
2806+ # ~43k dyncfg applications per replica and most of a 132 MB services.log.
2807+ # What this action is after is a flag changing under a running dataflow, and
2808+ # one change per second is plenty for that.
2809+ MIN_INTERVAL_SEC = 1.0
2810+
2811+ # Guards `last_flip`, which is shared across workers since each holds its
2812+ # own instance of this action.
2813+ interval_lock = threading .Lock ()
2814+ last_flip = 0.0
2815+
27712816 def __init__ (
27722817 self ,
27732818 rng : random .Random ,
@@ -2927,10 +2972,10 @@ def __init__(
29272972 BOOLEAN_FLAG_VALUES
29282973 )
29292974 self .flags_with_values ["cluster" ] = ["quickstart" , "dont_exist" ]
2930- # NOTE: enable_frontend_peek_sequencing is pinned off in
2931- # ADDITIONAL_SYSTEM_PARAMETER_DEFAULTS (frontend-peek read-hold vs
2932- # compaction race, https://linear.app/materializeinc/issue/SQL-520), so
2933- # it is not flipped here.
2975+ self . flags_with_values [ " enable_frontend_peek_sequencing" ] = [
2976+ "true" ,
2977+ "false" ,
2978+ ]
29342979 self .flags_with_values ["enable_frontend_subscribes" ] = [
29352980 "true" ,
29362981 "false" ,
@@ -3251,6 +3296,12 @@ def errors_to_ignore(self, exe: Executor) -> list[str]:
32513296 ] + super ().errors_to_ignore (exe )
32523297
32533298 def run (self , exe : Executor ) -> bool :
3299+ with FlipFlagsAction .interval_lock :
3300+ now = time .time ()
3301+ if now - FlipFlagsAction .last_flip < FlipFlagsAction .MIN_INTERVAL_SEC :
3302+ return False
3303+ FlipFlagsAction .last_flip = now
3304+
32543305 # A tenth of the time set a random cluster's arrangement dictionary
32553306 # compression instead of flipping a global flag. The per-cluster option
32563307 # and the global `enable_arrangement_dictionary_compression_alpha` flag
@@ -3895,6 +3946,20 @@ def run(self, exe: Executor) -> bool:
38953946 f"{ new_size } sat in-progress { read_ts_ms - int (record_deadline )} ms "
38963947 f"past its deadline { record_deadline } "
38973948 )
3949+ # A settled record is retained with its terminal status, so a
3950+ # NULL from the LEFT JOIN means the record this ALTER wrote is
3951+ # gone, which no code path is allowed to do. Every writer of the
3952+ # durable field keeps the record and only moves its status:
3953+ # `reshape_alter_cluster_managed` writes InProgress or, for an
3954+ # ALTER back to the realized shape, Cancelled;
3955+ # `cancel_carried_reconfiguration` mutates the status in place;
3956+ # and the controller's three writes are Finalized, TimedOut and
3957+ # ResourceExhausted. `ReconfigurationWrite.record` is an Option
3958+ # documented as "None to clear it", but nothing constructs that.
3959+ # The relation also drops a cluster that stops being managed, and
3960+ # environmentd's bootstrap is the one place that resets the field
3961+ # outright, hence the scenario carve-outs in `applicable` and no
3962+ # `SET (MANAGED ...)` anywhere in the workload.
38983963 if read_ts_ms > our_deadline_ms and status is None :
38993964 raise ValueError (
39003965 f"Reconfiguration record of cluster { cluster } is gone as of "
@@ -5624,10 +5689,16 @@ def run(self, exe: Executor) -> bool:
56245689 log = f"POST { url } Headers: { ', ' .join (headers_strs )} Body: { payload .encode ('utf-8' )} "
56255690 exe .log (log )
56265691 try :
5627- source .num_rows += 1
56285692 result = requests .post (url , data = payload .encode (), headers = headers )
56295693 if result .status_code != 200 :
56305694 raise QueryError (f"{ result .status_code } : { result .text } " , log )
5695+ # Count after the POST landed. A webhook source is append-only,
5696+ # so its row budget is never released, and counting rejected
5697+ # POSTs retires the source without a single row ever reaching
5698+ # it. A source whose POSTs all 404 (a concurrent cascading drop,
5699+ # a rename) then silently stops being posted to for the rest of
5700+ # the run.
5701+ source .num_rows += 1
56315702 except requests .exceptions .ConnectionError :
56325703 # Expected when Mz is killed
56335704 if exe .db .scenario not in (
@@ -6177,13 +6248,7 @@ def __init__(
61776248 (CopyToStdoutAction , 20 ),
61786249 (ShowAction , 10 ),
61796250 (SystemCatalogReadAction , 10 ),
6180- # TODO: Reenable once EXPLAIN FILTER PUSHDOWN can no longer panic the
6181- # coordinator when a referenced compute collection is concurrently
6182- # dropped. sequence_explain_pushdown -> acquire_read_holds().expect(
6183- # "missing compute collection") at read_policy.rs:389 (normal peeks and
6184- # EXPLAIN ANALYZE handle the drop gracefully).
6185- # See https://linear.app/materializeinc/issue/SQL-519
6186- # (ExplainFilterPushdownAction, 5),
6251+ (ExplainFilterPushdownAction , 5 ),
61876252 # PREPARED BUT DISABLED (see class docstrings): enabling these now just
61886253 # re-detects known-unfixed coordinator bugs rather than finding new ones.
61896254 # (DependencyConsistencyAction, 5), # TODO: enable once SQL-521 fixed
@@ -6279,9 +6344,8 @@ def __init__(
62796344 (DropLoadGeneratorSourceAction , 4 ),
62806345 (CreateMultiLoadGeneratorSourceAction , 2 ),
62816346 (DropMultiLoadGeneratorSourceAction , 2 ),
6282- # TODO: Reenable when https://linear.app/materializeinc/issue/SS-307 is fixed
6283- # (CreateMySqlSourceAction, 4),
6284- # (DropMySqlSourceAction, 4),
6347+ (CreateMySqlSourceAction , 4 ),
6348+ (DropMySqlSourceAction , 4 ),
62856349 (CreatePostgresSourceAction , 4 ),
62866350 (DropPostgresSourceAction , 4 ),
62876351 # TODO: Reenable when https://linear.app/materializeinc/issue/SS-290 is fixed
@@ -6343,10 +6407,7 @@ def __init__(
63436407 (DDLTransactionAction , 2 ),
63446408 (SystemCatalogReadAction , 4 ),
63456409 (ExplainAnalyzeAction , 4 ),
6346- # TODO: Reenable with EXPLAIN FILTER PUSHDOWN's coordinator panic on a
6347- # concurrently-dropped compute collection (read_policy.rs:389).
6348- # See https://linear.app/materializeinc/issue/SQL-519
6349- # (ExplainFilterPushdownAction, 2),
6410+ (ExplainFilterPushdownAction , 2 ),
63506411 (FlipFlagsAction , 2 ),
63516412 # TODO: Reenable when https://linear.app/materializeinc/issue/SQL-405 is fixed.
63526413 # (AlterTableAddColumnAction, 10),
0 commit comments