Skip to main content

risingwave_frontend/optimizer/plan_node/
logical_cdc_scan.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::rc::Rc;
16
17use pretty_xmlish::{Pretty, XmlNode};
18use risingwave_common::catalog::{CdcTableDesc, ColumnDesc};
19use risingwave_connector::source::cdc::CdcScanOptions;
20
21use super::generic::GenericPlanRef;
22use super::utils::{Distill, childless_record};
23use super::{
24    BatchPlanRef, ColPrunable, ExprRewritable, Logical, LogicalFilter, LogicalPlanRef as PlanRef,
25    LogicalProject, PlanBase, PredicatePushdown, StreamPlanRef, ToBatch, ToStream, generic,
26};
27use crate::catalog::ColumnId;
28use crate::error::Result;
29use crate::expr::{ExprRewriter, ExprVisitor};
30use crate::optimizer::optimizer_context::OptimizerContextRef;
31use crate::optimizer::plan_node::expr_visitable::ExprVisitable;
32use crate::optimizer::plan_node::{
33    ColumnPruningContext, PredicatePushdownContext, RewriteStreamContext, StreamCdcTableScan,
34    ToStreamContext,
35};
36use crate::optimizer::property::Order;
37use crate::utils::{ColIndexMapping, Condition};
38
39/// `LogicalCdcScan` reads rows of a table from an external upstream database
40#[derive(Debug, Clone, PartialEq, Eq, Hash)]
41pub struct LogicalCdcScan {
42    pub base: PlanBase<Logical>,
43    core: generic::CdcScan,
44}
45
46impl From<generic::CdcScan> for LogicalCdcScan {
47    fn from(core: generic::CdcScan) -> Self {
48        let base = PlanBase::new_logical_with_core(&core);
49        Self { base, core }
50    }
51}
52
53impl From<generic::CdcScan> for PlanRef {
54    fn from(core: generic::CdcScan) -> Self {
55        LogicalCdcScan::from(core).into()
56    }
57}
58
59impl LogicalCdcScan {
60    pub fn create(
61        table_name: String, // explain-only
62        cdc_table_desc: Rc<CdcTableDesc>,
63        ctx: OptimizerContextRef,
64        options: CdcScanOptions,
65    ) -> Self {
66        generic::CdcScan::new(
67            table_name,
68            (0..cdc_table_desc.columns.len()).collect(),
69            cdc_table_desc,
70            ctx,
71            options,
72        )
73        .into()
74    }
75
76    pub fn table_name(&self) -> &str {
77        &self.core.table_name
78    }
79
80    pub fn cdc_table_desc(&self) -> &CdcTableDesc {
81        self.core.cdc_table_desc.as_ref()
82    }
83
84    /// Get the descs of the output columns.
85    pub fn column_descs(&self) -> Vec<ColumnDesc> {
86        self.core.column_descs()
87    }
88
89    /// Get the ids of the output columns.
90    pub fn output_column_ids(&self) -> Vec<ColumnId> {
91        self.core.output_column_ids()
92    }
93
94    pub fn output_col_idx(&self) -> &Vec<usize> {
95        &self.core.output_col_idx
96    }
97}
98
99impl_plan_tree_node_for_leaf! { Logical, LogicalCdcScan}
100
101impl Distill for LogicalCdcScan {
102    fn distill<'a>(&self) -> XmlNode<'a> {
103        let verbose = self.base.ctx().is_explain_verbose();
104        let mut vec = Vec::with_capacity(5);
105        vec.push(("table", Pretty::from(self.table_name().to_owned())));
106        let key_is_columns = true;
107        let key = if key_is_columns {
108            "columns"
109        } else {
110            "output_columns"
111        };
112        vec.push((key, self.core.columns_pretty(verbose)));
113        if !key_is_columns {
114            vec.push((
115                "required_columns",
116                Pretty::Array(
117                    self.output_col_idx()
118                        .iter()
119                        .map(|i| {
120                            let col_name = &self.cdc_table_desc().columns[*i].name;
121                            Pretty::from(if verbose {
122                                format!("{}.{}", self.table_name(), col_name)
123                            } else {
124                                col_name.clone()
125                            })
126                        })
127                        .collect(),
128                ),
129            ));
130        }
131
132        childless_record("LogicalCdcScan", vec)
133    }
134}
135
136impl ColPrunable for LogicalCdcScan {
137    fn prune_col(&self, required_cols: &[usize], _ctx: &mut ColumnPruningContext) -> PlanRef {
138        // The CDC executor currently expects an identity projection over the complete external
139        // table descriptor. Keep the scan intact and let a project above it prune the user-visible
140        // columns. During stream rewrite, the project also keeps the stream key internally.
141        LogicalProject::with_out_col_idx(self.clone().into(), required_cols.iter().copied()).into()
142    }
143}
144
145impl ExprRewritable<Logical> for LogicalCdcScan {
146    fn has_rewritable_expr(&self) -> bool {
147        true
148    }
149
150    fn rewrite_exprs(&self, r: &mut dyn ExprRewriter) -> PlanRef {
151        let core = self.core.clone();
152        core.rewrite_exprs(r);
153        Self {
154            base: self.base.clone_with_new_plan_id(),
155            core,
156        }
157        .into()
158    }
159}
160
161impl ExprVisitable for LogicalCdcScan {
162    fn visit_exprs(&self, v: &mut dyn ExprVisitor) {
163        self.core.visit_exprs(v);
164    }
165}
166
167impl PredicatePushdown for LogicalCdcScan {
168    fn predicate_pushdown(
169        &self,
170        predicate: Condition,
171        _ctx: &mut PredicatePushdownContext,
172    ) -> PlanRef {
173        LogicalFilter::create(self.clone().into(), predicate)
174    }
175}
176
177impl ToBatch for LogicalCdcScan {
178    fn to_batch(&self) -> Result<BatchPlanRef> {
179        unreachable!()
180    }
181
182    fn to_batch_with_order_required(&self, _required_order: &Order) -> Result<BatchPlanRef> {
183        unreachable!()
184    }
185}
186
187impl ToStream for LogicalCdcScan {
188    fn to_stream(&self, _ctx: &mut ToStreamContext) -> Result<StreamPlanRef> {
189        Ok(StreamCdcTableScan::new(self.core.clone()).into())
190    }
191
192    fn logical_rewrite_for_stream(
193        &self,
194        _ctx: &mut RewriteStreamContext,
195    ) -> Result<(PlanRef, ColIndexMapping)> {
196        Ok((
197            self.clone().into(),
198            ColIndexMapping::identity(self.schema().len()),
199        ))
200    }
201}
202
203#[cfg(test)]
204mod tests {
205    use std::rc::Rc;
206
207    use risingwave_common::catalog::{CdcTableDesc, ColumnDesc, ColumnId, TableId};
208    use risingwave_common::id::SourceId;
209    use risingwave_common::types::DataType;
210    use risingwave_common::util::sort_util::{ColumnOrder, OrderType};
211    use risingwave_connector::source::cdc::CdcScanOptions;
212
213    use super::*;
214    use crate::expr::InputRef;
215    use crate::optimizer::optimizer_context::OptimizerContext;
216    use crate::optimizer::plan_node::PlanTreeNodeUnary;
217
218    #[test]
219    fn test_predicate_pushdown_preserves_filter() {
220        let desc = CdcTableDesc {
221            table_id: TableId::new(1),
222            source_id: SourceId::new(2),
223            external_table_name: "mydb.orders".to_owned(),
224            pk: vec![ColumnOrder::new(0, OrderType::ascending())],
225            pk_comparisons: vec![risingwave_common::catalog::CdcKeyComparison::Native],
226            columns: vec![
227                ColumnDesc::named("id", ColumnId::new(1), DataType::Int32),
228                ColumnDesc::named("selected", ColumnId::new(2), DataType::Boolean),
229            ],
230            stream_key: vec![0],
231            ..Default::default()
232        };
233        let scan: PlanRef = LogicalCdcScan::create(
234            "orders".to_owned(),
235            Rc::new(desc),
236            OptimizerContext::mock(),
237            CdcScanOptions {
238                disable_backfill: true,
239                ..Default::default()
240            },
241        )
242        .into();
243        let predicate = Condition::with_expr(InputRef::new(1, DataType::Boolean).into());
244        let mut ctx = PredicatePushdownContext::new(scan.clone());
245
246        let result = scan.predicate_pushdown(predicate.clone(), &mut ctx);
247
248        let filter = result
249            .as_logical_filter()
250            .expect("the predicate must remain above the CDC scan");
251        assert_eq!(filter.predicate(), &predicate);
252        assert!(filter.input().as_logical_cdc_scan().is_some());
253    }
254}