Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
120 commits
Select commit Hold shift + click to select a range
9dff86f
compute: two-runtime read isolation (Arc-batch spines, arrangement sh…
antiguru Jul 23, 2026
c61bba9
testdrive: bump ii_t4 arrangement-size bound for Arc-batch overhead
antiguru Jul 23, 2026
ad75d81
compute: release input read holds when a deferred interactive dataflo…
antiguru Jul 23, 2026
155ef3a
parallel-benchmark: add TwoRuntimeReadIsolation scenario
antiguru Jul 23, 2026
61abaf5
compute: wake interactive peeks waiting on a re-exported index's seal
antiguru Jul 23, 2026
5db39ff
doc: design for two-runtime arrangement sharing capture-based lifecycle
antiguru Jul 23, 2026
e63c50c
doc: two-runtime arrangement sharing lifecycle (protocol-split design)
antiguru Jul 23, 2026
494814c
doc: drop the Import command, use shared imports by role
antiguru Jul 23, 2026
f08c31c
doc: implementation plan for two-runtime arrangement sharing lifecycle
antiguru Jul 23, 2026
2da28b3
compute: single-source the shared-trace replay feed from the stream f…
antiguru Jul 23, 2026
534c424
compute: pre-allocated publication points (placeholder plus adopt-in-…
antiguru Jul 23, 2026
3384d49
compute: registry get-or-create for publication points
antiguru Jul 23, 2026
b0f1ce3
compute: close and evict never-adopted placeholders
antiguru Jul 23, 2026
f6b87fe
compute-client: route to interactive by the bounded-read predicate
antiguru Jul 23, 2026
c2e765c
compute: render interactive dataflows in arrival order with late-boun…
antiguru Jul 23, 2026
0f36b46
compute: assert since <= as_of on the interactive import path
antiguru Jul 23, 2026
5f51d14
compute: remove the dead live import(), keep import_snapshot_at
antiguru Jul 23, 2026
f0a0490
compute: harden interactive routing tripwire and document eviction ok…
antiguru Jul 24, 2026
ce61a15
compute: document that shared-import cleanup relies on the controller…
antiguru Jul 24, 2026
273df48
compute: fix broken intra-doc links flagged by lint-doc
antiguru Jul 24, 2026
6e4f834
compute: remove the dead insert/publish registry API superseded by adopt
antiguru Jul 24, 2026
b26a6cf
compute: collapse the now single-variant PendingWork enum
antiguru Jul 24, 2026
ecbf731
parallel-benchmark: encode HydrationChurn queries for psycopg execute
antiguru Jul 24, 2026
8be6613
compute: seed shared-trace importers with the full snapshot, not a fr…
antiguru Jul 24, 2026
3867c25
sqllogictest: rewrite relations.slt golden for shared-arrangement ope…
antiguru Jul 24, 2026
00a8591
compute: fire the shared-trace seal signal from the sink, not an upst…
antiguru Jul 24, 2026
9783c4e
compute: rewrite relations golden and fix doc link after the seal-sig…
antiguru Jul 24, 2026
f786665
compute: drop DD fork, drive publisher compaction from AllowCompaction
antiguru Jul 24, 2026
1550c97
parallel-benchmark: halve TwoRuntimeReadIsolation read rates
antiguru Jul 24, 2026
4e39de8
compute: fix interactive_import tests for AllowCompaction-driven publ…
antiguru Jul 24, 2026
c106b13
parallel-benchmark: lower TwoRuntimeReadIsolation read probe to 50/12
antiguru Jul 24, 2026
f3de43f
doc: consolidate two-runtime design into a single document
antiguru Jul 24, 2026
183fc2a
doc: record interactive-runtime observability open items
antiguru Jul 24, 2026
01b45b6
compute: remove enable_index_arrangement_sharing dyncfg
antiguru Jul 24, 2026
3799b93
compute-client: make transient_owner a BTreeSet
antiguru Jul 24, 2026
7a27fdd
clusterd-test-driver: use t-prefixed transient id in two-runtime spec
antiguru Jul 24, 2026
5cc4d33
compute: police the shared-trace cut and seed importers with the chai…
antiguru Aug 8, 2026
e24ce62
compute: cut the unexercised sharing machinery, arm the routing tripwire
antiguru Aug 8, 2026
2f960df
storage: port the Iceberg sink's batch stash to Arc-backed batches
antiguru Aug 8, 2026
02df0ee
doc: correct the two-runtime design record and record the arrangement…
antiguru Aug 8, 2026
270ca2f
compute: let interactive peeks use the peek stash
antiguru Aug 8, 2026
7b110ae
compute: walk fast-path index peeks off the serving worker
antiguru Aug 9, 2026
381972a
compute: document what the peek-placement flags actually change
antiguru Aug 9, 2026
7425f36
doc: experimental evaluation for the read-placement matrix
antiguru Aug 9, 2026
04001db
compute: scope the peek-placement flags to replicas, and plan the arm…
antiguru Aug 9, 2026
6c38d75
compute: drop the shared-trace cut precondition, it does not hold
antiguru Aug 9, 2026
39cf278
compute: drop live batches an importer's seed already covers
antiguru Aug 9, 2026
0d36c2b
doc: write down the protocol invariant the runtime split breaks
antiguru Aug 9, 2026
ce5db21
compute-client: cap compaction at what in-flight interactive dataflow…
antiguru Aug 9, 2026
11fd652
compute: fix the lint failures from the protocol work
antiguru Aug 9, 2026
f2d39c3
compute-client: never let capping regress a compaction frontier
antiguru Aug 9, 2026
2ec910d
controller: shorten the interactive port name, Kubernetes rejects 19 …
antiguru Aug 9, 2026
31de976
doc: record the staging results for E6 and E1
antiguru Aug 9, 2026
e735cc1
compute: count index peek walks by substrate
antiguru Aug 9, 2026
9329b87
doc: record that peeks broadcast regardless of replica targeting
antiguru Aug 9, 2026
7bfc5e6
doc: arms needing independent load need a cluster each
antiguru Aug 9, 2026
d5f0628
doc: E1 and E2 measured, the offload alone removes head-of-line blocking
antiguru Aug 9, 2026
9556581
doc: plan to bring the persist stash to the offloaded walk
antiguru Aug 10, 2026
474a17c
doc: E7 measures temporary dataflows, the second runtime's remaining …
antiguru Aug 10, 2026
e06d7f6
controller: scope enable_two_runtime_compute to replicas
antiguru Aug 10, 2026
53b00b1
compute: let an offloaded peek walk divert to the stash
antiguru Aug 10, 2026
fa9b1c3
doc: reruns with the stash on, and a first look at disk contention
antiguru Aug 10, 2026
c3254c0
doc: swap with a matched control shows no regression
antiguru Aug 10, 2026
3ba44c0
doc: E9 measures the console under load, latency and staleness apart
antiguru Aug 10, 2026
6cb0882
doc: hydration memory is a sawtooth that overshoots about fourfold
antiguru Aug 10, 2026
e0e11ac
doc: the observability itself is the result
antiguru Aug 10, 2026
958b427
doc: E11 measures the skewed point lookup, a real customer pattern
antiguru Aug 10, 2026
cb44957
doc: how the branch splits into landable pieces
antiguru Aug 10, 2026
26b282d
compute: fix defects found by adversarial review
antiguru Aug 10, 2026
43bb0e2
doc: design of record reflects the fork, and the open review findings
antiguru Aug 10, 2026
2680679
doc: systems framing, the incremental path, and the serving-layer con…
antiguru Aug 11, 2026
e39101f
doc: equal peer counts are a soundness requirement, not a sizing choice
antiguru Aug 11, 2026
2523ebe
doc: what the platform isolates, and the QoS constraint on pinning
antiguru Aug 11, 2026
98bd503
doc: correct the CPU isolation story, the class is Burstable
antiguru Aug 11, 2026
b957d1c
doc: problems and mechanisms as the spine
antiguru Aug 11, 2026
2dc11a3
doc: correct the mechanism table after adversarial review
antiguru Aug 11, 2026
c58dda2
doc: E12 measures the operator-activation mechanism, and inverts two …
antiguru Aug 11, 2026
ce67fc6
doc: add freshness as a mechanism, with predictions registered before…
antiguru Aug 11, 2026
5c909f1
doc: E13 measures the freshness dual, refuting one registered prediction
antiguru Aug 11, 2026
831fdf9
doc: complete the mechanism table, and say what is still missing from it
antiguru Aug 11, 2026
d09cc95
doc: SUBSCRIBE slow start, and the exclusion that makes it unmovable
antiguru Aug 12, 2026
ce59f54
doc: the import shape is not the obstacle, the read hold is
antiguru Aug 12, 2026
3b996d1
doc: the interactive read hold needs early release, not advancing
antiguru Aug 12, 2026
300d1c2
doc: compaction feedback exists, and two frozen holds defeat it
antiguru Aug 12, 2026
866c347
doc: separate the cause of the inert physical hold from how it manifests
antiguru Aug 12, 2026
1cb8f39
doc: the shared handle reports the request, not the grant
antiguru Aug 12, 2026
198fd52
doc: localise the physical-hold defect to two divergences from TraceA…
antiguru Aug 12, 2026
eeb22c7
doc: capture where the design exploration stands, and mark the split …
antiguru Aug 12, 2026
dad8727
compute: hold back shared-trace compaction correctly
antiguru Aug 12, 2026
05c1615
compute: make the interactive import's read hold downgradeable
antiguru Aug 12, 2026
2fc8abb
compute-client: retire interactive holds on confirmation, not on the …
antiguru Aug 12, 2026
afcb80c
doc: replace the TLA+ protocol sketch with a checked Lean 4 model
antiguru Aug 12, 2026
a336399
compute: give the shared import a hold that outlives construction
antiguru Aug 12, 2026
c34f402
compute: fix three defects in the shared-trace hold machinery
antiguru Aug 12, 2026
a5bd883
doc: design and TLA+ model for read holds across two runtimes
antiguru Aug 12, 2026
720b70c
compute-client: add the hold commands, without their semantics
antiguru Aug 12, 2026
2ea9ca6
doc: record the read-hold implementation sequence
antiguru Aug 12, 2026
939099b
compute: install and maintain command-acquired read holds
antiguru Aug 12, 2026
dd1678f
doc: record what steps 1 and 2 settled, and that G3 dissolves
antiguru Aug 12, 2026
b106806
compute-client: synthesize the hold commands and delete the cap
antiguru Aug 12, 2026
15a10fa
doc: bring read-holds.md up to what was built
antiguru Aug 12, 2026
3c6f3e8
doc: make the TLA+ model check what was built
antiguru Aug 12, 2026
87075f8
compute: format the rebase conflict resolution
antiguru Aug 12, 2026
516c48e
clusterd: depend on mz-compute-client again
antiguru Aug 12, 2026
93c9f53
compute: stop linking module docs to private items
antiguru Aug 12, 2026
9383dea
compute: derive physical compaction from readers' cut floors
antiguru Aug 13, 2026
357e12f
compute: model broadcast compaction in the holds spec
antiguru Aug 13, 2026
f279fe5
compute: design broadcast compaction, superseding the hold protocol
antiguru Aug 13, 2026
4e5a1ed
compute: broadcast compaction and assert I1c
antiguru Aug 13, 2026
0c14e23
compute: delete the command-acquired hold layer
antiguru Aug 13, 2026
7e93173
compute: remodel the holds spec around stream positions
antiguru Aug 13, 2026
fa56da4
compute: do not broadcast a transient collection's compaction
antiguru Aug 13, 2026
775ea1f
compute: report the standing hold when an import is refused
antiguru Aug 13, 2026
4fbddd4
compute: rename the flag to enable_compute_interactive_runtime
antiguru Aug 17, 2026
b1b562f
design: park peek placement as an orthogonal problem
antiguru Aug 17, 2026
15e8be3
design: move the design docs and models to their own PR
antiguru Aug 17, 2026
0406c04
compute: move the peek offload to its own branch
antiguru Aug 17, 2026
3e03995
compute: reconcile with the merged ComputeRuntimeRole and error-disti…
antiguru Aug 21, 2026
f73b5ce
compute: walk fast-path index peeks off the serving worker
antiguru Aug 17, 2026
f8fc314
design: move the peek-placement documents to the offload work
antiguru Aug 21, 2026
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
4 changes: 4 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
# Bringing the persist stash to the offloaded walk

**Parked with the rest of the peek-placement work, and implemented.** This describes the gate
that made the offload unreachable in a production configuration and how it was removed. It
is retained because the gate is the reason two rounds of staging measurement showed nothing,
which is worth not rediscovering. Whether the offload survives at all is
[peek-placement.md](peek-placement.md)'s question.

## Why

E1 measured the offload cutting point-lookup tail latency 32-fold behind a 2.2 second scan.
That measurement required `enable_compute_peek_response_stash = false` for the environment.
With the stash on, which is how production runs, `should_offload_peek` declines every peek the stash could take, and the win is unreachable.

The gate is one branch.

```rust
fn should_offload_peek(&self, peek_stash_usable: bool) -> bool {
if peek_stash_usable {
return false;
}
```

`peek_stash_usable` is a capability rather than a prediction.
It is `finishing.is_streamable(arity) && ENABLE_PEEK_RESPONSE_STASH && location_available`, with no reference to result size, so a peek returning three rows declines the offload over a divergence that will never happen.
The reachable domain today is peeks with an `ORDER BY` or a non-identity projection, which is not the traffic the offload was built for.

## What the two mechanisms actually do

They are usually described as alternatives.
They are not, and seeing that is what makes the fix small.

| | Walk runs on | Response path | Fixes |
|---|---|---|---|
| Peek response stash | the timely worker | persist, uploaded by an async task fed over a channel | response size |
| Index peek offload | a blocking task | inline, returned whole | head-of-line blocking |

The stash never moved the walk.
`pump_rows` keeps stepping the row iterator on the worker thread and pushes batches into `rows_tx`, and the comment on that field says exactly why: a trace cursor is not `Send`, so the iterator cannot be handed to the upload task.

That constraint is the one the offload has already solved.

## The two facts that make this cheap

Both are already in the tree, neither needs changing.

* `StashingPeek::start_upload` takes `Box<dyn Iterator<Item = Result<(Row, NonZeroI64), String>>>`.
It is already independent of which trace produced the rows, because the maintenance path walks a `TraceBundle` and the interactive path a registry handle.
* `spawn_offloaded_walk` is bounded on `TraceCursor<OksTr>: Send` and `TraceStorage<OksTr>: Send`, and `OffloadSnapshot::Ready` carries owned cursors and storage.
The offloaded walk's row source crosses threads by construction.

So on the offload path the reason `pump_rows` exists does not apply.
An offloaded walk can own the persist writer directly and stream into it as it walks, with no channel hop and no worker involvement at all.
Unifying these makes the stash simpler on the path that matters, not more complex.

## Plan

### Phase 1: let the offloaded walk stash

Pass the stash configuration into `spawn_offloaded_walk`: eligibility, threshold, persist location, and `batch_max_runs`.
The task walks its snapshot, accumulates as it does today, and on crossing the threshold opens the upload and streams the remainder into it.
`IndexOffloadPeek` then resolves to either a rows response or a stashed response, and `PendingPeek::IndexOffload` retires both the same way it retires a rows response now.

The mid-walk divergence that `PeekStatus::UsePeekStash` expresses as a return value becomes an ordinary branch inside the task.
Nothing has to be communicated back to the worker to make the decision, which is what forced the gate in the first place.

### Phase 2: delete the gate

`should_offload_peek` stops consulting `peek_stash_usable` and keeps only the in-flight cap.
The offload becomes reachable for ordinary streamable peeks, which is where the measured win lives.

Phases 1 and 2 land together behind `enable_index_peek_offload`, which is already off by default in production.

### Phase 3: demote the worker-pumped stash to a fallback

The inline stash still has to exist, because the offload declines in cases the stash does not.
Those are the in-flight cap being reached, `OffloadSnapshot::NotReady`, and any trace whose cursor cannot produce a `Send` snapshot.
This is a demotion rather than a deletion, and it should wait until E4 has been re-run, because the cap is what decides how often the fallback is taken.

## What to be careful about

**Compaction holds get longer.**
An offloaded walk pins its batches for the duration of the walk.
A stashing walk pins them for the walk plus the upload, which is seconds rather than milliseconds.
`index_peek_offload_max_inflight` now bounds a much longer-lived hold, and its interaction with the controller read holds and the capping in the multiplexer needs checking before the cap default is trusted.

**The finishing rule must not drift.**
`is_streamable` decides whether rows can be shipped as they are produced.
A non-streamable finishing has to accumulate and sort before anything is emitted, and the offloaded path must apply the same rule rather than assuming it can always stream.

**The size limit must stay identical.**
`max_result_size` is checked during the inline walk.
If the offloaded walk checks it differently, a peek that errors on one substrate succeeds on the other, which is a correctness difference visible to users.
The existing test asserting the offloaded walk returns exactly what the inline walk returns is the right place to extend.

**Backpressure gets simpler, not harder.**
Today the bounded `rows_tx` channel throttles the worker against the upload.
With walk and upload on one task it becomes an ordinary await, and the channel and its capacity constant can go.

## Alternative considered

Let the offloaded walk abandon on divergence and re-run inline.
Rejected.
It walks the data twice, and the restart lands on the timely worker exactly when the result is large, which is precisely the case where blocking the worker hurts most.
It would reintroduce the tail the offload exists to remove, at the worst moment.

## Validation

* Re-run E1 with the stash left on.
The flat tail has to survive, since surviving a production configuration is the entire point of the work.
* Confirm engagement through `mz_index_peek_walks_total{substrate}` rather than inferring it from latency.
* Extend the offloaded-equals-inline walk test to cover a result large enough to cross the stash threshold, on both cursor sources.
* Re-run E4, because the cap now governs upload concurrency and compaction hold duration as well as walk concurrency.
101 changes: 101 additions & 0 deletions doc/developer/design/20260720_two_runtime_compute/peek-placement.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
# Where a peek's walk runs

**Status: parked, deliberately.** The problem is real and M1's evidence is the strongest in
this directory. What is not settled is which remedy to keep, and that is answerable by
experiment on fixtures that already exist. Implementing every candidate, merging them, and
then learning which one was needed is the expensive order.

Nothing here is gated on the interactive runtime, and the interactive runtime does not need
any of it. See [design.md](design.md) for the mechanism decomposition these solutions are
scored against, and CPU-217 for this work's tracking issue.

## Two axes, not one

An earlier draft of this proposed a single `compute_peek_substrate` parameter enumerating
placements, on the argument that a walk runs in exactly one place so the choice is
mutually exclusive. **That argument is wrong.** Where the walk runs and whether it yields
are independent:

| | run to completion | yields |
|---|---|---|
| **on the serving worker** | today's default | S1, cooperative slicing, PR #38040 |
| **off the worker** | S5, a blocking task | coherent, unmeasured, unimplemented |

A walk on a blocking task can perfectly well have yield points. It would gain a place to
observe a cancellation request, which is S2's prerequisite, and a bound on how long one
walk pins the batches its cursor covers. So the fourth cell is not a nonsense state that a
config should forbid, it is a candidate nobody has tried.

The single-selector shape would have made that cell unreachable, which is the failure mode
of collapsing two questions into one name. design.md already had this right, listing S1 and
S5 as separate solutions. Any future parameter follows the axes: placement and preemption
are separate settings, or one setting over the cross, but not one setting over placement
alone.

## What is measured, per cell

| Cell | M1, peek behind peek | M2, peek behind an activation | M8, freshness |
|---|---|---|---|
| on-worker, run to completion | baseline: E1 max 5783.7 ms, E11 58 of 261 slow | baseline: E12 p90 129.5 ms | baseline: E13 2340 ms peak |
| on-worker, yielding (S1) | **predicted only** | not measured | E13 365 ms peak, light write load |
| off-worker (S5) | E1 max 180.4 ms flat, E11 0 of 261, E8b 29151.7 to 152.4 ms | **worse**: E12 p90 148.2 ms | E13 about 101 ms |
| off-worker, yielding | nothing | nothing | nothing |

Two entries in that table decide the question and neither is filled in.

**S1 on M1 is predicted, not measured.** It is the presumed default and the row with the
most cells empty. Everything the decision rests on is an argument about its quantum.

**Yielding does not rescue S5 on M2.** Worth stating because the new cell invites the
hope. E12's regression is dispatch timing: the snapshot is taken in `process_peeks`, which
runs only after `step_or_park` returns, and retirement costs another step. That cost is
paid before the walk starts, so giving the walk yield points cannot recover it. S9,
size- or residency-aware routing, is the candidate that addresses it.

## The experiment that decides

The cheap decisive question is whether S1 alone matches S5 on the three fixtures that
carry S5's case. All three exist and are described in the project document "Interactive read isolation:
experimental evaluation":

* **E1**, a point lookup behind three concurrent scans at walk costs of 23, 190 and 2170 ms.
* **E11**, the skewed point lookup, a hot key holding millions of values under open-loop
arrivals.
* **E8b**, a lookup on a resident index behind swap-resident walks at matched swap depth.

Decision rule, registered before running: **if S1 matches S5 within noise on all three,
S5 has no remaining justification and should be deleted rather than merged.** If S5 wins on
E8b alone, its case is the unattributed swap-walk duration (M11) and nothing else, which is
a much narrower claim than the one it was built on.

This needs PR #38040 and the existing fixtures. It does not need the interactive runtime,
which is why parking costs nothing.

## Why this is parked rather than dropped

The evidence that *something* is needed is not in doubt. E11 is a field-reported shape, and
E8b turned a 29.2 second victim latency into 152 ms. What is in doubt is which of four cells
to keep, and there are two reasons not to answer that by building:

The candidates are substitutes, not complements. S1 and S5 both create a preemption point
for the same walk, so shipping both means carrying two mechanisms where measurement is
expected to retire one. S5 already has no axis on which it is uniquely best: E2 showed it
is not needed for the peek-tail win once the walk moves, E12 measured it behind doing
nothing, and S1 reaches M8 as well.

And a merged mechanism is harder to delete than an unmerged one. S5 is about 700 lines
including `local_snapshot.rs`, `PendingPeek::IndexOffload` and its metric. Landing that to
discover S1 subsumes it means removing it afterwards, from a tree where something may
already depend on it.

## Consequence for the interactive-runtime work

S5's code currently rides on the interactive-runtime branch, where it is a conditional on
top of an inline walk rather than load-bearing: interactive peeks resolve through
`shared_index_peek_response`, and the offload branch is entered only when the parameter
selects it. So it can come out without touching how the interactive runtime serves reads.

Removing it costs the branch E1 and E11 as supporting evidence, which is correct: E2 showed
those results belong to the walk substrate rather than to the second runtime, and the
interactive runtime's own case is E7 and E9. A branch that keeps S5 is a branch arguing for
two things at once, and one of them is the one we just agreed to decide by experiment.
7 changes: 7 additions & 0 deletions doc/user/data/metrics.yml
Original file line number Diff line number Diff line change
Expand Up @@ -1095,6 +1095,13 @@ metrics:
help: Total time processing index peeks, from process_peek entry to response. Excluding peeks that use the peek response stash.
source: src/compute/src/metrics.rs
visibility: internal
- name: mz_index_peek_walks_total
help: The total number of fast-path index peek walks, by the substrate that ran them.
labels:
- substrate
- worker_id
source: src/compute/src/metrics.rs
visibility: internal
- name: mz_kafka_partition_offset_max
help: High watermark offset on broker for partition
labels:
Expand Down
11 changes: 11 additions & 0 deletions misc/python/materialize/mzcompose/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -297,6 +297,16 @@ def get_variable_system_parameters(
"true",
["true", "false"],
),
VariableSystemParameter(
"enable_index_peek_offload",
"true",
["true", "false"],
),
VariableSystemParameter(
"enable_compute_interactive_runtime",
"true",
["true", "false"],
),
VariableSystemParameter(
"enable_upsert_v2",
"false",
Expand Down Expand Up @@ -584,6 +594,7 @@ def get_default_system_parameters(
# all. Only add it in UNINTERESTING_SYSTEM_PARAMETERS if none of the above
# apply.
UNINTERESTING_SYSTEM_PARAMETERS = [
"index_peek_offload_max_inflight",
"enable_compute_half_join2",
"enable_mz_join_core",
"linear_join_yielding",
Expand Down
18 changes: 17 additions & 1 deletion misc/python/materialize/mzcompose/services/clusterd.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ def __init__(
workers: int = 1,
process_names: list[str] = [],
mz_service: str = "materialized",
interactive_compute: bool = False,
) -> None:
environment = [
"CLUSTERD_LOG_FILTER",
Expand Down Expand Up @@ -78,6 +79,21 @@ def __init__(
f"CLUSTERD_STORAGE_TIMELY_CONFIG={storage_timely_config}",
]

# When set, clusterd runs a second, interactive compute runtime alongside the
# maintenance one (see `--interactive-compute-timely-config` in
# `src/clusterd/src/lib.rs`). It must span the same number of Timely peers as
# the maintenance compute config, so it reuses `process_names`/`workers`; its
# addresses use a distinct port (2104) so the two runtimes don't collide.
ports = [2100, 2101, 6878]
if interactive_compute:
interactive_compute_timely_config = timely_config(
process_names, 2104, workers, DEFAULT_COMPUTE_EXERT_PROPORTIONALITY
)
environment += [
f"CLUSTERD_INTERACTIVE_COMPUTE_TIMELY_CONFIG={interactive_compute_timely_config}"
]
ports += [2104]

options = ["clusterd", f"--scratch-directory={scratch_directory}", *options]

config: ServiceConfig = {}
Expand Down Expand Up @@ -106,7 +122,7 @@ def __init__(
config.update(
{
"command": options,
"ports": [2100, 2101, 6878],
"ports": ports,
"environment": environment,
"volumes": volumes or DEFAULT_MZ_VOLUMES,
"restart": restart,
Expand Down
Loading
Loading