1use std::collections::HashMap;
16use std::ops::Deref;
17use std::rc::Rc;
18
19use iceberg::spec::{Operation, TableMetadata};
20use itertools::Itertools;
21use risingwave_common::bail_not_implemented;
22use risingwave_common::catalog::{
23 CdcTableDesc, ColumnCatalog, Engine, Field, RISINGWAVE_ICEBERG_COMMIT_EPOCH,
24 RISINGWAVE_ICEBERG_ROW_ID, ROW_ID_COLUMN_NAME, Schema,
25};
26use risingwave_common::constants::log_store::{
27 EPOCH_COLUMN_NAME, INSERT_OP_CODE, ROW_OP_COLUMN_NAME, encode_epoch,
28};
29use risingwave_common::session_config::IcebergQueryStorageMode;
30use risingwave_common::types::{DataType, Interval, ScalarImpl};
31use risingwave_common::util::iter_util::ZipEqFast;
32use risingwave_connector::WithOptionsSecResolved;
33use risingwave_connector::source::ConnectorProperties;
34use risingwave_connector::source::cdc::CdcScanOptions;
35use risingwave_connector::source::iceberg::IcebergTimeTravelInfo;
36use risingwave_sqlparser::ast::AsOf;
37use thiserror_ext::AsReport;
38
39use crate::TableCatalog;
40use crate::binder::{
41 BoundBaseTable, BoundGapFill, BoundIcebergMetadataTable, BoundJoin, BoundMatchRecognize,
42 BoundShare, BoundShareInput, BoundSource, BoundSystemTable, BoundWatermark,
43 BoundWindowTableFunction, Relation, WindowTableFunctionKind,
44};
45use crate::catalog::source_catalog::SourceCatalog;
46use crate::error::{ErrorCode, Result};
47use crate::expr::{CastContext, Expr, ExprImpl, ExprType, FunctionCall, InputRef, Literal};
48use crate::handler::cdc::derive_with_options_for_cdc_table;
49use crate::optimizer::IcebergSnapshotInfo;
50use crate::optimizer::plan_node::generic::{self, GenericPlanRef, SourceNodeKind};
51use crate::optimizer::plan_node::utils::to_iceberg_time_travel_as_of;
52use crate::optimizer::plan_node::{
53 LogicalApply, LogicalCdcScan, LogicalGapFill, LogicalHopWindow, LogicalIcebergIntermediateScan,
54 LogicalIcebergMetadataScan, LogicalJoin, LogicalMatchRecognize, LogicalPlanRef as PlanRef,
55 LogicalProject, LogicalScan, LogicalShare, LogicalSource, LogicalSysScan, LogicalTableFunction,
56 LogicalUnion, LogicalValues,
57};
58use crate::optimizer::property::Cardinality;
59use crate::planner::{PlanFor, Planner};
60use crate::utils::{ColIndexMapping, Condition};
61
62const ERROR_WINDOW_SIZE_ARG: &str =
63 "The size arg of window table function should be an interval literal.";
64
65impl Planner {
66 pub fn plan_relation(&mut self, relation: Relation) -> Result<PlanRef> {
67 match relation {
68 Relation::BaseTable(t) => self.plan_base_table(&t),
69 Relation::SystemTable(st) => self.plan_sys_table(*st),
70 Relation::IcebergMetadataTable(table) => self.plan_iceberg_metadata_table(*table),
71 Relation::Subquery(q) => Ok(self.plan_query(q.query)?.into_unordered_subplan()),
73 Relation::Join(join) => self.plan_join(*join),
74 Relation::Apply(join) => self.plan_apply(*join),
75 Relation::WindowTableFunction(tf) => self.plan_window_table_function(*tf),
76 Relation::Source(s) => self.plan_source(*s),
77 Relation::TableFunction {
78 expr: tf,
79 with_ordinality,
80 } => self.plan_table_function(tf, with_ordinality),
81 Relation::Watermark(tf) => self.plan_watermark(*tf),
82 Relation::Share(share) => self.plan_share(*share),
83 Relation::GapFill(bound_gap_fill) => self.plan_gap_fill(*bound_gap_fill),
84 Relation::MatchRecognize(mr) => self.plan_match_recognize(*mr),
85 }
86 }
87
88 pub(crate) fn plan_sys_table(&mut self, sys_table: BoundSystemTable) -> Result<PlanRef> {
89 Ok(LogicalSysScan::create(
90 sys_table.sys_table_catalog,
91 self.ctx(),
92 Cardinality::unknown(), )
94 .into())
95 }
96
97 fn plan_iceberg_metadata_table(&mut self, table: BoundIcebergMetadataTable) -> Result<PlanRef> {
98 let timezone = self.ctx().get_session_timezone();
99 let time_travel_info = to_iceberg_time_travel_as_of(&table.as_of, &timezone)?;
100 let core = crate::optimizer::plan_node::generic::IcebergMetadataScan {
101 metadata_type: table.metadata_type,
102 properties: table.properties,
103 secret_refs: table.secret_refs,
104 time_travel_info,
105 ctx: self.ctx(),
106 };
107 Ok(LogicalIcebergMetadataScan::new(core).into())
108 }
109
110 pub(super) fn plan_base_table(&mut self, base_table: &BoundBaseTable) -> Result<PlanRef> {
111 let as_of = base_table.as_of.clone();
112 let scan = LogicalScan::from_base_table(base_table, self.ctx(), as_of.clone());
113
114 match base_table.table_catalog.engine {
115 Engine::Hummock => {
116 match as_of {
117 None
118 | Some(AsOf::ProcessTime)
119 | Some(AsOf::ProcessTimeBroadcast)
120 | Some(AsOf::TimestampNum(_))
121 | Some(AsOf::TimestampString(_))
122 | Some(AsOf::ProcessTimeWithInterval(_)) => {}
123 Some(AsOf::VersionNum(_)) | Some(AsOf::VersionString(_)) => {
124 bail_not_implemented!("As Of Version is not supported yet.")
125 }
126 };
127 Ok(scan.into())
128 }
129 Engine::Iceberg => self.plan_iceberg_table(base_table, scan, as_of),
130 }
131 }
132
133 fn plan_iceberg_table(
134 &mut self,
135 base_table: &BoundBaseTable,
136 scan: LogicalScan,
137 as_of: Option<AsOf>,
138 ) -> Result<PlanRef> {
139 let is_append_only = base_table.table_catalog.append_only;
140 let iceberg_query_storage_mode = self
141 .ctx()
142 .session_ctx()
143 .config()
144 .iceberg_query_storage_mode();
145
146 enum PlanTarget {
147 TableScan,
148 Source,
149 IntermediateScan,
150 }
151 let plan_target = match self.plan_for() {
152 PlanFor::StreamIcebergEngineInternal => PlanTarget::TableScan,
153 PlanFor::BatchDql => match iceberg_query_storage_mode {
154 IcebergQueryStorageMode::Hummock if !is_append_only => PlanTarget::TableScan,
157 _ => PlanTarget::IntermediateScan,
158 },
159 PlanFor::Stream => {
160 if is_append_only {
161 PlanTarget::Source
162 } else {
163 PlanTarget::TableScan
164 }
165 }
166 PlanFor::Batch => {
167 if is_append_only {
168 PlanTarget::IntermediateScan
169 } else {
170 PlanTarget::TableScan
171 }
172 }
173 };
174 match as_of {
175 None
176 | Some(AsOf::VersionNum(_))
177 | Some(AsOf::TimestampString(_))
178 | Some(AsOf::TimestampNum(_)) => {}
179 Some(AsOf::ProcessTime)
180 | Some(AsOf::ProcessTimeBroadcast)
181 | Some(AsOf::ProcessTimeWithInterval(_)) => {
182 bail_not_implemented!("As Of ProcessTime() is not supported yet.")
183 }
184 Some(AsOf::VersionString(_)) => {
185 bail_not_implemented!("As Of Version is not supported yet.")
186 }
187 }
188
189 if matches!(plan_target, PlanTarget::TableScan) {
190 return Ok(scan.into());
191 }
192
193 let source_catalog = self.get_iceberg_source_by_table_catalog(&base_table.table_catalog)
194 .ok_or_else(|| {
195 ErrorCode::BindError(format!(
196 "failed to plan an iceberg engine table: {}. Can't find the corresponding iceberg source. Maybe you need to recreate the table",
197 base_table.table_catalog.name()
198 ))
199 })?;
200
201 let mut table_column_type_mapping = HashMap::new();
206 let table_column_map: HashMap<&str, &DataType> = base_table
207 .table_catalog
208 .columns
209 .iter()
210 .map(|c| (c.name.as_str(), &c.column_desc.data_type))
211 .collect();
212 for source_col in &source_catalog.columns {
213 let source_name = source_col.name();
214 let table_name = if source_name == RISINGWAVE_ICEBERG_ROW_ID {
215 ROW_ID_COLUMN_NAME
216 } else {
217 source_name
218 };
219 if let Some(&table_type) = table_column_map.get(table_name)
220 && source_col.column_desc.data_type != *table_type
221 {
222 table_column_type_mapping.insert(source_name.to_owned(), table_type.clone());
223 }
224 }
225
226 let column_map: HashMap<String, (usize, ColumnCatalog)> = source_catalog
227 .columns
228 .clone()
229 .into_iter()
230 .enumerate()
231 .map(|(i, column)| (column.name().to_owned(), (i, column)))
232 .collect();
233 let exprs = scan
237 .table()
238 .column_schema()
239 .fields()
240 .iter()
241 .map(|field| {
242 let source_filed_name = if field.name == ROW_ID_COLUMN_NAME {
243 RISINGWAVE_ICEBERG_ROW_ID
244 } else {
245 &field.name
246 };
247 if let Some((i, source_column)) = column_map.get(source_filed_name) {
248 let input_type = &source_column.column_desc.data_type;
249 if matches!(plan_target, PlanTarget::Source) && input_type != &field.data_type {
250 let mut input_ref =
251 ExprImpl::InputRef(InputRef::new(*i, input_type.clone()).into());
252 FunctionCall::cast_mut(
253 &mut input_ref,
254 &field.data_type,
255 CastContext::Explicit,
256 )
257 .unwrap();
258 input_ref
259 } else {
260 ExprImpl::InputRef(InputRef::new(*i, field.data_type.clone()).into())
261 }
262 } else {
263 ExprImpl::Literal(Literal::new(None, field.data_type.clone()).into())
265 }
266 })
267 .collect_vec();
268
269 let table_col_index: HashMap<&str, usize> = base_table
272 .table_catalog
273 .columns
274 .iter()
275 .enumerate()
276 .map(|(i, c)| (c.name.as_str(), i))
277 .collect();
278 let source_to_table_mapping = ColIndexMapping::new(
279 source_catalog
280 .columns
281 .iter()
282 .map(|c| {
283 let table_name = if c.name() == RISINGWAVE_ICEBERG_ROW_ID {
284 ROW_ID_COLUMN_NAME
285 } else {
286 c.name()
287 };
288 table_col_index.get(table_name).copied()
289 })
290 .collect(),
291 base_table.table_catalog.columns.len(),
292 );
293
294 let logical_source = LogicalSource::with_catalog(
295 Rc::new(source_catalog),
296 SourceNodeKind::CreateMViewOrBatch,
297 self.ctx(),
298 as_of.clone(),
299 )?;
300 if matches!(plan_target, PlanTarget::Source) {
301 return Ok(LogicalProject::new(logical_source.into(), exprs).into());
302 }
303
304 let read_pending_log_store =
308 matches!(self.plan_for(), PlanFor::BatchDql) && is_append_only && as_of.is_none();
309 if read_pending_log_store {
310 drop(self.ctx().session_ctx().pinned_snapshot());
311 }
312
313 let mut logical_iceberg_intermediate_scan = self.plan_iceberg_intermediate_scan(
314 &logical_source,
315 table_column_type_mapping,
316 source_to_table_mapping,
317 )?;
318
319 if read_pending_log_store {
320 let snapshot_info = self.fetch_current_snapshot_info(&logical_source)?;
321 let committed_epoch = match snapshot_info {
325 None => Some(None),
326 Some(IcebergSnapshotInfo {
327 commit_epoch: Some(epoch),
328 ..
329 }) => Some(Some(epoch)),
330 Some(_) => None,
331 };
332 if let Some(committed_epoch) = committed_epoch
333 && let Some(log_store_scan) = self.plan_iceberg_log_store_scan(
334 &base_table.table_catalog,
335 logical_iceberg_intermediate_scan.schema(),
336 &logical_source.core.column_catalog,
337 committed_epoch,
338 )?
339 {
340 logical_iceberg_intermediate_scan = LogicalUnion::create(
341 true,
342 vec![logical_iceberg_intermediate_scan, log_store_scan],
343 );
344 }
345 }
346
347 Ok(LogicalProject::new(logical_iceberg_intermediate_scan, exprs).into())
348 }
349
350 pub(super) fn plan_source(&mut self, source: BoundSource) -> Result<PlanRef> {
351 if source.catalog.is_cdc_table_source() {
352 if !matches!(self.plan_for(), PlanFor::Stream) {
353 return Err(ErrorCode::NotSupported(
354 "batch queries on CDC table sources are not supported".to_owned(),
355 "Create a materialized view to consume the CDC table source".to_owned(),
356 )
357 .into());
358 }
359 if source.as_of.is_some() {
360 return Err(ErrorCode::NotSupported(
361 "AS OF on CDC table sources is not supported".to_owned(),
362 "Remove the AS OF clause".to_owned(),
363 )
364 .into());
365 }
366
367 let catalog_desc = source
368 .catalog
369 .info
370 .external_table
371 .as_ref()
372 .expect("checked by is_cdc_table_source");
373 let desc = CdcTableDesc::from_protobuf(catalog_desc)?;
374 let upstream_source = {
375 let session = self.ctx.session_ctx();
376 let catalog_reader = session.env().catalog_reader().read_guard();
377 catalog_reader
378 .get_source_by_id_with_db(&session.database(), desc.source_id)?
379 .clone()
380 };
381 let desc = resolve_current_cdc_table_desc(desc, &upstream_source.with_properties)?;
382 let scan = LogicalCdcScan::create(
383 source.catalog.name.clone(),
384 Rc::new(desc),
385 self.ctx(),
386 CdcScanOptions {
387 disable_backfill: true,
388 ..Default::default()
389 },
390 );
391 Ok(scan.into())
392 } else if source.is_shareable_cdc_connector() {
393 Err(ErrorCode::InternalError(
394 "Should not create MATERIALIZED VIEW or SELECT directly on shared CDC source. HINT: create TABLE from the source instead.".to_owned(),
395 )
396 .into())
397 } else {
398 let as_of = source.as_of.clone();
399 match as_of {
400 None
401 | Some(AsOf::VersionNum(_))
402 | Some(AsOf::TimestampString(_))
403 | Some(AsOf::TimestampNum(_)) => {}
404 Some(AsOf::ProcessTime)
405 | Some(AsOf::ProcessTimeBroadcast)
406 | Some(AsOf::ProcessTimeWithInterval(_)) => {
407 bail_not_implemented!("As Of ProcessTime() is not supported yet.")
408 }
409 Some(AsOf::VersionString(_)) => {
410 bail_not_implemented!("As Of Version is not supported yet.")
411 }
412 }
413 let is_iceberg = source.catalog.is_iceberg_connector();
414
415 if matches!(self.plan_for(), PlanFor::Stream) {
418 let has_pk =
419 source.catalog.row_id_index.is_some() || !source.catalog.pk_col_ids.is_empty();
420 if !has_pk {
421 debug_assert!(is_iceberg);
424 if is_iceberg {
425 return Err(ErrorCode::BindError(format!(
426 "Cannot create a stream job from an iceberg source without a primary key.\nThe iceberg source might be created in an older version of RisingWave. Please try recreating the source.\nSource: {:?}",
427 source.catalog
428 ))
429 .into());
430 } else {
431 return Err(ErrorCode::BindError(format!(
432 "Cannot create a stream job from a source without a primary key.
433This is a bug. We would appreciate a bug report at:
434https://github.com/risingwavelabs/risingwave/issues/new?labels=type%2Fbug&template=bug_report.yml
435
436source: {:?}",
437 source.catalog
438 ))
439 .into());
440 }
441 }
442 }
443
444 let source = LogicalSource::with_catalog(
445 Rc::new(source.catalog),
446 SourceNodeKind::CreateMViewOrBatch,
447 self.ctx(),
448 as_of,
449 )?;
450 if is_iceberg && !matches!(self.plan_for(), PlanFor::Stream) {
451 let num_cols = source.core.column_catalog.len();
452 let intermediate_scan = self.plan_iceberg_intermediate_scan(
453 &source,
454 HashMap::new(),
455 ColIndexMapping::identity(num_cols),
456 )?;
457 Ok(intermediate_scan)
458 } else {
459 Ok(source.into())
460 }
461 }
462 }
463
464 pub(super) fn plan_join(&mut self, join: BoundJoin) -> Result<PlanRef> {
465 let left = self.plan_relation(join.left)?;
466 let right = self.plan_relation(join.right)?;
467 let join_type = join.join_type;
468 let on_clause = join.cond;
469 if on_clause.has_subquery() {
470 bail_not_implemented!("Subquery in join on condition");
471 } else {
472 Ok(LogicalJoin::create(left, right, join_type, on_clause))
473 }
474 }
475
476 pub(super) fn plan_apply(&mut self, mut join: BoundJoin) -> Result<PlanRef> {
477 let join_type = join.join_type;
478 let on_clause = join.cond;
479 if on_clause.has_subquery() {
480 bail_not_implemented!("Subquery in join on condition");
481 }
482
483 let correlated_id = self.ctx.next_correlated_id();
484 let correlated_indices = join
485 .right
486 .collect_correlated_indices_by_depth_and_assign_id(0, correlated_id);
487 let left = self.plan_relation(join.left)?;
488 let right = self.plan_relation(join.right)?;
489
490 Ok(LogicalApply::create(
491 left,
492 right,
493 join_type,
494 Condition::with_expr(on_clause),
495 correlated_id,
496 correlated_indices,
497 false,
498 ))
499 }
500
501 pub(super) fn plan_window_table_function(
502 &mut self,
503 table_function: BoundWindowTableFunction,
504 ) -> Result<PlanRef> {
505 use WindowTableFunctionKind::*;
506 match table_function.kind {
507 Tumble => self.plan_tumble_window(
508 table_function.input,
509 table_function.time_col,
510 table_function.args,
511 ),
512 Hop => self.plan_hop_window(
513 table_function.input,
514 table_function.time_col,
515 table_function.args,
516 ),
517 }
518 }
519
520 pub(super) fn plan_table_function(
521 &mut self,
522 table_function: ExprImpl,
523 with_ordinality: bool,
524 ) -> Result<PlanRef> {
525 match table_function {
527 ExprImpl::TableFunction(tf) => {
528 Ok(LogicalTableFunction::new(*tf, with_ordinality, self.ctx()).into())
529 }
530 expr => {
531 let schema = Schema {
532 fields: vec![Field::unnamed(expr.return_type())],
534 };
535 let expr_return_type = expr.return_type();
536 let root = LogicalValues::create(vec![vec![expr]], schema, self.ctx());
537 let input_ref = ExprImpl::from(InputRef::new(0, expr_return_type.clone()));
538 let mut exprs = if let DataType::Struct(st) = expr_return_type {
539 st.iter()
540 .enumerate()
541 .map(|(i, (_, ty))| {
542 let idx = ExprImpl::literal_int(i.try_into().unwrap());
543 let args = vec![input_ref.clone(), idx];
544 FunctionCall::new_unchecked(ExprType::Field, args, ty.clone()).into()
545 })
546 .collect()
547 } else {
548 vec![input_ref]
549 };
550 if with_ordinality {
551 exprs.push(ExprImpl::literal_bigint(1));
552 }
553 Ok(LogicalProject::create(root, exprs))
554 }
555 }
556 }
557
558 pub(super) fn plan_share(&mut self, share: BoundShare) -> Result<PlanRef> {
559 match share.input {
560 BoundShareInput::Query(query) => {
561 let id = share.share_id;
562 match self.share_cache.get(&id) {
563 None => {
564 let result = self.plan_query(query)?.into_unordered_subplan();
565 let logical_share = LogicalShare::create(result);
566 self.share_cache.insert(id, logical_share.clone());
567 Ok(logical_share)
568 }
569 Some(result) => Ok(result.clone()),
570 }
571 }
572 BoundShareInput::ChangeLog {
573 relation,
574 key_indices,
575 } => {
576 let id = share.share_id;
577 let result = self.plan_changelog(relation, key_indices)?;
578 let logical_share = LogicalShare::create(result);
579 self.share_cache.insert(id, logical_share.clone());
580 Ok(logical_share)
581 }
582 }
583 }
584
585 pub(super) fn plan_watermark(&mut self, _watermark: BoundWatermark) -> Result<PlanRef> {
586 todo!("plan watermark");
587 }
588
589 pub(super) fn plan_gap_fill(&mut self, gap_fill: BoundGapFill) -> Result<PlanRef> {
590 let input = self.plan_relation(gap_fill.input)?;
591 Ok(LogicalGapFill::new(
592 input,
593 gap_fill.time_col,
594 gap_fill.interval,
595 gap_fill.fill_strategies,
596 gap_fill.partition_by_cols,
597 )
598 .into())
599 }
600
601 pub(super) fn plan_match_recognize(&mut self, mr: BoundMatchRecognize) -> Result<PlanRef> {
602 let input = self.plan_relation(mr.input)?;
603 Ok(LogicalMatchRecognize::new(
604 input,
605 mr.partition_by,
606 mr.order_by,
607 mr.measures,
608 mr.rows_per_match,
609 mr.after_match_skip,
610 mr.pattern,
611 mr.defines,
612 mr.within,
613 mr.within_deadline,
614 )
615 .into())
616 }
617
618 fn collect_col_data_types_for_tumble_window(relation: &Relation) -> Result<Vec<DataType>> {
619 let col_data_types = match relation {
620 Relation::Source(s) => s
621 .catalog
622 .columns
623 .iter()
624 .map(|col| col.data_type().clone())
625 .collect(),
626 Relation::BaseTable(t) => t
627 .table_catalog
628 .columns
629 .iter()
630 .map(|col| col.data_type().clone())
631 .collect(),
632 Relation::Subquery(q) => q.query.schema().data_types(),
633 Relation::Share(share) => share
634 .input
635 .fields()?
636 .into_iter()
637 .map(|(_, f)| f.data_type)
638 .collect(),
639 r => {
640 return Err(ErrorCode::BindError(format!(
641 "Invalid input relation to tumble: {r:?}"
642 ))
643 .into());
644 }
645 };
646 Ok(col_data_types)
647 }
648
649 fn plan_tumble_window(
650 &mut self,
651 input: Relation,
652 time_col: InputRef,
653 args: Vec<ExprImpl>,
654 ) -> Result<PlanRef> {
655 let mut args = args.into_iter();
656 let col_data_types: Vec<_> = Self::collect_col_data_types_for_tumble_window(&input)?;
657
658 match (args.next(), args.next(), args.next()) {
659 (Some(window_size @ ExprImpl::Literal(_)), None, None) => {
660 let mut exprs = Vec::with_capacity(col_data_types.len() + 2);
661 for (idx, col_dt) in col_data_types.iter().enumerate() {
662 exprs.push(InputRef::new(idx, col_dt.clone()).into());
663 }
664 let window_start: ExprImpl = FunctionCall::new(
665 ExprType::TumbleStart,
666 vec![ExprImpl::InputRef(Box::new(time_col)), window_size.clone()],
667 )?
668 .into();
669 let window_end =
673 FunctionCall::new(ExprType::Add, vec![window_start.clone(), window_size])?
674 .into();
675 exprs.push(window_start);
676 exprs.push(window_end);
677 let base = self.plan_relation(input)?;
678 let project = LogicalProject::create(base, exprs);
679 Ok(project)
680 }
681 (
682 Some(window_size @ ExprImpl::Literal(_)),
683 Some(window_offset @ ExprImpl::Literal(_)),
684 None,
685 ) => {
686 let mut exprs = Vec::with_capacity(col_data_types.len() + 2);
687 for (idx, col_dt) in col_data_types.iter().enumerate() {
688 exprs.push(InputRef::new(idx, col_dt.clone()).into());
689 }
690 let window_start: ExprImpl = FunctionCall::new(
691 ExprType::TumbleStart,
692 vec![
693 ExprImpl::InputRef(Box::new(time_col)),
694 window_size.clone(),
695 window_offset,
696 ],
697 )?
698 .into();
699 let window_end =
703 FunctionCall::new(ExprType::Add, vec![window_start.clone(), window_size])?
704 .into();
705 exprs.push(window_start);
706 exprs.push(window_end);
707 let base = self.plan_relation(input)?;
708 let project = LogicalProject::create(base, exprs);
709 Ok(project)
710 }
711 _ => Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into()),
712 }
713 }
714
715 fn plan_hop_window(
716 &mut self,
717 input: Relation,
718 time_col: InputRef,
719 args: Vec<ExprImpl>,
720 ) -> Result<PlanRef> {
721 let input = self.plan_relation(input)?;
722 let mut args = args.into_iter();
723 let Some((ExprImpl::Literal(window_slide), ExprImpl::Literal(window_size))) =
724 args.next_tuple()
725 else {
726 return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into());
727 };
728
729 let Some(ScalarImpl::Interval(window_slide)) = *window_slide.get_data() else {
730 return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into());
731 };
732 let Some(ScalarImpl::Interval(window_size)) = *window_size.get_data() else {
733 return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into());
734 };
735
736 let window_offset = match (args.next(), args.next()) {
737 (Some(ExprImpl::Literal(window_offset)), None) => match *window_offset.get_data() {
738 Some(ScalarImpl::Interval(window_offset)) => window_offset,
739 _ => return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into()),
740 },
741 (None, None) => Interval::from_month_day_usec(0, 0, 0),
742 _ => return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into()),
743 };
744
745 if !window_size.is_positive() || !window_slide.is_positive() {
746 return Err(ErrorCode::BindError(format!(
747 "window_size {} and window_slide {} must be positive",
748 window_size, window_slide
749 ))
750 .into());
751 }
752
753 if window_size.exact_div(&window_slide).is_none() {
754 return Err(ErrorCode::BindError(format!("Invalid arguments for HOP window function: window_size {} cannot be divided by window_slide {}",window_size, window_slide)).into());
755 }
756
757 Ok(LogicalHopWindow::create(
758 input,
759 time_col,
760 window_slide,
761 window_size,
762 window_offset,
763 ))
764 }
765
766 fn plan_iceberg_intermediate_scan(
767 &self,
768 source: &LogicalSource,
769 table_column_type_mapping: HashMap<String, DataType>,
770 source_to_table_mapping: ColIndexMapping,
771 ) -> Result<PlanRef> {
772 let timezone = self.ctx().get_session_timezone();
774 let mut time_travel_info = to_iceberg_time_travel_as_of(&source.core.as_of, &timezone)?;
775 if time_travel_info.is_none() {
776 time_travel_info = self
777 .fetch_current_snapshot_info(source)?
778 .map(|info| IcebergTimeTravelInfo::Version(info.snapshot_id));
779 }
780 let Some(time_travel_info) = time_travel_info else {
781 let mut schema = source.schema().clone();
782 for field in &mut schema.fields {
783 if let Some(target_type) = table_column_type_mapping.get(&field.name) {
784 field.data_type = target_type.clone();
785 }
786 }
787 return Ok(LogicalValues::new(vec![], schema, self.ctx()).into());
788 };
789 let intermediate_scan = LogicalIcebergIntermediateScan::new(
790 source,
791 time_travel_info,
792 table_column_type_mapping,
793 source_to_table_mapping,
794 );
795 Ok(intermediate_scan.into())
796 }
797
798 fn fetch_current_snapshot_info(
799 &self,
800 source: &LogicalSource,
801 ) -> Result<Option<IcebergSnapshotInfo>> {
802 let mut map = self.ctx.iceberg_snapshot_info_map();
803 let catalog = source.source_catalog().ok_or_else(|| {
804 crate::error::ErrorCode::InternalError(
805 "Iceberg source must have a valid source catalog".to_owned(),
806 )
807 })?;
808 if let Some(&snapshot_info) = map.get(&catalog.id) {
809 return Ok(snapshot_info);
810 }
811
812 #[cfg(madsim)]
813 return Err(crate::error::ErrorCode::BindError(
814 "iceberg source time travel can't be used in the madsim mode".to_string(),
815 )
816 .into());
817
818 #[cfg(not(madsim))]
819 {
820 let ConnectorProperties::Iceberg(prop) =
821 ConnectorProperties::extract(catalog.with_properties.clone(), false)?
822 else {
823 return Err(crate::error::ErrorCode::InternalError(
824 "Iceberg source must have Iceberg connector properties".to_owned(),
825 )
826 .into());
827 };
828
829 let snapshot_info = tokio::task::block_in_place(|| {
830 crate::utils::FRONTEND_RUNTIME.block_on(async {
831 prop.load_table().await.map(|table| {
832 let metadata = table.metadata();
833 metadata
834 .current_snapshot()
835 .map(|snapshot| IcebergSnapshotInfo {
836 snapshot_id: snapshot.snapshot_id(),
837 commit_epoch: risingwave_iceberg_commit_epoch(
838 metadata,
839 snapshot.snapshot_id(),
840 ),
841 })
842 })
843 })
844 })?;
845 map.insert(catalog.id, snapshot_info);
846 Ok(snapshot_info)
847 }
848 }
849
850 fn plan_iceberg_log_store_scan(
853 &self,
854 table_catalog: &TableCatalog,
855 target_schema: &Schema,
856 source_columns: &[ColumnCatalog],
857 committed_epoch: Option<u64>,
858 ) -> Result<Option<PlanRef>> {
859 let Some(sink_name) = table_catalog.iceberg_sink_name() else {
860 return Ok(None);
861 };
862 let log_store_table = {
863 let catalog_reader = self.ctx.session_ctx().env().catalog_reader().read_guard();
864 let Ok(schema) =
865 catalog_reader.get_schema_by_id(table_catalog.database_id, table_catalog.schema_id)
866 else {
867 return Ok(None);
868 };
869 let Some(sink) = schema.get_created_sink_by_name(&sink_name) else {
870 return Ok(None);
871 };
872 schema
873 .iter_internal_table()
874 .find(|table| {
875 table.job_id == Some(sink.id.as_job_id())
876 && table
877 .columns()
878 .iter()
879 .any(|column| column.name() == EPOCH_COLUMN_NAME)
880 && table
881 .columns()
882 .iter()
883 .any(|column| column.name() == ROW_OP_COLUMN_NAME)
884 })
885 .cloned()
886 };
887 let Some(log_store_table) = log_store_table else {
888 return Ok(None);
889 };
890
891 let Some(epoch_idx) = log_store_table
892 .columns()
893 .iter()
894 .position(|column| column.name() == EPOCH_COLUMN_NAME)
895 else {
896 return Ok(None);
897 };
898 let Some(row_op_idx) = log_store_table
899 .columns()
900 .iter()
901 .position(|column| column.name() == ROW_OP_COLUMN_NAME)
902 else {
903 return Ok(None);
904 };
905 let payload_col_idx = log_store_table
909 .columns()
910 .iter()
911 .enumerate()
912 .skip(row_op_idx + 1)
913 .filter_map(|(idx, column)| (!column.is_hidden).then_some(idx))
914 .collect_vec();
915 if source_columns.len() != target_schema.len()
916 || payload_col_idx.len()
917 != source_columns
918 .iter()
919 .filter(|column| !column.is_hidden)
920 .count()
921 {
922 return Ok(None);
923 }
924
925 let row_op_predicate: ExprImpl = FunctionCall::new(
926 ExprType::Equal,
927 vec![
928 InputRef::new(row_op_idx, DataType::Int16).into(),
929 Literal::new(Some(ScalarImpl::Int16(INSERT_OP_CODE)), DataType::Int16).into(),
930 ],
931 )?
932 .into();
933 let mut predicate = Condition::with_expr(row_op_predicate);
934 if let Some(epoch) = committed_epoch {
935 let epoch_predicate: ExprImpl = FunctionCall::new(
936 ExprType::GreaterThan,
937 vec![
938 InputRef::new(epoch_idx, DataType::Int64).into(),
939 Literal::new(
940 Some(ScalarImpl::Int64(encode_epoch(epoch))),
941 DataType::Int64,
942 )
943 .into(),
944 ],
945 )?
946 .into();
947 predicate = predicate.and(Condition::with_expr(epoch_predicate));
948 }
949
950 let log_store_scan: PlanRef = LogicalScan::from(generic::TableScan::new(
951 payload_col_idx,
952 log_store_table,
953 vec![],
954 vec![],
955 self.ctx(),
956 predicate,
957 None,
958 ))
959 .into();
960
961 let source_types = log_store_scan.schema().data_types();
962 let target_types = target_schema.data_types();
963 let mut payload_idx = 0;
964 let project_exprs = source_columns
965 .iter()
966 .zip_eq_fast(target_types)
967 .map(|(source_column, target_type)| {
968 if source_column.is_hidden {
969 return Ok(Literal::new(None, target_type).into());
970 }
971
972 let source_type = source_types[payload_idx].clone();
973 let mut expr: ExprImpl = InputRef::new(payload_idx, source_type).into();
974 payload_idx += 1;
975 FunctionCall::cast_mut(&mut expr, &target_type, CastContext::Explicit).map_err(
976 |error| {
977 ErrorCode::InternalError(format!(
978 "failed to align Iceberg log-store payload type: {}",
979 error.as_report()
980 ))
981 },
982 )?;
983 Ok(expr)
984 })
985 .collect::<Result<Vec<_>>>()?;
986 Ok(Some(LogicalProject::create(log_store_scan, project_exprs)))
987 }
988
989 fn get_iceberg_source_by_table_catalog(
990 &self,
991 table_catalog: &TableCatalog,
992 ) -> Option<SourceCatalog> {
993 let catalog_reader = self.ctx.session_ctx().env().catalog_reader().read_guard();
994
995 let iceberg_source_name = table_catalog.iceberg_source_name()?;
996 let schema = catalog_reader
997 .get_schema_by_id(table_catalog.database_id, table_catalog.schema_id)
998 .ok()?;
999 let source_catalog = schema.get_source_by_name(&iceberg_source_name)?;
1000 Some(source_catalog.deref().clone())
1001 }
1002}
1003
1004fn risingwave_iceberg_commit_epoch(metadata: &TableMetadata, snapshot_id: i64) -> Option<u64> {
1009 let mut snapshot = metadata.snapshot_by_id(snapshot_id)?;
1010 loop {
1011 if let Some(epoch) = snapshot
1012 .summary()
1013 .additional_properties
1014 .get(RISINGWAVE_ICEBERG_COMMIT_EPOCH)
1015 {
1016 return epoch.parse().ok();
1017 }
1018 if snapshot.summary().operation != Operation::Replace {
1019 return None;
1020 }
1021 snapshot = metadata.snapshot_by_id(snapshot.parent_snapshot_id()?)?;
1022 }
1023}
1024
1025fn resolve_current_cdc_table_desc(
1026 mut desc: CdcTableDesc,
1027 upstream_properties: &WithOptionsSecResolved,
1028) -> Result<CdcTableDesc> {
1029 let (current_properties, normalized_external_table_name) =
1030 derive_with_options_for_cdc_table(upstream_properties, desc.external_table_name)?;
1031 let (connect_properties, secret_refs) = current_properties.into_parts();
1032 desc.external_table_name = normalized_external_table_name;
1033 desc.connect_properties = connect_properties;
1034 desc.secret_refs = secret_refs;
1035 Ok(desc)
1036}
1037
1038#[cfg(test)]
1039mod tests {
1040 use std::collections::BTreeMap;
1041
1042 use iceberg::spec::{
1043 FormatVersion, MAIN_BRANCH, NestedField, PrimitiveType, Schema as IcebergSchema, Snapshot,
1044 SortOrder, Summary, TableMetadataBuilder, Type, UnboundPartitionSpec,
1045 };
1046
1047 use super::*;
1048
1049 #[test]
1050 fn test_risingwave_commit_epoch_inherits_only_across_replace_snapshots() {
1051 let append = snapshot(
1052 1,
1053 None,
1054 Operation::Append,
1055 HashMap::from([(RISINGWAVE_ICEBERG_COMMIT_EPOCH.to_owned(), "42".to_owned())]),
1056 );
1057 let replace = snapshot(2, Some(1), Operation::Replace, HashMap::new());
1058 let metadata = metadata_with_current(vec![append, replace]);
1059 assert_eq!(risingwave_iceberg_commit_epoch(&metadata, 2), Some(42));
1060
1061 let external_append = snapshot(3, Some(2), Operation::Append, HashMap::new());
1062 let metadata = metadata_with_current(vec![
1063 snapshot(
1064 1,
1065 None,
1066 Operation::Append,
1067 HashMap::from([(RISINGWAVE_ICEBERG_COMMIT_EPOCH.to_owned(), "42".to_owned())]),
1068 ),
1069 snapshot(2, Some(1), Operation::Replace, HashMap::new()),
1070 external_append,
1071 ]);
1072 assert_eq!(risingwave_iceberg_commit_epoch(&metadata, 3), None);
1073 }
1074
1075 fn metadata_with_current(mut snapshots: Vec<Snapshot>) -> TableMetadata {
1076 let current = snapshots.pop().unwrap();
1077 let mut builder = TableMetadataBuilder::new(
1078 IcebergSchema::builder()
1079 .with_fields(vec![
1080 NestedField::new(1, "id", Type::Primitive(PrimitiveType::Long), false).into(),
1081 ])
1082 .build()
1083 .unwrap(),
1084 UnboundPartitionSpec::builder().build(),
1085 SortOrder::unsorted_order(),
1086 "s3://warehouse/db/table".to_owned(),
1087 FormatVersion::V2,
1088 HashMap::new(),
1089 )
1090 .unwrap();
1091 for snapshot in snapshots {
1092 builder = builder.add_snapshot(snapshot).unwrap();
1093 }
1094 builder
1095 .set_branch_snapshot(current, MAIN_BRANCH)
1096 .unwrap()
1097 .build()
1098 .unwrap()
1099 .metadata
1100 }
1101
1102 fn snapshot(
1103 snapshot_id: i64,
1104 parent_snapshot_id: Option<i64>,
1105 operation: Operation,
1106 additional_properties: HashMap<String, String>,
1107 ) -> Snapshot {
1108 Snapshot::builder()
1109 .with_snapshot_id(snapshot_id)
1110 .with_parent_snapshot_id(parent_snapshot_id)
1111 .with_sequence_number(snapshot_id)
1112 .with_timestamp_ms(snapshot_id)
1113 .with_manifest_list(format!("/snap-{snapshot_id}.avro"))
1114 .with_summary(Summary {
1115 operation,
1116 additional_properties,
1117 })
1118 .with_schema_id(0)
1119 .build()
1120 }
1121
1122 #[test]
1123 fn test_resolve_current_cdc_table_desc_replaces_stale_properties() {
1124 let desc = CdcTableDesc {
1125 external_table_name: "risedev.orders".to_owned(),
1126 connect_properties: BTreeMap::from([("hostname".to_owned(), "old-host".to_owned())]),
1127 ..Default::default()
1128 };
1129 let upstream_properties = WithOptionsSecResolved::new(
1130 BTreeMap::from([
1131 ("connector".to_owned(), "mysql-cdc".to_owned()),
1132 ("database.name".to_owned(), "risedev".to_owned()),
1133 ("hostname".to_owned(), "new-host".to_owned()),
1134 ]),
1135 BTreeMap::new(),
1136 );
1137
1138 let desc = resolve_current_cdc_table_desc(desc, &upstream_properties).unwrap();
1139
1140 assert_eq!(desc.connect_properties.get("hostname").unwrap(), "new-host");
1141 assert_eq!(desc.connect_properties.get("table.name").unwrap(), "orders");
1142 assert_eq!(
1143 desc.connect_properties.get("database.name").unwrap(),
1144 "risedev"
1145 );
1146 }
1147}