1use std::marker::PhantomData;
13use std::sync::Arc;
14use std::time::{Duration, Instant};
15
16use differential_dataflow::lattice::Lattice;
17use mz_compute_types::ComputeInstanceId;
18use mz_compute_types::plan::LirRelationExpr;
19use mz_compute_types::sinks::{ComputeSinkConnection, ComputeSinkDesc, SubscribeSinkConnection};
20use mz_expr::{ColumnOrder, MirRelationExpr};
21use mz_ore::soft_assert_or_log;
22use mz_repr::{GlobalId, RelationDesc, Timestamp};
23use mz_sql::optimizer_metrics::OptimizerMetrics;
24use mz_sql::plan::{HirToMirConfig, SubscribeFrom, SubscribePlan};
25use mz_transform::TransformCtx;
26use mz_transform::dataflow::{DataflowMetainfo, optimize_dataflow_snapshot};
27use mz_transform::normalize_lets::normalize_lets;
28use mz_transform::typecheck::{SharedTypecheckingContext, empty_typechecking_context};
29use timely::progress::Antichain;
30
31use crate::CollectionIdBundle;
32use crate::optimize::dataflows::{
33 ComputeInstanceSnapshot, DataflowBuilder, ExprPrep, ExprPrepMaintained,
34 dataflow_import_id_bundle,
35};
36use crate::optimize::{
37 LirDataflowDescription, MirDataflowDescription, Optimize, OptimizeMode, OptimizerCatalog,
38 OptimizerConfig, OptimizerError, optimize_mir_local, trace_plan,
39};
40
41pub struct Optimizer {
42 typecheck_ctx: SharedTypecheckingContext,
44 catalog: Arc<dyn OptimizerCatalog>,
46 compute_instance: ComputeInstanceSnapshot,
48 sink_id: GlobalId,
50 view_id: GlobalId,
53 with_snapshot: bool,
55 up_to: Option<Timestamp>,
57 debug_name: String,
59 config: OptimizerConfig,
61 metrics: OptimizerMetrics,
63 duration: Duration,
65}
66
67impl std::fmt::Debug for Optimizer {
73 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
74 f.debug_struct("Optimizer")
75 .field("config", &self.config)
76 .finish_non_exhaustive()
77 }
78}
79
80impl Optimizer {
81 pub fn new(
82 catalog: Arc<dyn OptimizerCatalog>,
83 compute_instance: ComputeInstanceSnapshot,
84 view_id: GlobalId,
85 sink_id: GlobalId,
86 with_snapshot: bool,
87 up_to: Option<Timestamp>,
88 debug_name: String,
89 config: OptimizerConfig,
90 metrics: OptimizerMetrics,
91 ) -> Self {
92 Self {
93 typecheck_ctx: empty_typechecking_context(),
94 catalog,
95 compute_instance,
96 view_id,
97 sink_id,
98 with_snapshot,
99 up_to,
100 debug_name,
101 config,
102 metrics,
103 duration: Default::default(),
104 }
105 }
106
107 pub fn cluster_id(&self) -> ComputeInstanceId {
108 self.compute_instance.instance_id()
109 }
110
111 pub fn up_to(&self) -> Option<Timestamp> {
112 self.up_to.clone()
113 }
114
115 pub fn sink_id(&self) -> GlobalId {
116 self.sink_id
117 }
118
119 pub fn optimize_query(
127 &mut self,
128 expr: MirRelationExpr,
129 from_desc: RelationDesc,
130 output: Vec<ColumnOrder>,
131 ) -> Result<GlobalMirPlan<Unresolved>, OptimizerError> {
132 self.optimize_inner(SubscribeSource::Query { expr, from_desc }, output)
133 }
134
135 fn optimize_inner(
139 &mut self,
140 source: SubscribeSource,
141 output: Vec<ColumnOrder>,
142 ) -> Result<GlobalMirPlan<Unresolved>, OptimizerError> {
143 let time = Instant::now();
144
145 let mut df_builder = {
146 let compute = self.compute_instance.clone();
147 DataflowBuilder::new(&*self.catalog, compute).with_config(&self.config)
148 };
149 let mut df_desc = MirDataflowDescription::new(self.debug_name.clone());
150 let mut df_meta = DataflowMetainfo::default();
151
152 let (from, from_desc) = match source {
153 SubscribeSource::Id(from_id) => {
154 let from_desc = self
155 .catalog
156 .get_entry(&from_id)
157 .relation_desc()
158 .expect("subscribes can only be run on items with descs")
159 .into_owned();
160
161 df_builder.import_into_dataflow(&from_id, &mut df_desc, &self.config.features)?;
162
163 (from_id, from_desc)
164 }
165 SubscribeSource::Query { expr, from_desc } => {
166 let mut transform_ctx = TransformCtx::local(
168 &self.config.features,
169 &self.typecheck_ctx,
170 &mut df_meta,
171 Some(&mut self.metrics),
172 Some(self.view_id),
173 );
174 let expr = optimize_mir_local(expr, &mut transform_ctx)?;
175
176 df_builder.import_view_into_dataflow(
177 &self.view_id,
178 &expr,
179 &mut df_desc,
180 &self.config.features,
181 )?;
182
183 (self.view_id, from_desc)
184 }
185 };
186 df_builder.maybe_reoptimize_imported_views(&mut df_desc, &self.config)?;
187
188 let sink_description = ComputeSinkDesc {
190 from,
191 from_desc,
192 connection: ComputeSinkConnection::Subscribe(SubscribeSinkConnection { output }),
193 with_snapshot: self.with_snapshot,
194 up_to: self.up_to.map(Antichain::from_elem).unwrap_or_default(),
195 non_null_assertions: vec![],
197 refresh_schedule: None,
199 };
200 df_desc.export_sink(self.sink_id, sink_description);
201
202 let style = ExprPrepMaintained;
204 df_desc.visit_children(
205 |r| style.prep_relation_expr(r),
206 |s| style.prep_scalar_expr(s),
207 )?;
208
209 let mut transform_ctx = TransformCtx::global(
211 &df_builder,
212 &mz_transform::EmptyStatisticsOracle, &self.config.features,
214 &self.typecheck_ctx,
215 &mut df_meta,
216 Some(&mut self.metrics),
217 );
218 mz_transform::optimize_dataflow(&mut df_desc, &mut transform_ctx, false)?;
220
221 if self.config.mode == OptimizeMode::Explain {
222 trace_plan!(at: "global", &df_meta.used_indexes(&df_desc));
224 }
225
226 self.duration += time.elapsed();
227
228 Ok(GlobalMirPlan {
230 df_desc,
231 df_meta,
232 phantom: PhantomData::<Unresolved>,
233 })
234 }
235}
236
237enum SubscribeSource {
239 Id(GlobalId),
241 Query {
243 expr: MirRelationExpr,
244 from_desc: RelationDesc,
245 },
246}
247
248#[derive(Clone, Debug)]
254pub struct GlobalMirPlan<T: Clone> {
255 df_desc: MirDataflowDescription,
256 df_meta: DataflowMetainfo,
257 phantom: PhantomData<T>,
258}
259
260impl<T: Clone> GlobalMirPlan<T> {
261 pub fn id_bundle(&self, compute_instance_id: ComputeInstanceId) -> CollectionIdBundle {
263 dataflow_import_id_bundle(&self.df_desc, compute_instance_id)
264 }
265}
266
267#[derive(Clone, Debug)]
270pub struct GlobalLirPlan {
271 df_desc: LirDataflowDescription,
272 df_meta: DataflowMetainfo,
273}
274
275impl GlobalLirPlan {
276 pub fn sink_id(&self) -> GlobalId {
282 self.df_desc.sink_id()
283 }
284}
285
286#[derive(Clone, Debug)]
289pub struct Unresolved;
290
291#[derive(Clone, Debug)]
297pub struct Resolved;
298
299impl Optimize<SubscribePlan> for Optimizer {
300 type To = GlobalMirPlan<Unresolved>;
301
302 fn optimize(&mut self, plan: SubscribePlan) -> Result<Self::To, OptimizerError> {
303 let output = plan.output.row_order().to_vec();
304
305 match plan.from {
306 SubscribeFrom::Id(from_id) => self.optimize_inner(SubscribeSource::Id(from_id), output),
307 SubscribeFrom::Query { expr, desc } => {
308 let time = Instant::now();
317 let expr = expr.lower(HirToMirConfig::from(&self.config), None)?;
318 self.duration += time.elapsed();
319
320 self.optimize_inner(
321 SubscribeSource::Query {
322 expr,
323 from_desc: desc,
324 },
325 output,
326 )
327 }
328 }
329 }
330}
331
332impl GlobalMirPlan<Unresolved> {
333 pub fn resolve(mut self, as_of: Antichain<Timestamp>) -> GlobalMirPlan<Resolved> {
339 soft_assert_or_log!(
342 self.df_desc.index_exports.is_empty(),
343 "unexpectedly setting until for a DataflowDescription with an index",
344 );
345
346 self.df_desc.set_as_of(as_of);
348
349 self.df_desc.until = Antichain::from_elem(Timestamp::MIN);
353 for (_, sink) in &self.df_desc.sink_exports {
354 self.df_desc.until.join_assign(&sink.up_to);
355 }
356
357 GlobalMirPlan {
358 df_desc: self.df_desc,
359 df_meta: self.df_meta,
360 phantom: PhantomData::<Resolved>,
361 }
362 }
363}
364
365impl Optimize<GlobalMirPlan<Resolved>> for Optimizer {
366 type To = GlobalLirPlan;
367
368 fn optimize(&mut self, plan: GlobalMirPlan<Resolved>) -> Result<Self::To, OptimizerError> {
369 let time = Instant::now();
370
371 let GlobalMirPlan {
372 mut df_desc,
373 df_meta,
374 phantom: _,
375 } = plan;
376
377 for build in df_desc.objects_to_build.iter_mut() {
379 normalize_lets(&mut build.plan.0, &self.config.features)?
380 }
381
382 if self.config.subscribe_snapshot_optimization {
383 optimize_dataflow_snapshot(&mut df_desc)?;
385 }
386
387 let df_desc = LirRelationExpr::finalize_dataflow(
391 df_desc,
392 &self.config.features,
393 Some(self.metrics.lowering()),
394 )?;
395
396 self.duration += time.elapsed();
397 self.metrics
398 .observe_e2e_optimization_time("subscribe", self.duration);
399
400 Ok(GlobalLirPlan { df_desc, df_meta })
402 }
403}
404
405impl GlobalLirPlan {
406 pub fn unapply(self) -> (LirDataflowDescription, DataflowMetainfo) {
408 (self.df_desc, self.df_meta)
409 }
410}