mz_compute_types/plan.rs
1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! An explicit representation of a rendering plan for provided dataflows.
11
12#![warn(missing_debug_implementations)]
13
14use std::collections::{BTreeMap, BTreeSet};
15
16use columnar::Columnar;
17use mz_expr::{
18 CollectionPlan, EvalError, Id, LetRecLimit, LocalId, MapFilterProject, MfpPlan,
19 OptimizedMirRelationExpr, SafeMfpPlan, TableFunc,
20};
21use mz_ore::metric;
22use mz_ore::metrics::MetricsRegistry;
23use mz_ore::metrics::raw::IntCounterVec;
24use mz_ore::soft_assert_eq_no_log;
25use mz_ore::str::Indent;
26use mz_repr::explain::text::text_string_at;
27use mz_repr::explain::{DummyHumanizer, ExplainConfig, ExprHumanizer, PlanRenderingContext};
28use mz_repr::optimize::OptimizerFeatures;
29use mz_repr::{Diff, GlobalId, StableRow, Timestamp};
30use serde::{Deserialize, Serialize};
31
32use crate::dataflows::DataflowDescription;
33use crate::plan::join::JoinPlan;
34use crate::plan::reduce::{KeyValPlan, ReducePlan};
35use crate::plan::scalar::LirScalarExpr;
36use crate::plan::threshold::ThresholdPlan;
37use crate::plan::top_k::TopKPlan;
38use crate::plan::transform::{Transform, TransformConfig};
39
40mod lowering;
41
42pub mod interpret;
43pub mod join;
44pub mod reduce;
45pub mod render_plan;
46pub mod scalar;
47pub mod threshold;
48pub mod top_k;
49pub mod transform;
50
51/// Metrics collected during MIR to LIR lowering.
52#[derive(Debug, Clone)]
53pub struct LoweringMetrics {
54 /// Counts non-`None` results of `MapFilterProject::literal_constraints` during lowering,
55 /// labeled by the call site (`"get"` or `"mfp"`).
56 literal_constraints: IntCounterVec,
57}
58
59impl LoweringMetrics {
60 /// Registers the lowering metrics into `registry`.
61 pub fn register_into(registry: &MetricsRegistry) -> Self {
62 Self {
63 literal_constraints: registry.register(metric!(
64 name: "mz_optimizer_lowering_literal_constraints_total",
65 help: "How often the MFP-based literal-constraint detector succeeded, by call site.",
66 var_labels: ["case"],
67 )),
68 }
69 }
70
71 /// Records that a `literal_constraints` call at `case` produced a usable constraint.
72 pub fn inc_literal_constraints(&self, case: &str) {
73 self.literal_constraints.with_label_values(&[case]).inc();
74 }
75}
76
77/// The forms in which an operator's output is available.
78///
79/// These forms may include "raw", meaning as a streamed collection, but also any
80/// number of "arranged" representations.
81///
82/// Each arranged representation is described by a `(to_key, permutation, thinning)`
83/// triple, built by `permutation_for_arrangement`. `to_key` (length `K`) are the key
84/// expressions over a row. `permutation` (length `A`, the raw/unthinned arity) maps
85/// each row column to its position in the `(key, value)` concatenation. `thinning`
86/// (length `M`) lists the row columns that form the value, in value order, so a value
87/// datum at concatenation position `c >= K` came from row column `thinning[c - K]`.
88///
89/// This triple is unrelated to `KeyValRowMapping`'s `(to_key, to_val, to_row)` fields
90/// despite the visual resemblance: `permutation` here plays the role of
91/// `KeyValRowMapping::to_row`, and `thinning` plays the role of
92/// `KeyValRowMapping::to_val`. Do not assume the same field order.
93#[derive(
94 Clone,
95 Debug,
96 Default,
97 Deserialize,
98 Eq,
99 Ord,
100 PartialEq,
101 PartialOrd,
102 Serialize
103)]
104pub struct AvailableCollections {
105 /// Whether the collection exists in unarranged form.
106 pub raw: bool,
107 /// The list of available arrangements, each a `(to_key, permutation, thinning)`
108 /// triple. See the struct-level documentation for field semantics.
109 pub arranged: Vec<(Vec<LirScalarExpr>, Vec<usize>, Vec<usize>)>,
110}
111
112impl AvailableCollections {
113 /// Represent a collection that has no arrangements.
114 pub fn new_raw() -> Self {
115 Self {
116 raw: true,
117 arranged: Vec::new(),
118 }
119 }
120
121 /// Represent a collection that is arranged in the specified ways.
122 pub fn new_arranged(arranged: Vec<(Vec<LirScalarExpr>, Vec<usize>, Vec<usize>)>) -> Self {
123 assert!(
124 !arranged.is_empty(),
125 "Invariant violated: at least one collection must exist"
126 );
127 Self {
128 raw: false,
129 arranged,
130 }
131 }
132
133 /// Get some arrangement, if one exists.
134 pub fn arbitrary_arrangement(&self) -> Option<&(Vec<LirScalarExpr>, Vec<usize>, Vec<usize>)> {
135 assert!(
136 self.raw || !self.arranged.is_empty(),
137 "Invariant violated: at least one collection must exist"
138 );
139 self.arranged.get(0)
140 }
141}
142
143/// How to render the arrangements requested by an `ArrangeBy`.
144///
145/// Decided during LIR lowering and consumed by the renderer. The variant says what the
146/// renderer will do, not what it knows about the input.
147#[derive(
148 Clone,
149 Copy,
150 Debug,
151 Deserialize,
152 Eq,
153 Ord,
154 PartialEq,
155 PartialOrd,
156 Serialize
157)]
158pub enum ArrangementStrategy {
159 /// Form arrangements directly from the input collection.
160 Direct,
161 /// Insert temporal bucketing in front of the arrangement, to delay future-stamped
162 /// updates (e.g., from `mz_now()` MFPs) until their bucket boundary releases them.
163 /// Honoured only when `ENABLE_COMPUTE_TEMPORAL_BUCKETING` is set; otherwise behaves like
164 /// `Direct`.
165 TemporalBucketing,
166}
167
168impl std::fmt::Display for ArrangementStrategy {
169 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
170 match self {
171 ArrangementStrategy::Direct => write!(f, "Direct"),
172 ArrangementStrategy::TemporalBucketing => write!(f, "TemporalBucketing"),
173 }
174 }
175}
176
177/// An identifier for an LIR node.
178#[derive(
179 Clone,
180 Copy,
181 Debug,
182 Deserialize,
183 Eq,
184 Ord,
185 PartialEq,
186 PartialOrd,
187 Serialize,
188 Columnar
189)]
190pub struct LirId(u64);
191
192impl LirId {
193 fn as_u64(&self) -> u64 {
194 self.0
195 }
196}
197
198impl From<LirId> for u64 {
199 fn from(value: LirId) -> Self {
200 value.as_u64()
201 }
202}
203
204impl std::fmt::Display for LirId {
205 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
206 write!(f, "{}", self.0)
207 }
208}
209
210/// Version of the stable LIR serialization format.
211///
212/// Once a version has shipped, meaning a released Materialize durably stores
213/// pinned plans in it, bump this when the serialized representation of
214/// [`LirRelationExpr`] or anything it transitively contains changes, or when
215/// a scalar function's declared properties change. Whether the current
216/// version has shipped is recorded in [`LIR_VERSION_POLICY`]. Two snapshot
217/// tests key off it: the schema snapshot in `tests/lir_schema.rs` against
218/// `tests/snapshots/lir_v{LIR_VERSION}.json`, and the function property
219/// registry in `tests/func_registry.rs` against
220/// `tests/snapshots/func_registry.json` with a per-version digest in
221/// `tests/snapshots/func_registry_digests.json`.
222pub const LIR_VERSION: u64 = 1;
223
224/// What a change to the stable LIR format must do at the current
225/// [`LIR_VERSION`]. The snapshot tests print it when a snapshot goes stale.
226///
227/// This is the one place that records whether the current version has
228/// shipped. Once a released Materialize durably stores pinned plans, replace
229/// the text with the shipped policy: bump `LIR_VERSION` rather than rewriting
230/// the shipped version's snapshots, so the old schema stays available to
231/// migration tooling.
232pub const LIR_VERSION_POLICY: &str = "\
233The current LIR version has NOT shipped: no released Materialize stores pinned
234LIR plans yet, so nothing depends on the checked-in snapshots. Regenerate them
235in place and include the diff in your PR.
236
237Once pinned plans are durably stored, a change like this must instead bump
238LIR_VERSION in src/compute-types/src/plan.rs so it lands as a new version.
239See doc/developer/design/20260826_pinned_lir.md.";
240
241pub use constant_rows_serde::ConstantRows;
242
243/// Serializes `LirRelationNode::Constant`'s rows through the named
244/// [`ConstantRows`] mirror enum instead of std `Result`.
245///
246/// The stable LIR schema registry maps each container name to a single
247/// format, and `Result` would clash with the differently instantiated
248/// `Result` in `LirScalarExpr::Literal`. The mirror has the same variant
249/// order as `Result`, so the encoded bytes are unchanged.
250mod constant_rows_serde {
251 use mz_expr::{EvalError, StableEvalError, StableEvalErrorRef};
252 use mz_repr::{Diff, StableRow, Timestamp};
253 use serde::{Deserialize, Deserializer, Serialize, Serializer};
254
255 /// The serialized form of `LirRelationNode::Constant`'s rows.
256 #[derive(Debug, Serialize, Deserialize)]
257 pub enum ConstantRows {
258 /// See `Result::Ok`.
259 Ok(Vec<(StableRow, Timestamp, Diff)>),
260 /// See `Result::Err`.
261 Err(StableEvalError),
262 }
263
264 /// Borrowing mirror of [`ConstantRows`], to serialize without cloning.
265 #[derive(Serialize)]
266 #[serde(rename = "ConstantRows")]
267 enum ConstantRowsRef<'a> {
268 Ok(&'a Vec<(StableRow, Timestamp, Diff)>),
269 Err(StableEvalErrorRef<'a>),
270 }
271
272 pub fn serialize<S: Serializer>(
273 rows: &Result<Vec<(StableRow, Timestamp, Diff)>, EvalError>,
274 serializer: S,
275 ) -> Result<S::Ok, S::Error> {
276 let mirror = match rows {
277 Ok(rows) => ConstantRowsRef::Ok(rows),
278 Err(err) => ConstantRowsRef::Err(StableEvalErrorRef(err)),
279 };
280 mirror.serialize(serializer)
281 }
282
283 pub fn deserialize<'de, D: Deserializer<'de>>(
284 deserializer: D,
285 ) -> Result<Result<Vec<(StableRow, Timestamp, Diff)>, EvalError>, D::Error> {
286 Ok(match ConstantRows::deserialize(deserializer)? {
287 ConstantRows::Ok(rows) => Ok(rows),
288 ConstantRows::Err(err) => Err(err.0),
289 })
290 }
291}
292
293/// A rendering plan with as much conditional logic as possible removed.
294#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
295pub struct LirRelationExpr {
296 /// A dataflow-local identifier.
297 pub lir_id: LirId,
298 /// The underlying operator.
299 pub node: LirRelationNode,
300}
301
302/// The actual AST node of the `LirRelationExpr`.
303#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
304pub enum LirRelationNode {
305 /// A collection containing a pre-determined collection.
306 Constant {
307 /// Explicit update triples for the collection.
308 #[serde(with = "constant_rows_serde")]
309 rows: Result<Vec<(StableRow, Timestamp, Diff)>, EvalError>,
310 },
311 /// A reference to a bound collection.
312 ///
313 /// This is commonly either an external reference to an existing source or
314 /// maintained arrangement, or an internal reference to a `Let` identifier.
315 Get {
316 /// A global or local identifier naming the collection.
317 id: Id,
318 /// Arrangements that will be available.
319 ///
320 /// The collection will also be loaded if available, which it will
321 /// not be for imported data, but which it may be for locally defined
322 /// data.
323 // TODO: Be more explicit about whether a collection is available,
324 // although one can always produce it from an arrangement, and it
325 // seems generally advantageous to do that instead (to avoid cloning
326 // rows, by using `mfp` first on borrowed data).
327 keys: AvailableCollections,
328 /// The actions to take when introducing the collection.
329 plan: GetPlan,
330 },
331 /// Binds `value` to `id`, and then results in `body` with that binding.
332 ///
333 /// This stage has the effect of sharing `value` across multiple possible
334 /// uses in `body`, and is the only mechanism we have for sharing collection
335 /// information across parts of a dataflow.
336 ///
337 /// The binding is not available outside of `body`.
338 Let {
339 /// The local identifier to be used, available to `body` as `Id::Local(id)`.
340 id: LocalId,
341 /// The collection that should be bound to `id`.
342 value: Box<LirRelationExpr>,
343 /// The collection that results, which is allowed to contain `Get` stages
344 /// that reference `Id::Local(id)`.
345 body: Box<LirRelationExpr>,
346 },
347 /// Binds `values` to `ids`, evaluates them potentially recursively, and returns `body`.
348 ///
349 /// All bindings are available to all bindings, and to `body`.
350 /// The contents of each binding are initially empty, and then updated through a sequence
351 /// of iterations in which each binding is updated in sequence, from the most recent values
352 /// of all bindings.
353 LetRec {
354 /// The local identifiers to be used, available to `body` as `Id::Local(id)`.
355 ids: Vec<LocalId>,
356 /// The collection that should be bound to `id`.
357 values: Vec<LirRelationExpr>,
358 /// Maximum number of iterations. See further info on the MIR `LetRec`.
359 limits: Vec<Option<LetRecLimit>>,
360 /// The collection that results, which is allowed to contain `Get` stages
361 /// that reference `Id::Local(id)`.
362 body: Box<LirRelationExpr>,
363 },
364 /// Map, Filter, and Project operators.
365 ///
366 /// This stage contains work that we would ideally like to fuse to other plan
367 /// stages, but for practical reasons cannot. For example: threshold, topk,
368 /// and sometimes reduce stages are not able to absorb this operator.
369 Mfp {
370 /// The input collection.
371 input: Box<LirRelationExpr>,
372 /// Linear operator to apply to each record.
373 mfp: MfpPlan<LirScalarExpr>,
374 /// Whether the input is from an arrangement, and if so,
375 /// whether we can seek to a specific value therein
376 input_key_val: Option<(Vec<LirScalarExpr>, Option<StableRow>)>,
377 },
378 /// A variable number of output records for each input record.
379 ///
380 /// This stage is a bit of a catch-all for logic that does not easily fit in
381 /// map stages. This includes table valued functions, but also functions of
382 /// multiple arguments, and functions that modify the sign of updates.
383 ///
384 /// This stage allows a `MapFilterProject` operator to be fused to its output,
385 /// and this can be very important as otherwise the output of `func` is just
386 /// appended to the input record, for as many outputs as it has. This has the
387 /// unpleasant default behavior of repeating potentially large records that
388 /// are being unpacked, producing quadratic output in those cases. Instead,
389 /// in these cases use a `mfp` member that projects away these large fields.
390 FlatMap {
391 /// The particular arrangement of the input we expect to use,
392 /// if any
393 input_key: Option<Vec<LirScalarExpr>>,
394 /// The input collection.
395 input: Box<LirRelationExpr>,
396 /// Expressions that for each row prepare the arguments to `func`.
397 exprs: Vec<LirScalarExpr>,
398 /// The variable-record emitting function.
399 func: TableFunc,
400 /// Linear operator to apply to each record produced by `func`.
401 mfp_after: MfpPlan<LirScalarExpr>,
402 },
403 /// A multiway relational equijoin, with fused map, filter, and projection.
404 ///
405 /// This stage performs a multiway join among `inputs`, using the equality
406 /// constraints expressed in `plan`. The plan also describes the implementation
407 /// strategy we will use, and any pushed down per-record work.
408 Join {
409 /// An ordered list of inputs that will be joined.
410 inputs: Vec<LirRelationExpr>,
411 /// Detailed information about the implementation of the join.
412 ///
413 /// This includes information about the implementation strategy, but also
414 /// any map, filter, project work that we might follow the join with, but
415 /// potentially pushed down into the implementation of the join.
416 plan: JoinPlan,
417 },
418 /// Aggregation by key.
419 Reduce {
420 /// The particular arrangement of the input we expect to use,
421 /// if any
422 input_key: Option<Vec<LirScalarExpr>>,
423 /// The input collection.
424 input: Box<LirRelationExpr>,
425 /// A plan for changing input records into key, value pairs.
426 key_val_plan: KeyValPlan,
427 /// A plan for performing the reduce.
428 ///
429 /// The implementation of reduction has several different strategies based
430 /// on the properties of the reduction, and the input itself. Please check
431 /// out the documentation for this type for more detail.
432 plan: ReducePlan,
433 /// An MFP that must be applied to results. The projection part of this
434 /// MFP must preserve the key for the reduction; otherwise, the results
435 /// become undefined. Additionally, the MFP is guaranteed to be free from
436 /// temporal predicates so that it can be readily evaluated.
437 mfp_after: SafeMfpPlan<LirScalarExpr>,
438 /// Strategy for forming the internal input arrangement built by `Reduce`
439 /// (materialized via `key_val_plan`).
440 ///
441 /// Set by the lowering from the input's `has_future_updates` flag. The
442 /// renderer applies it to the keyed `(key, val)` stream feeding the
443 /// reduce. See `render_reduce` for the rationale on why this is
444 /// plumbed through `Reduce` rather than handled at the arrangement site.
445 ///
446 /// Note: unrelated to the hash buckets used by hierarchical reductions
447 /// (e.g. `ReducePlan::Hierarchical`'s `buckets`), which are an internal
448 /// sharding scheme for `min`/`max`-style aggregations. Here "bucketing"
449 /// refers exclusively to temporal (time-domain) bucketing of
450 /// future-stamped updates.
451 temporal_bucketing_strategy: ArrangementStrategy,
452 },
453 /// Key-based "Top K" operator, retaining the first K records in each group.
454 TopK {
455 /// The input collection.
456 input: Box<LirRelationExpr>,
457 /// A plan for performing the Top-K.
458 ///
459 /// The implementation of reduction has several different strategies based
460 /// on the properties of the reduction, and the input itself. Please check
461 /// out the documentation for this type for more detail.
462 top_k_plan: TopKPlan,
463 /// Strategy for bucketing the input collection ahead of the Top-K operator.
464 ///
465 /// Set by the lowering from the input's `has_future_updates` flag. The
466 /// renderer applies it to the per-row input stream at the top of
467 /// `render_topk`, covering all three `TopKPlan` arms uniformly. See
468 /// `LirRelationNode::Reduce::temporal_bucketing_strategy` for the underlying
469 /// convention.
470 temporal_bucketing_strategy: ArrangementStrategy,
471 },
472 /// Inverts the sign of each update.
473 Negate {
474 /// The input collection.
475 input: Box<LirRelationExpr>,
476 },
477 /// Filters records that accumulate negatively.
478 ///
479 /// Although the operator suppresses updates, it is a stateful operator taking
480 /// resources proportional to the number of records with non-zero accumulation.
481 Threshold {
482 /// The input collection.
483 input: Box<LirRelationExpr>,
484 /// A plan for performing the threshold.
485 ///
486 /// The implementation of reduction has several different strategies based
487 /// on the properties of the reduction, and the input itself. Please check
488 /// out the documentation for this type for more detail.
489 threshold_plan: ThresholdPlan,
490 },
491 /// Adds the contents of the input collections.
492 ///
493 /// Importantly, this is *multiset* union, so the multiplicities of records will
494 /// add. This is in contrast to *set* union, where the multiplicities would be
495 /// capped at one. A set union can be formed with `Union` followed by `Reduce`
496 /// implementing the "distinct" operator.
497 Union {
498 /// The input collections
499 inputs: Vec<LirRelationExpr>,
500 /// Whether to consolidate the output, e.g., cancel negated records.
501 consolidate_output: bool,
502 /// Per-input bucketing strategies. Lockstep with `inputs`: index `i` is the
503 /// strategy applied to `inputs[i]` before concatenation.
504 ///
505 /// Set by the lowering from each input's `has_future_updates` flag. Only
506 /// consolidating Unions (`consolidate_output: true`) carry non-`Direct`
507 /// entries, because bucketing only pays off ahead of a consolidating
508 /// downstream operator. See `LirRelationNode::Reduce::temporal_bucketing_strategy`
509 /// for the underlying convention.
510 temporal_bucketing_strategies: Vec<ArrangementStrategy>,
511 },
512 /// The `input` plan, but with additional arrangements.
513 ///
514 /// This operator does not change the logical contents of `input`, but ensures
515 /// that certain arrangements are available in the results. This operator can
516 /// be important for e.g. the `Join` stage which benefits from multiple arrangements
517 /// or to cap a `LirRelationExpr` so that indexes can be exported.
518 ArrangeBy {
519 /// The key that must be used to access the input.
520 input_key: Option<Vec<LirScalarExpr>>,
521 /// The input collection.
522 input: Box<LirRelationExpr>,
523 /// The MFP that must be applied to the input.
524 input_mfp: MfpPlan<LirScalarExpr>,
525 /// A list of arrangement keys, and possibly a raw collection,
526 /// that will be added to those of the input. Does not include
527 /// any other existing arrangements.
528 forms: AvailableCollections,
529 /// How the renderer should form the arrangements requested by `forms`.
530 strategy: ArrangementStrategy,
531 },
532}
533
534impl LirRelationNode {
535 /// Iterates through references to child expressions.
536 pub fn children(&self) -> impl Iterator<Item = &LirRelationExpr> {
537 let mut first = None;
538 let mut second = None;
539 let mut rest = None;
540 let mut last = None;
541
542 use LirRelationNode::*;
543 match self {
544 Constant { .. } | Get { .. } => (),
545 Let { value, body, .. } => {
546 first = Some(&**value);
547 second = Some(&**body);
548 }
549 LetRec { values, body, .. } => {
550 rest = Some(values);
551 last = Some(&**body);
552 }
553 Mfp { input, .. }
554 | FlatMap { input, .. }
555 | Reduce { input, .. }
556 | TopK { input, .. }
557 | Negate { input, .. }
558 | Threshold { input, .. }
559 | ArrangeBy { input, .. } => {
560 first = Some(&**input);
561 }
562 Join { inputs, .. } | Union { inputs, .. } => {
563 rest = Some(inputs);
564 }
565 }
566
567 first
568 .into_iter()
569 .chain(second)
570 .chain(rest.into_iter().flatten())
571 .chain(last)
572 }
573
574 /// Iterates through mutable references to child expressions.
575 pub fn children_mut(&mut self) -> impl Iterator<Item = &mut LirRelationExpr> {
576 let mut first = None;
577 let mut second = None;
578 let mut rest = None;
579 let mut last = None;
580
581 use LirRelationNode::*;
582 match self {
583 Constant { .. } | Get { .. } => (),
584 Let { value, body, .. } => {
585 first = Some(&mut **value);
586 second = Some(&mut **body);
587 }
588 LetRec { values, body, .. } => {
589 rest = Some(values);
590 last = Some(&mut **body);
591 }
592 Mfp { input, .. }
593 | FlatMap { input, .. }
594 | Reduce { input, .. }
595 | TopK { input, .. }
596 | Negate { input, .. }
597 | Threshold { input, .. }
598 | ArrangeBy { input, .. } => {
599 first = Some(&mut **input);
600 }
601 Join { inputs, .. } | Union { inputs, .. } => {
602 rest = Some(inputs);
603 }
604 }
605
606 first
607 .into_iter()
608 .chain(second)
609 .chain(rest.into_iter().flatten())
610 .chain(last)
611 }
612}
613
614impl LirRelationNode {
615 /// Attach an `lir_id` to a `LirRelationNode` to make a complete `LirRelationExpr`.
616 pub fn as_plan(self, lir_id: LirId) -> LirRelationExpr {
617 LirRelationExpr { lir_id, node: self }
618 }
619}
620
621impl LirRelationExpr {
622 /// Pretty-print this [LirRelationExpr] to a string.
623 pub fn pretty(&self) -> String {
624 let config = ExplainConfig::default();
625 self.debug_explain(&config, None)
626 }
627
628 /// Pretty-print this [LirRelationExpr] to a string using a custom
629 /// [ExplainConfig] and an optionally provided [ExprHumanizer].
630 /// This is intended for debugging and tests, not users.
631 pub fn debug_explain(
632 &self,
633 config: &ExplainConfig,
634 humanizer: Option<&dyn ExprHumanizer>,
635 ) -> String {
636 text_string_at(self, || PlanRenderingContext {
637 indent: Indent::default(),
638 humanizer: humanizer.unwrap_or(&DummyHumanizer),
639 annotations: BTreeMap::default(),
640 config,
641 ambiguous_ids: BTreeSet::default(),
642 })
643 }
644}
645
646/// How a `Get` stage will be rendered.
647#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq, Ord, PartialOrd)]
648pub enum GetPlan {
649 /// Simply pass input arrangements on to the next stage.
650 PassArrangements,
651 /// Using the supplied key, optionally seek the row, and apply the MFP.
652 Arrangement(
653 Vec<LirScalarExpr>,
654 Option<StableRow>,
655 MfpPlan<LirScalarExpr>,
656 ),
657 /// Scan the input collection (unarranged) and apply the MFP.
658 Collection(MfpPlan<LirScalarExpr>),
659}
660
661impl LirRelationExpr {
662 /// Convert the dataflow description into one that uses render plans.
663 #[mz_ore::instrument(
664 target = "optimizer",
665 level = "debug",
666 fields(path.segment = "finalize_dataflow")
667 )]
668 pub fn finalize_dataflow(
669 desc: DataflowDescription<OptimizedMirRelationExpr>,
670 features: &OptimizerFeatures,
671 metrics: Option<&LoweringMetrics>,
672 ) -> Result<DataflowDescription<Self>, String> {
673 fail::fail_point!("finalize_dataflow");
674
675 // First, we lower the dataflow description from MIR to LIR. Lowering
676 // also moves common parts of the MFPs pushed onto each source's reads into the source
677 // itself (see `Context::refine_source_mfps`).
678 let mut dataflow = Self::lower_dataflow(desc, features, metrics)?;
679
680 // Note: `consolidate_output` for `Union` and per-input
681 // `temporal_bucketing_strategies` are decided at lowering time (see the
682 // `Union` arm of `lower_mir_expr_stack_safe`). The pre-existing
683 // `refine_union_negate_consolidation` pass — which used to flip
684 // `consolidate_output` to `true` for Unions with a `Negate` child — has
685 // been folded into the lowering, since lowering is the only point where
686 // the bucketing decision (which depends on `has_future_updates`) is
687 // available.
688
689 if dataflow.is_single_time() {
690 // The relaxation of the `must_consolidate` flag performs an LIR-based
691 // analysis and transform under checked recursion. By a similar argument
692 // made in `from_mir`, we do not expect the recursion limit to be hit.
693 // However, if that happens, we propagate an error to the caller.
694 // To apply the transform, we first obtain monotonic source and index
695 // global IDs and add them to a `TransformConfig` instance.
696 let monotonic_ids = dataflow
697 .source_imports
698 .iter()
699 .filter_map(|(id, source_import)| source_import.monotonic.then_some(*id))
700 .chain(
701 dataflow
702 .index_imports
703 .iter()
704 .filter_map(|(_id, index_import)| {
705 if index_import.monotonic {
706 Some(index_import.desc.on_id)
707 } else {
708 None
709 }
710 }),
711 )
712 .collect::<BTreeSet<_>>();
713
714 let config = TransformConfig { monotonic_ids };
715 Self::refine_single_time_consolidation(&mut dataflow, &config)?;
716
717 // For non-recursive delta joins in single-time dataflows, only the delta path for the
718 // first relation produces updates: the other paths discard updates at the as-of, which
719 // is the only time present. We keep just that path and, where the first input's
720 // arrangement existed solely to seed it, drop the arrangement and consume the input as a
721 // raw collection. This requires rewriting the path's initial closure to address the raw
722 // row layout rather than the arranged `(key, value)` layout.
723 for build_desc in dataflow.objects_to_build.iter_mut() {
724 // Worklist of plan nodes. `LetRec` bodies are explored but `LetRec` values are not,
725 // which excludes recursive (WMR) joins from this transform.
726 let mut todo = vec![&mut build_desc.plan];
727 while let Some(expr) = todo.pop() {
728 match &mut expr.node {
729 // TODO: also handle binary differential joins, which can likewise shed a
730 // bespoke arrangement on their first input.
731 LirRelationNode::Join {
732 inputs,
733 plan: JoinPlan::Delta(plan),
734 } => {
735 // Only the first relation's path survives at a single time.
736 plan.path_plans.truncate(1);
737
738 let source_relation = plan.path_plans[0].source_relation;
739 // Replace the source input's bespoke arrangement with a raw collection,
740 // but only when the surviving path's source is fed by an `ArrangeBy`
741 // that exists solely to build that arrangement. A source backed directly
742 // by an arranged import has no `ArrangeBy` node here, so this guard skips
743 // it and the path keeps reading it arranged.
744 if let Some(source_key) = plan.path_plans[0].source_key.clone() {
745 if let LirRelationNode::ArrangeBy { forms, .. } =
746 &mut inputs[source_relation].node
747 {
748 // Drop arrangement forms other than the source key, which the
749 // remaining path no longer needs.
750 forms.arranged.retain(|(key, _, _)| key == &source_key);
751 if let Some((to_key, permutation, thinning)) =
752 forms.arranged.pop()
753 {
754 // Make the input a raw collection and unset the source key.
755 // Clearing every arrangement form is safe: `truncate(1)`
756 // already dropped the sibling paths that were the only other
757 // consumers, and this `ArrangeBy` is private to this join
758 // input. What remains is the input's raw collection.
759 forms.raw = true;
760 forms.arranged.clear();
761 plan.path_plans[0].source_key = None;
762
763 // The initial closure addresses the arranged `(key, value)`
764 // layout: columns `[0, K)` are key datums and columns
765 // `[K, K + M)` are the thinned value datums. We rewrite it to
766 // address the raw row instead.
767 //
768 // `to_key` (length `K`) are the key expressions over a row.
769 // `permutation` (length `A`, the raw arity) maps each row
770 // column to its position in the `(key, value)` concatenation.
771 // `thinning` (length `M`) lists the row columns that form the
772 // value.
773 let key_len = to_key.len();
774 let row_arity = permutation.len();
775 let closure = &mut plan.path_plans[0].initial_closure;
776
777 // Step 1: rewrite `ready_equivalences`, which reference the
778 // arranged layout. A key column becomes its defining
779 // expression. A value column becomes the row column it was
780 // projected from.
781 for class in closure.ready_equivalences.iter_mut() {
782 for expr in class.iter_mut() {
783 let mut todo = vec![expr];
784 while let Some(expr) = todo.pop() {
785 if let LirScalarExpr::Column(c, _) = expr {
786 if let Some(key_expr) = to_key.get(*c) {
787 *expr = key_expr.clone();
788 } else {
789 *c = thinning[*c - key_len];
790 }
791 } else {
792 todo.extend(expr.children_mut());
793 }
794 }
795 }
796 }
797
798 // Step 2: rewrite the `before` MFP. Starting from a raw row,
799 // materialize the key datums and project to the arranged
800 // `(key, value)` layout the original MFP expects, then apply
801 // it.
802 let (m, f, p) = closure.before.as_map_filter_project();
803 let mfp = MapFilterProject::new(row_arity)
804 .map(to_key)
805 .project(
806 (row_arity..row_arity + key_len).chain(thinning),
807 )
808 .map(m)
809 .filter(f)
810 .project(p);
811 closure.before =
812 mfp.into_plan().unwrap().into_nontemporal().unwrap();
813 }
814 }
815 }
816
817 todo.extend(inputs.iter_mut());
818 }
819 LirRelationNode::LetRec { body, .. } => {
820 todo.push(body);
821 }
822 x => {
823 todo.extend(x.children_mut());
824 }
825 }
826 }
827 }
828 }
829
830 soft_assert_eq_no_log!(dataflow.check_invariants(), Ok(()));
831
832 mz_repr::explain::trace_plan(&dataflow);
833
834 Ok(dataflow)
835 }
836
837 /// Lowers the dataflow description from MIR to LIR. To this end, the
838 /// method collects all available arrangements and based on this information
839 /// creates plans for every object to be built for the dataflow.
840 #[mz_ore::instrument(
841 target = "optimizer",
842 level = "debug",
843 fields(path.segment ="mir_to_lir")
844 )]
845 fn lower_dataflow(
846 desc: DataflowDescription<OptimizedMirRelationExpr>,
847 features: &OptimizerFeatures,
848 metrics: Option<&LoweringMetrics>,
849 ) -> Result<DataflowDescription<Self>, String> {
850 let context = lowering::Context::new(desc.debug_name.clone(), features, metrics);
851 let dataflow = context.lower(desc)?;
852
853 mz_repr::explain::trace_plan(&dataflow);
854
855 Ok(dataflow)
856 }
857
858 /// Refines the plans of objects to be built as part of a single-time `dataflow` to relax
859 /// the setting of the `must_consolidate` attribute of monotonic operators, if necessary,
860 /// whenever the input is deemed to be physically monotonic.
861 #[mz_ore::instrument(
862 target = "optimizer",
863 level = "debug",
864 fields(path.segment = "refine_single_time_consolidation")
865 )]
866 fn refine_single_time_consolidation(
867 dataflow: &mut DataflowDescription<Self>,
868 config: &TransformConfig,
869 ) -> Result<(), String> {
870 // We should only reach here if we have a one-shot SELECT query, i.e.,
871 // a single-time dataflow.
872 assert!(dataflow.is_single_time());
873
874 let transform = transform::RelaxMustConsolidate;
875 for build_desc in dataflow.objects_to_build.iter_mut() {
876 transform
877 .transform(config, &mut build_desc.plan)
878 .map_err(|_| "Maximum recursion limit error in consolidation relaxation.")?;
879 }
880 mz_repr::explain::trace_plan(dataflow);
881 Ok(())
882 }
883}
884
885impl CollectionPlan for LirRelationNode {
886 fn depends_on_into(&self, out: &mut BTreeSet<GlobalId>) {
887 match self {
888 LirRelationNode::Constant { rows: _ } => (),
889 LirRelationNode::Get {
890 id,
891 keys: _,
892 plan: _,
893 } => match id {
894 Id::Global(id) => {
895 out.insert(*id);
896 }
897 Id::Local(_) => (),
898 },
899 LirRelationNode::Let { id: _, value, body } => {
900 value.depends_on_into(out);
901 body.depends_on_into(out);
902 }
903 LirRelationNode::LetRec {
904 ids: _,
905 values,
906 limits: _,
907 body,
908 } => {
909 for value in values.iter() {
910 value.depends_on_into(out);
911 }
912 body.depends_on_into(out);
913 }
914 LirRelationNode::Join { inputs, plan: _ }
915 | LirRelationNode::Union {
916 inputs,
917 consolidate_output: _,
918 temporal_bucketing_strategies: _,
919 } => {
920 for input in inputs {
921 input.depends_on_into(out);
922 }
923 }
924 LirRelationNode::Mfp {
925 input,
926 mfp: _,
927 input_key_val: _,
928 }
929 | LirRelationNode::FlatMap {
930 input_key: _,
931 input,
932 exprs: _,
933 func: _,
934 mfp_after: _,
935 }
936 | LirRelationNode::ArrangeBy {
937 input_key: _,
938 input,
939 input_mfp: _,
940 forms: _,
941 strategy: _,
942 }
943 | LirRelationNode::Reduce {
944 input_key: _,
945 input,
946 key_val_plan: _,
947 plan: _,
948 mfp_after: _,
949 temporal_bucketing_strategy: _,
950 }
951 | LirRelationNode::TopK {
952 input,
953 top_k_plan: _,
954 temporal_bucketing_strategy: _,
955 }
956 | LirRelationNode::Negate { input }
957 | LirRelationNode::Threshold {
958 input,
959 threshold_plan: _,
960 } => {
961 input.depends_on_into(out);
962 }
963 }
964 }
965}
966
967impl CollectionPlan for LirRelationExpr {
968 fn depends_on_into(&self, out: &mut BTreeSet<GlobalId>) {
969 self.node.depends_on_into(out);
970 }
971}
972
973/// Returns bucket sizes, descending, suitable for hierarchical decomposition of an operator, based
974/// on the expected number of rows that will have the same group key.
975fn bucketing_of_expected_group_size(expected_group_size: Option<u64>) -> Vec<u64> {
976 // NOTE(vmarcos): The fan-in of 16 defined below is used in the tuning advice built-in view
977 // mz_introspection.mz_expected_group_size_advice.
978 let mut buckets = vec![];
979 let mut current = 16;
980
981 // Plan for 4B records in the expected case if the user didn't specify a group size.
982 let limit = expected_group_size.unwrap_or(4_000_000_000);
983
984 // Distribute buckets in powers of 16, so that we can strike a balance between how many inputs
985 // each layer gets from the preceding layer, while also limiting the number of layers.
986 while current < limit {
987 buckets.push(current);
988 current = current.saturating_mul(16);
989 }
990
991 buckets.reverse();
992 buckets
993}