1use 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
41pub 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() }
60}
61
62impl OptimizerTrace {
63 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 .with(PlanTrace::<String>::new(filter()))
76 .with(PlanTrace::<HirScalarExpr>::new(filter()))
77 .with(PlanTrace::<MirScalarExpr>::new(filter()))
78 .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 .with(PlanTrace::<FastPathPlan>::new(None))
87 .with(PlanTrace::<UsedIndexes>::new(None))
88 .with(tracing::level_filters::LevelFilter::TRACE);
97
98 OptimizerTrace(dispatcher::Dispatch::new(subscriber))
99 } else {
100 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 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 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 let rows = collect_all(format)?
166 .0
167 .into_iter()
168 .map(|entry| {
169 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 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 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 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 AdapterError::Unstructured(anyhow::anyhow!(format!(
278 "cannot EXPLAIN {stage} FOR {stmt_kind}"
279 )))
280 } else {
281 AdapterError::Internal(format!(
283 "stage `{path}` not present in the collected optimizer trace",
284 ))
285 }
286 })?;
287 vec![row]
288 }
289 };
290
291 drop(self);
301 tracing_core::callsite::rebuild_interest_cache();
302 Ok(rows)
303 }
304
305 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 Ok(rows.into_element().into_element().unwrap_str().into())
335 }
336
337 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 let mut context = ExplainContext {
354 config,
355 features,
356 humanizer,
357 cardinality_stats: Default::default(), 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 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 let mut context = ExplainContext {
377 config,
378 features,
379 humanizer,
380 cardinality_stats: Default::default(), 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 results.extend(itertools::chain!(
406 self.collect_scalar_entries::<HirScalarExpr>(),
407 self.collect_scalar_entries::<MirScalarExpr>(),
408 self.collect_string_entries(),
409 ));
410
411 results.sort_by_key(|x| x.instant);
415
416 Ok(TraceEntries(results))
417 }
418
419 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 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 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 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 context.duration = entry.full_duration;
456
457 let used_indexes = used_indexes_trace.map(|t| used_indexes_for(t, &entry.path));
459
460 let plan = if let Some(mut used_indexes) = used_indexes {
462 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 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 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 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
522pub struct DispatchGuard<'a> {
524 _tracing_guard: tracing::subscriber::DefaultGuard,
525 _life: std::marker::PhantomData<&'a ()>,
526}
527
528pub struct TraceEntries<T>(pub Vec<TraceEntry<T>>);
530
531impl<T> TraceEntries<T> {
532 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
540fn used_indexes_for(trace: &PlanTrace<UsedIndexes>, plan_path: &str) -> UsedIndexes {
545 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 let entry = match path {
555 Some(path) => trace.find(path.path()),
556 None => None,
557 };
558 entry.map_or_else(Default::default, |e| e.plan)
561}