risingwave_frontend/optimizer/plan_node/
logical_cdc_scan.rs1use 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#[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, 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 pub fn column_descs(&self) -> Vec<ColumnDesc> {
86 self.core.column_descs()
87 }
88
89 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 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}