Skip to main content

mz_adapter/explain/
optimizer_trace.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//! Tracing utilities for explainable plans.
11
12use std::fmt::{Debug, Display};
13use std::sync::Arc;
14
15use mz_catalog::memory::objects::Cluster;
16use mz_compute_types::dataflows::DataflowDescription;
17use mz_compute_types::plan::LirRelationExpr;
18use mz_expr::explain::ExplainContext;
19use mz_expr::{MirRelationExpr, MirScalarExpr, OptimizedMirRelationExpr, RowSetFinishing};
20use mz_ore::collections::CollectionExt;
21use mz_repr::explain::tracing::{PlanTrace, TraceEntry};
22use mz_repr::explain::{
23    Explain, ExplainConfig, ExplainError, ExplainFormat, ExprHumanizer, UsedIndexes,
24};
25use mz_repr::optimize::OptimizerFeatures;
26use mz_repr::{Datum, Row};
27use mz_sql::ast::display::AstDisplay;
28use mz_sql::plan::{self, HirRelationExpr, HirScalarExpr};
29use mz_sql_parser::ast::{ExplainStage, NamedPlan};
30use mz_transform::dataflow::DataflowMetainfo;
31use mz_transform::notice::RawOptimizerNotice;
32use smallvec::SmallVec;
33use tracing::dispatcher;
34use tracing_subscriber::prelude::*;
35
36use crate::AdapterError;
37use crate::coord::peek::FastPathPlan;
38use crate::explain::Explainable;
39use crate::explain::insights::{self, PlanInsightsContext};
40
41/// Provides functionality for tracing plans generated by the execution of an
42/// optimization pipeline.
43///
44/// Internally, this will create a layered [`tracing::subscriber::Subscriber`]
45/// consisting of one layer for each supported plan type `T` and wrap it into a
46/// [`dispatcher::Dispatch`] instance.
47///
48/// Use [`OptimizerTrace::as_guard`] to activate the [`dispatcher::Dispatch`]
49/// and collect a trace.
50///
51/// Use [`OptimizerTrace::into_rows`] or [`OptimizerTrace::into_plan_insights`]
52/// to cleanly destroy the [`OptimizerTrace`] instance and obtain the tracing
53/// result.
54pub struct OptimizerTrace(dispatcher::Dispatch);
55
56impl std::fmt::Debug for OptimizerTrace {
57    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
58        f.debug_tuple("OptimizerTrace").finish() // Skip the dispatch field
59    }
60}
61
62impl OptimizerTrace {
63    /// Create a new [`OptimizerTrace`].
64    ///
65    /// The instance will only accumulate [`TraceEntry`] instances along
66    /// the prefix of the given `path` if `path` is present, or it will
67    /// accumulate all [`TraceEntry`] instances otherwise.
68    pub fn new(filter: Option<SmallVec<[NamedPlan; 4]>>) -> OptimizerTrace {
69        let filter = filter.map(|named| named.iter().map(NamedPlan::path).collect());
70        let filter = || filter.clone();
71        if let Some(global_subscriber) = mz_ore::tracing::GLOBAL_SUBSCRIBER.get() {
72            let subscriber = Arc::clone(global_subscriber)
73                // Collect `explain_plan` types that are not used in the regular explain
74                // path, but are useful when instrumenting code for debugging purposes.
75                .with(PlanTrace::<String>::new(filter()))
76                .with(PlanTrace::<HirScalarExpr>::new(filter()))
77                .with(PlanTrace::<MirScalarExpr>::new(filter()))
78                // Collect `explain_plan` types that are used in the regular explain path.
79                .with(PlanTrace::<HirRelationExpr>::new(filter()))
80                .with(PlanTrace::<MirRelationExpr>::new(filter()))
81                .with(PlanTrace::<DataflowDescription<OptimizedMirRelationExpr>>::new(filter()))
82                .with(PlanTrace::<DataflowDescription<LirRelationExpr>>::new(
83                    filter(),
84                ))
85                // Don't filter for FastPathPlan entries (there can be at most one).
86                .with(PlanTrace::<FastPathPlan>::new(None))
87                .with(PlanTrace::<UsedIndexes>::new(None))
88                // All optimizer spans are `TRACE` and up. Technically this slows down the system
89                // by skipping the tracing fast path DURING an `EXPLAIN`, but we haven't
90                // seen this be a problem (yet).
91                //
92                // Note that we typically do NOT use global filters like this, preferring
93                // per-layer ones, but we are forced to because per-layer filters
94                // require an `Arc<dyn Subscriber + LookupSpan>`, which isn't a trait
95                // exposed by tracing, for now.
96                .with(tracing::level_filters::LevelFilter::TRACE);
97
98            OptimizerTrace(dispatcher::Dispatch::new(subscriber))
99        } else {
100            // This codepath should not be taken except in tests, and is left here as a
101            // convenience.
102            let subscriber = tracing_subscriber::registry()
103                .with(PlanTrace::<String>::new(filter()))
104                .with(PlanTrace::<HirScalarExpr>::new(filter()))
105                .with(PlanTrace::<MirScalarExpr>::new(filter()))
106                .with(PlanTrace::<HirRelationExpr>::new(filter()))
107                .with(PlanTrace::<MirRelationExpr>::new(filter()))
108                .with(PlanTrace::<DataflowDescription<OptimizedMirRelationExpr>>::new(filter()))
109                .with(PlanTrace::<DataflowDescription<LirRelationExpr>>::new(
110                    filter(),
111                ))
112                .with(PlanTrace::<FastPathPlan>::new(None))
113                .with(PlanTrace::<UsedIndexes>::new(None))
114                .with(tracing::level_filters::LevelFilter::TRACE);
115
116            OptimizerTrace(dispatcher::Dispatch::new(subscriber))
117        }
118    }
119
120    /// Enter this [`OptimizerTrace`]'s tracing [`dispatcher::Dispatch`], returning a guard.
121    ///
122    /// Linked to this [`OptimizerTrace`] with a lifetime to ensure
123    /// [`OptimizerTrace::into_rows`] isn't called until the guard is dropped.
124    pub fn as_guard<'s>(&'s self) -> DispatchGuard<'s> {
125        let dispatch = self.0.clone();
126        let tracing_guard = tracing::dispatcher::set_default(&dispatch);
127
128        DispatchGuard {
129            _tracing_guard: tracing_guard,
130            _life: std::marker::PhantomData,
131        }
132    }
133
134    /// Convert the optimizer trace into a vector or rows that can be returned
135    /// to the client.
136    pub async fn into_rows(
137        self,
138        format: ExplainFormat,
139        config: &ExplainConfig,
140        features: &OptimizerFeatures,
141        humanizer: &dyn ExprHumanizer,
142        row_set_finishing: Option<RowSetFinishing>,
143        target_cluster: Option<&Cluster>,
144        dataflow_metainfo: DataflowMetainfo,
145        stage: ExplainStage,
146        stmt_kind: plan::ExplaineeStatementKind,
147        insights_ctx: Option<Box<PlanInsightsContext>>,
148    ) -> Result<Vec<Row>, AdapterError> {
149        let collect_all = |format| {
150            self.collect_all(
151                format,
152                config,
153                features,
154                humanizer,
155                row_set_finishing.clone(),
156                target_cluster.map(|c| c.name.as_str()),
157                dataflow_metainfo.clone(),
158            )
159        };
160
161        let rows = match stage {
162            ExplainStage::Trace => {
163                // For the `Trace` (pseudo-)stage, return the entire trace as
164                // triples of (time, path, plan) values.
165                let rows = collect_all(format)?
166                    .0
167                    .into_iter()
168                    .map(|entry| {
169                        // The trace would have to take over 584 years to overflow a u64.
170                        let span_duration = u64::try_from(entry.span_duration.as_nanos());
171                        Row::pack_slice(&[
172                            Datum::from(span_duration.unwrap_or(u64::MAX)),
173                            Datum::from(entry.path.as_str()),
174                            Datum::from(entry.plan.as_str()),
175                        ])
176                    })
177                    .collect();
178                rows
179            }
180            ExplainStage::PlanInsights => {
181                if format != ExplainFormat::Json {
182                    coord_bail!("EXPLAIN PLAN INSIGHTS only supports JSON format");
183                }
184
185                let mut text_traces = collect_all(ExplainFormat::Text)?;
186                let mut json_traces = collect_all(ExplainFormat::Json)?;
187                let global_plan = self.collect_global_plan();
188                let fast_path_plan = self.collect_fast_path_plan();
189
190                // Plans can be very large and exhaust the json serialization recursion limit.
191                // Convert those into error objects.
192                let mut get_plan = |name: NamedPlan| {
193                    let text_plan = match text_traces.remove(name.path()) {
194                        None => "<unknown>".into(),
195                        Some(entry) => entry.plan,
196                    };
197                    let json_plan = match json_traces.remove(name.path()) {
198                        None => serde_json::Value::Null,
199                        Some(entry) => serde_json::from_str(&entry.plan).unwrap_or_else(|e| {
200                            serde_json::json!({
201                                "error": format!("internal error: {e}"),
202                            })
203                        }),
204                    };
205                    serde_json::json!({
206                        "text": text_plan,
207                        "json": json_plan,
208                    })
209                };
210
211                let is_fast_path = fast_path_plan.is_some();
212                let mut plan_insights =
213                    insights::plan_insights(humanizer, global_plan, fast_path_plan);
214                let mut redacted_sql = None;
215                if let Some(insights_ctx) = insights_ctx {
216                    redacted_sql = insights_ctx
217                        .stmt
218                        .as_ref()
219                        .map(|s| Some(s.to_ast_string_redacted()));
220                    if let (Some(plan_insights), false) = (plan_insights.as_mut(), is_fast_path) {
221                        if insights_ctx.enable_re_optimize {
222                            plan_insights
223                                .compute_fast_path_clusters(humanizer, insights_ctx)
224                                .await;
225                        }
226                    }
227                }
228                let cluster = target_cluster.map(|c| {
229                    serde_json::json!({
230                        "name": c.name,
231                        "id": c.id,
232                    })
233                });
234
235                let output = serde_json::json!({
236                    "plans": {
237                        "raw": get_plan(NamedPlan::Raw),
238                        "optimized": {
239                            "global": get_plan(NamedPlan::Global),
240                            "fast_path": get_plan(NamedPlan::FastPath),
241                        }
242                    },
243                    "insights": plan_insights,
244                    "cluster": cluster,
245                    "redacted_sql": redacted_sql,
246                });
247                let output = serde_json::to_string_pretty(&output).expect("JSON string");
248                vec![Row::pack_slice(&[Datum::from(output.as_str())])]
249            }
250            _ => {
251                // For everything else, return the plan for the stage identified
252                // by the corresponding path.
253
254                let path = stage
255                    .paths()
256                    .map(|path| path.into_element().path())
257                    .ok_or_else(|| {
258                        AdapterError::Internal("explain stage unexpectedly missing path".into())
259                    })?;
260                let mut traces = collect_all(format)?;
261
262                // For certain stages we want to return the resulting fast path
263                // plan instead of the selected stage if it is present.
264                let plan = if stage.show_fast_path() && !config.no_fast_path {
265                    traces
266                        .remove(NamedPlan::FastPath.path())
267                        .or_else(|| traces.remove(path))
268                } else {
269                    traces.remove(path)
270                };
271
272                let row = plan
273                    .map(|entry| Row::pack_slice(&[Datum::from(entry.plan.as_str())]))
274                    .ok_or_else(|| {
275                        if !stmt_kind.supports(&stage) {
276                            // Print a nicer error for unsupported stages.
277                            AdapterError::Unstructured(anyhow::anyhow!(format!(
278                                "cannot EXPLAIN {stage} FOR {stmt_kind}"
279                            )))
280                        } else {
281                            // We don't expect this stage to be missing.
282                            AdapterError::Internal(format!(
283                                "stage `{path}` not present in the collected optimizer trace",
284                            ))
285                        }
286                    })?;
287                vec![row]
288            }
289        };
290
291        // We assume that any `Dispatch` cloned from this `OptimizerTrace` has long been dropped
292        // (`as_guard` tries to ensure this.). We rebuild the tracing interest cache, as
293        // this `OptimizerTrace` is acting like a reload-layer, and tracing needs to
294        // recalculate what the max level is, using this often-unknown
295        // API. Note that the reference to the `Dispatch` in self MUST be dropped before
296        // re-calculating interest.
297        //
298        // Before this is dropped and rebuilt, there is small extra cost to all `DEBUG` spans and
299        // events, if the other layers (otel and stderr) are only interested in `INFO`.
300        drop(self);
301        tracing_core::callsite::rebuild_interest_cache();
302        Ok(rows)
303    }
304
305    /// Collect a [`insights::PlanInsights`] with insights about the the
306    /// optimized plans rendered as a JSON `String`.
307    pub async fn into_plan_insights(
308        self,
309        features: &OptimizerFeatures,
310        humanizer: &dyn ExprHumanizer,
311        row_set_finishing: Option<RowSetFinishing>,
312        target_cluster: Option<&Cluster>,
313        dataflow_metainfo: DataflowMetainfo,
314        insights_ctx: Option<Box<PlanInsightsContext>>,
315    ) -> Result<String, AdapterError> {
316        let rows = self
317            .into_rows(
318                ExplainFormat::Json,
319                &ExplainConfig::default(),
320                features,
321                humanizer,
322                row_set_finishing,
323                target_cluster,
324                dataflow_metainfo,
325                ExplainStage::PlanInsights,
326                plan::ExplaineeStatementKind::Select,
327                insights_ctx,
328            )
329            .await?;
330
331        // When using `ExplainStage::PlanInsights`, we're guaranteed that the
332        // output is a single row containing a single column containing the plan
333        // insights as a string.
334        Ok(rows.into_element().into_element().unwrap_str().into())
335    }
336
337    /// Collect all traced plans for all plan types `T` that are available in
338    /// the wrapped [`dispatcher::Dispatch`].
339    fn collect_all(
340        &self,
341        format: ExplainFormat,
342        config: &ExplainConfig,
343        features: &OptimizerFeatures,
344        humanizer: &dyn ExprHumanizer,
345        row_set_finishing: Option<RowSetFinishing>,
346        target_cluster: Option<&str>,
347        dataflow_metainfo: DataflowMetainfo,
348    ) -> Result<TraceEntries<String>, ExplainError> {
349        let mut results = vec![];
350
351        // First, create an ExplainContext without `used_indexes`. We'll use this to, e.g., collect
352        // HIR plans.
353        let mut context = ExplainContext {
354            config,
355            features,
356            humanizer,
357            cardinality_stats: Default::default(), // empty stats
358            used_indexes: Default::default(),
359            finishing: row_set_finishing.clone(),
360            duration: Default::default(),
361            target_cluster,
362            optimizer_notices: RawOptimizerNotice::explain(
363                &dataflow_metainfo.optimizer_notices,
364                humanizer,
365                config.redacted,
366            )?,
367        };
368
369        // Collect trace entries of types produced by local optimizer stages.
370        results.extend(itertools::chain!(
371            self.collect_explainable_entries::<HirRelationExpr>(&format, &mut context)?,
372            self.collect_explainable_entries::<MirRelationExpr>(&format, &mut context)?,
373        ));
374
375        // Collect trace entries of types produced by global optimizer stages.
376        let mut context = ExplainContext {
377            config,
378            features,
379            humanizer,
380            cardinality_stats: Default::default(), // empty stats
381            used_indexes: Default::default(),
382            finishing: row_set_finishing,
383            duration: Default::default(),
384            target_cluster,
385            optimizer_notices: RawOptimizerNotice::explain(
386                &dataflow_metainfo.optimizer_notices,
387                humanizer,
388                config.redacted,
389            )?,
390        };
391        results.extend(itertools::chain!(
392            self.collect_explainable_entries::<DataflowDescription<OptimizedMirRelationExpr>>(
393                &format,
394                &mut context,
395            )?,
396            self.collect_explainable_entries::<DataflowDescription<LirRelationExpr>>(
397                &format,
398                &mut context
399            )?,
400            self.collect_explainable_entries::<FastPathPlan>(&format, &mut context)?,
401        ));
402
403        // Collect trace entries of type String, HirScalarExpr, MirScalarExpr
404        // which are useful for ad-hoc debugging.
405        results.extend(itertools::chain!(
406            self.collect_scalar_entries::<HirScalarExpr>(),
407            self.collect_scalar_entries::<MirScalarExpr>(),
408            self.collect_string_entries(),
409        ));
410
411        // sort plans by instant (TODO: this can be implemented in a more
412        // efficient way, as we can assume that each of the runs that are used
413        // to `*.extend` the `results` vector is already sorted).
414        results.sort_by_key(|x| x.instant);
415
416        Ok(TraceEntries(results))
417    }
418
419    /// Collects the global optimized plan from the trace, if it exists.
420    fn collect_global_plan(&self) -> Option<DataflowDescription<OptimizedMirRelationExpr>> {
421        self.0
422            .downcast_ref::<PlanTrace<DataflowDescription<OptimizedMirRelationExpr>>>()
423            .and_then(|trace| trace.find(NamedPlan::Global.path()))
424            .map(|entry| entry.plan)
425    }
426
427    /// Collects the fast path plan from the trace, if it exists.
428    fn collect_fast_path_plan(&self) -> Option<FastPathPlan> {
429        self.0
430            .downcast_ref::<PlanTrace<FastPathPlan>>()
431            .and_then(|trace| trace.find(NamedPlan::FastPath.path()))
432            .map(|entry| entry.plan)
433    }
434
435    /// Collect all trace entries of a plan type `T` that implements
436    /// [`Explainable`].
437    fn collect_explainable_entries<T>(
438        &self,
439        format: &ExplainFormat,
440        context: &mut ExplainContext,
441    ) -> Result<Vec<TraceEntry<String>>, ExplainError>
442    where
443        T: Clone + Debug + 'static,
444        for<'a> Explainable<'a, T>: Explain<'a, Context = ExplainContext<'a>>,
445    {
446        if let Some(trace) = self.0.downcast_ref::<PlanTrace<T>>() {
447            // Get a handle of the associated `PlanTrace<UsedIndexes>`.
448            let used_indexes_trace = self.0.downcast_ref::<PlanTrace<UsedIndexes>>();
449
450            trace
451                .collect_as_vec()
452                .into_iter()
453                .map(|mut entry| {
454                    // Update the context with the current time.
455                    context.duration = entry.full_duration;
456
457                    // Try to find the UsedIndexes instance for this entry.
458                    let used_indexes = used_indexes_trace.map(|t| used_indexes_for(t, &entry.path));
459
460                    // Render the EXPLAIN output string for this entry.
461                    let plan = if let Some(mut used_indexes) = used_indexes {
462                        // Temporary swap the found UsedIndexes with the default
463                        // one in the ExplainContext while explaining the plan
464                        // for this entry.
465                        std::mem::swap(&mut context.used_indexes, &mut used_indexes);
466                        let plan = Explainable::new(&mut entry.plan).explain(format, context)?;
467                        std::mem::swap(&mut context.used_indexes, &mut used_indexes);
468                        plan
469                    } else {
470                        // No UsedIndexes instance for this entry found - use
471                        // the default UsedIndexes in the ExplainContext.
472                        Explainable::new(&mut entry.plan).explain(format, context)?
473                    };
474
475                    Ok(TraceEntry {
476                        instant: entry.instant,
477                        span_duration: entry.span_duration,
478                        full_duration: entry.full_duration,
479                        path: entry.path,
480                        plan,
481                    })
482                })
483                .collect()
484        } else {
485            unreachable!("collect_explainable_entries called with wrong plan type T");
486        }
487    }
488
489    /// Collect all trace entries of a plan type `T`.
490    fn collect_scalar_entries<T>(&self) -> Vec<TraceEntry<String>>
491    where
492        T: Clone + Debug + 'static,
493        T: Display,
494    {
495        if let Some(trace) = self.0.downcast_ref::<PlanTrace<T>>() {
496            trace
497                .collect_as_vec()
498                .into_iter()
499                .map(|entry| TraceEntry {
500                    instant: entry.instant,
501                    span_duration: entry.span_duration,
502                    full_duration: entry.full_duration,
503                    path: entry.path,
504                    plan: entry.plan.to_string(),
505                })
506                .collect()
507        } else {
508            vec![]
509        }
510    }
511
512    /// Collect all trace entries with plans of type [`String`].
513    fn collect_string_entries(&self) -> Vec<TraceEntry<String>> {
514        if let Some(trace) = self.0.downcast_ref::<PlanTrace<String>>() {
515            trace.collect_as_vec()
516        } else {
517            vec![]
518        }
519    }
520}
521
522/// A wrapper around a `tracing::subscriber::DefaultGuard`.
523pub struct DispatchGuard<'a> {
524    _tracing_guard: tracing::subscriber::DefaultGuard,
525    _life: std::marker::PhantomData<&'a ()>,
526}
527
528/// A collection of optimizer trace entries with convenient accessor methods.
529pub struct TraceEntries<T>(pub Vec<TraceEntry<T>>);
530
531impl<T> TraceEntries<T> {
532    // Removes the first (and by assumption the only) trace that matches the
533    // given path from the collected trace.
534    pub fn remove(&mut self, path: &'static str) -> Option<TraceEntry<T>> {
535        let index = self.0.iter().position(|entry| entry.path == path);
536        index.map(|index| self.0.remove(index))
537    }
538}
539
540/// Get the [`UsedIndexes`] corresponding to the given `plan_path`.
541///
542/// Note that the path under which a `UsedIndexes` entry is traced might
543/// differ from the path of the `plan_path` of the plan that needs it.
544fn used_indexes_for(trace: &PlanTrace<UsedIndexes>, plan_path: &str) -> UsedIndexes {
545    // Compute the path from which we are going to lookup the `UsedIndexes`
546    // instance from the requested path.
547    let path = match NamedPlan::of_path(plan_path) {
548        Some(NamedPlan::Global) => Some(NamedPlan::Global),
549        Some(NamedPlan::Physical) => Some(NamedPlan::Global),
550        Some(NamedPlan::FastPath) => Some(NamedPlan::FastPath),
551        _ => None,
552    };
553    // Find the `TraceEntry` wrapping the `UsedIndexes` instance.
554    let entry = match path {
555        Some(path) => trace.find(path.path()),
556        None => None,
557    };
558    // Either return the `UsedIndexes` wrapped by the found entry or a
559    // default `UsedIndexes` instance if such entry was not found.
560    entry.map_or_else(Default::default, |e| e.plan)
561}