1use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
16use std::sync::{Arc, LazyLock};
17
18use anyhow::Context;
19use either::Either;
20use iceberg::arrow::type_to_arrow_type;
21use iceberg::spec::Transform;
22use itertools::Itertools;
23use maplit::{convert_args, hashmap, hashset};
24use pgwire::pg_response::{PgResponse, StatementType};
25use risingwave_common::array::arrow::IcebergArrowConvert;
26use risingwave_common::array::arrow::arrow_schema_iceberg::DataType as ArrowDataType;
27use risingwave_common::bail;
28use risingwave_common::catalog::{
29 ColumnCatalog, ICEBERG_SINK_PREFIX, ObjectId, RISINGWAVE_ICEBERG_ROW_ID, ROW_ID_COLUMN_NAME,
30 Schema,
31};
32use risingwave_common::license::Feature;
33use risingwave_common::secret::LocalSecretManager;
34use risingwave_common::system_param::reader::SystemParamsRead;
35use risingwave_common::types::{DataType, Timestamptz};
36use risingwave_common::util::epoch::Epoch;
37use risingwave_connector::sink::catalog::{SinkCatalog, SinkFormatDesc};
38use risingwave_connector::sink::file_sink::s3::SnowflakeSink;
39use risingwave_connector::sink::iceberg::{ICEBERG_SINK, IcebergConfig};
40use risingwave_connector::sink::kafka::KAFKA_SINK;
41use risingwave_connector::sink::snowflake_redshift::redshift::RedshiftSink;
42use risingwave_connector::sink::snowflake_redshift::snowflake::SnowflakeV2Sink;
43use risingwave_connector::sink::{
44 CONNECTOR_TYPE_KEY, SINK_SNAPSHOT_OPTION, SINK_TYPE_OPTION, SINK_TYPE_UPSERT,
45 SINK_USER_FORCE_APPEND_ONLY_OPTION, SINK_USER_IGNORE_DELETE_OPTION, Sink, enforce_secret_sink,
46 sink_is_exactly_once,
47};
48use risingwave_connector::{
49 AUTO_SCHEMA_CHANGE_KEY, SINK_CREATE_TABLE_IF_NOT_EXISTS_KEY, SINK_INTERMEDIATE_TABLE_NAME,
50 SINK_TARGET_TABLE_NAME, WithPropertiesExt,
51};
52use risingwave_pb::catalog::connection_params::PbConnectionType;
53use risingwave_pb::telemetry::TelemetryDatabaseObject;
54use risingwave_sqlparser::ast::{
55 CreateSink, CreateSinkStatement, EmitMode, Encode, ExplainOptions, Format, FormatEncodeOptions,
56 ObjectName, Query, Statement,
57};
58use risingwave_sqlparser::parser::Parser;
59
60use super::RwPgResponse;
61use super::create_mv::get_column_names;
62use super::create_source::UPSTREAM_SOURCE_KEY;
63use super::util::gen_query_from_table_name;
64use crate::binder::{Binder, Relation};
65use crate::catalog::root_catalog::SchemaPath;
66use crate::catalog::table_catalog::TableType;
67use crate::error::{ErrorCode, Result, RwError};
68use crate::expr::{ExprImpl, InputRef, rewrite_now_to_proctime};
69use crate::handler::HandlerArgs;
70use crate::handler::alter_table_column::fetch_table_catalog_for_alter;
71use crate::handler::create_mv::{
72 extract_streaming_job_resource_options, parse_column_names, resolve_streaming_job_resource_type,
73};
74use crate::handler::util::{
75 LongRunningNotificationAction, check_connector_match_connection_type,
76 ensure_connection_type_allowed, ensure_local_fs_connector_allowed,
77 execute_with_long_running_notification, get_table_catalog_by_table_name,
78 reject_internal_table_dependencies,
79};
80use crate::optimizer::backfill_order_strategy::plan_backfill_order;
81use crate::optimizer::plan_node::{
82 IcebergPartitionInfo, LogicalSource, PartitionComputeInfo, StreamPlanRef as PlanRef,
83 StreamProject, ensure_sync_log_store_fragment_root, generic,
84};
85use crate::optimizer::{OptimizerContext, RelationCollectorVisitor};
86use crate::scheduler::streaming_manager::CreatingStreamingJobInfo;
87use crate::session::SessionImpl;
88use crate::session::current::notice_to_user;
89use crate::stream_fragmenter::{GraphJobType, build_graph_with_strategy};
90use crate::utils::{resolve_connection_ref_and_secret_ref, resolve_privatelink_in_with_option};
91use crate::{Explain, Planner, TableCatalog, WithOptions, WithOptionsSecResolved};
92
93static SINK_ALLOWED_CONNECTION_CONNECTOR: LazyLock<HashSet<PbConnectionType>> =
94 LazyLock::new(|| {
95 hashset! {
96 PbConnectionType::Unspecified,
97 PbConnectionType::Kafka,
98 PbConnectionType::Iceberg,
99 PbConnectionType::Elasticsearch,
100 }
101 });
102
103static SINK_ALLOWED_CONNECTION_SCHEMA_REGISTRY: LazyLock<HashSet<PbConnectionType>> =
104 LazyLock::new(|| {
105 hashset! {
106 PbConnectionType::Unspecified,
107 PbConnectionType::SchemaRegistry,
108 }
109 });
110
111const SINK_SINCE_TIMESTAMP_OPTION: &str = "since_timestamp";
112
113pub struct SinkPlanContext {
115 pub query: Box<Query>,
116 pub sink_plan: PlanRef,
117 pub sink_catalog: SinkCatalog,
118 pub target_table_catalog: Option<Arc<TableCatalog>>,
119 pub dependencies: HashSet<ObjectId>,
120 pub since_timestamp_epoch: Option<u64>,
121}
122
123fn maybe_fill_intermediate_table_name(
124 with_options: &mut WithOptionsSecResolved,
125 connector: &str,
126 sink_name: &str,
127) -> Result<()> {
128 let should_fill = with_options
129 .value_eq_ignore_case(SINK_CREATE_TABLE_IF_NOT_EXISTS_KEY, "true")
130 && with_options.value_eq_ignore_case(SINK_TYPE_OPTION, SINK_TYPE_UPSERT)
131 && matches!(
132 connector,
133 RedshiftSink::SINK_NAME | SnowflakeV2Sink::SINK_NAME
134 );
135 if !should_fill || with_options.contains_key(SINK_INTERMEDIATE_TABLE_NAME) {
136 return Ok(());
137 }
138
139 let table_name = with_options
140 .get(SINK_TARGET_TABLE_NAME)
141 .ok_or_else(|| ErrorCode::BindError("'table.name' option must be specified.".to_owned()))?;
142 let intermediate_table_name =
143 format!("rw_{}_{}_{}", sink_name, table_name, uuid::Uuid::new_v4());
144 with_options.insert(
145 SINK_INTERMEDIATE_TABLE_NAME.to_owned(),
146 intermediate_table_name,
147 );
148 Ok(())
149}
150
151pub async fn gen_sink_plan(
152 handler_args: HandlerArgs,
153 stmt: CreateSinkStatement,
154 explain_options: Option<ExplainOptions>,
155 is_iceberg_engine_internal: bool,
156) -> Result<SinkPlanContext> {
157 let session = handler_args.session.clone();
158 let session = session.as_ref();
159 let user_specified_columns = !stmt.columns.is_empty();
160 let db_name = &session.database();
161 let (sink_schema_name, sink_table_name) =
162 Binder::resolve_schema_qualified_name(db_name, &stmt.sink_name)?;
163
164 let mut with_options = handler_args.with_options.clone();
165 extract_streaming_job_resource_options(&mut with_options);
168
169 if session
170 .env()
171 .system_params_manager()
172 .get_params()
173 .load()
174 .enforce_secret()
175 && Feature::SecretManagement.check_available().is_ok()
176 {
177 enforce_secret_sink(&with_options)?;
178 }
179
180 resolve_privatelink_in_with_option(&mut with_options)?;
181 let (mut resolved_with_options, connection_type, connector_conn_ref) =
182 resolve_connection_ref_and_secret_ref(
183 with_options,
184 session,
185 Some(TelemetryDatabaseObject::Sink),
186 )?;
187
188 let since_timestamp_epoch = resolved_with_options
189 .remove(SINK_SINCE_TIMESTAMP_OPTION)
190 .map(|value| {
191 let timestamp = value.parse::<Timestamptz>().map_err(|err| {
192 ErrorCode::InvalidInputSyntax(format!(
193 "invalid value {value:?} of '{SINK_SINCE_TIMESTAMP_OPTION}' option: {err}; \
194 expected a timestamptz string with an explicit time zone, \
195 for example '2024-01-01 00:00:00Z'"
196 ))
197 })?;
198 let timestamp_millis = u64::try_from(timestamp.timestamp_millis()).unwrap_or(0);
199 Ok::<_, RwError>(Epoch::from_unix_millis_or_earliest(timestamp_millis).0)
200 })
201 .transpose()?;
202 if since_timestamp_epoch.is_some() {
203 Feature::SinkSinceTimestamp.check_available()?;
204 }
205
206 ensure_connection_type_allowed(connection_type, &SINK_ALLOWED_CONNECTION_CONNECTOR)?;
207
208 if !matches!(connection_type, PbConnectionType::Unspecified) {
210 let Some(connector) = resolved_with_options.get_connector() else {
211 return Err(RwError::from(ErrorCode::ProtocolError(format!(
212 "missing field '{}' in WITH clause",
213 CONNECTOR_TYPE_KEY
214 ))));
215 };
216 check_connector_match_connection_type(connector.as_str(), &connection_type)?;
217 }
218
219 let partition_info = get_partition_compute_info(&resolved_with_options).await?;
220
221 let context = if let Some(explain_options) = explain_options {
222 OptimizerContext::new(handler_args.clone(), explain_options)
223 } else {
224 OptimizerContext::from_handler_args(handler_args.clone())
225 };
226 let is_auto_schema_change = resolved_with_options
227 .get(AUTO_SCHEMA_CHANGE_KEY)
228 .map(|value| {
229 value.parse::<bool>().map_err(|_| {
230 ErrorCode::InvalidInputSyntax(format!(
231 "invalid value {} of '{}' option, expect",
232 value, AUTO_SCHEMA_CHANGE_KEY
233 ))
234 })
235 })
236 .transpose()?
237 .unwrap_or(false);
238
239 if is_auto_schema_change && !is_iceberg_engine_internal {
240 Feature::SinkAutoSchemaChange.check_available()?;
241 }
242
243 let sink_into_table_name = stmt.into_table_name.as_ref().map(|name| name.real_value());
244 if sink_into_table_name.is_some() {
245 let prev = resolved_with_options.insert(CONNECTOR_TYPE_KEY.to_owned(), "table".to_owned());
246
247 if prev.is_some() {
248 return Err(RwError::from(ErrorCode::BindError(
249 "In the case of sinking into table, the 'connector' parameter should not be provided.".to_owned(),
250 )));
251 }
252 }
253 let connector = resolved_with_options
254 .get(CONNECTOR_TYPE_KEY)
255 .cloned()
256 .ok_or_else(|| ErrorCode::BindError(format!("missing field '{CONNECTOR_TYPE_KEY}'")))?;
257 ensure_local_fs_connector_allowed(session, &connector)?;
258
259 let sink_from_table_name;
261 let direct_sink_from_name: Option<(ObjectName, bool)>;
264 let mut query = match stmt.sink_from {
265 CreateSink::From(from_name) => {
266 sink_from_table_name = from_name.0.last().unwrap().real_value();
267 direct_sink_from_name = Some((from_name.clone(), is_auto_schema_change));
268 if is_auto_schema_change && sink_into_table_name.is_some() {
269 return Err(RwError::from(ErrorCode::InvalidInputSyntax(
270 "auto schema change not supported for sink-into-table".to_owned(),
271 )));
272 }
273 maybe_fill_intermediate_table_name(
274 &mut resolved_with_options,
275 &connector,
276 &sink_table_name,
277 )?;
278 Box::new(gen_query_from_table_name(from_name))
279 }
280 CreateSink::AsQuery(query) => {
281 if is_auto_schema_change {
282 return Err(RwError::from(ErrorCode::InvalidInputSyntax(
283 "auto schema change not supported for CREATE SINK AS QUERY".to_owned(),
284 )));
285 }
286 sink_from_table_name = sink_table_name.clone();
287 direct_sink_from_name = None;
288 query
289 }
290 };
291
292 if is_iceberg_engine_internal && let Some((from_name, _)) = &direct_sink_from_name {
293 let (table, _) = get_table_catalog_by_table_name(session, from_name)?;
294 let pk_names = table.pk_column_names();
295 if pk_names.len() == 1 && pk_names[0].eq(ROW_ID_COLUMN_NAME) {
296 let [stmt]: [_; 1] = Parser::parse_sql(&format!(
297 "select {} as {}, * from {}",
298 ROW_ID_COLUMN_NAME, RISINGWAVE_ICEBERG_ROW_ID, from_name
299 ))
300 .context("unable to parse query")?
301 .try_into()
302 .unwrap();
303 let Statement::Query(parsed_query) = stmt else {
304 panic!("unexpected statement: {:?}", stmt);
305 };
306 query = parsed_query;
307 }
308 }
309
310 let (sink_database_id, sink_schema_id) =
311 session.get_database_and_schema_id_for_create(sink_schema_name.clone())?;
312
313 if since_timestamp_epoch.is_some() {
314 if sink_into_table_name.is_some() {
315 return Err(ErrorCode::BindError(format!(
316 "`{SINK_SINCE_TIMESTAMP_OPTION}` does not support `CREATE SINK INTO TABLE`"
317 ))
318 .into());
319 }
320 if is_iceberg_engine_internal {
321 return Err(ErrorCode::BindError(format!(
322 "`{SINK_SINCE_TIMESTAMP_OPTION}` does not support iceberg engine internal sinks"
323 ))
324 .into());
325 }
326 if let Some((from_name, _)) = &direct_sink_from_name {
327 let (table, _) = get_table_catalog_by_table_name(session, from_name)?;
328 if table.database_id != sink_database_id {
329 return Err(ErrorCode::NotSupported(
330 format!(
331 "`{SINK_SINCE_TIMESTAMP_OPTION}` does not support cross-database sinks"
332 ),
333 "Please create the sink in the same database as the upstream table.".to_owned(),
334 )
335 .into());
336 }
337 }
338 }
339
340 let (
341 dependent_relations,
342 dependent_udfs,
343 dependent_secrets,
344 bound,
345 auto_refresh_schema_from_table,
346 ) = {
347 let mut binder = Binder::new_for_stream(session);
348 let auto_refresh_schema_from_table = if let Some((from_name, true)) = &direct_sink_from_name
349 {
350 let from_relation = binder.bind_relation_by_name(from_name, None, None, true)?;
351 if let Relation::BaseTable(table) = from_relation {
352 if table.table_catalog.table_type != TableType::Table {
353 return Err(ErrorCode::InvalidInputSyntax(format!(
354 "auto schema change is supported only on TABLE, but got {:?}",
355 table.table_catalog.table_type
356 ))
357 .into());
358 }
359 if table.table_catalog.database_id != sink_database_id {
360 return Err(ErrorCode::InvalidInputSyntax(
361 "auto schema change sinks do not support cross-database tables".to_owned(),
362 )
363 .into());
364 }
365 for col in &table.table_catalog.columns {
366 if !col.is_hidden() && (col.is_generated() || col.is_rw_sys_column()) {
367 return Err(ErrorCode::InvalidInputSyntax(format!(
368 "auto schema change is not supported for tables with visible generated columns or visible system columns, but found column {}",
369 col.name()
370 ))
371 .into());
372 }
373 }
374 Some(table.table_catalog)
375 } else {
376 return Err(RwError::from(ErrorCode::NotSupported(
377 "auto schema change only supported for TABLE".to_owned(),
378 "try recreating the sink from table".to_owned(),
379 )));
380 }
381 } else {
382 None
383 };
384
385 let bound = binder.bind_query(&query)?;
386
387 (
388 binder.included_relations().clone(),
389 binder.included_udfs().clone(),
390 binder.included_secrets().clone(),
391 bound,
392 auto_refresh_schema_from_table,
393 )
394 };
395
396 reject_internal_table_dependencies(session, &dependent_relations, "CREATE SINK")?;
397
398 let col_names = if sink_into_table_name.is_some() {
399 parse_column_names(&stmt.columns)
400 } else {
401 get_column_names(&bound, stmt.columns)?
403 };
404
405 let emit_on_window_close = stmt.emit_mode == Some(EmitMode::OnWindowClose);
406 if emit_on_window_close {
407 context.warn_to_user("EMIT ON WINDOW CLOSE is currently an experimental feature. Please use it with caution.");
408 }
409
410 let format_desc = match stmt.sink_schema {
411 Some(f) => {
413 validate_compatibility(&connector, &f)?;
414 Some(bind_sink_format_desc(session,f)?)
415 }
416 None => match resolved_with_options.get(SINK_TYPE_OPTION) {
417 Some(t) => SinkFormatDesc::from_legacy_type(&connector, t)?.map(|mut f| {
419 session.notice_to_user("Consider using the newer syntax `FORMAT ... ENCODE ...` instead of `type = '...'`.");
420 if let Some(v) = resolved_with_options.get(SINK_USER_FORCE_APPEND_ONLY_OPTION) {
421 f.options.insert(SINK_USER_FORCE_APPEND_ONLY_OPTION.into(), v.into());
422 }
423 if let Some(v) = resolved_with_options.get(SINK_USER_IGNORE_DELETE_OPTION) {
424 f.options.insert(SINK_USER_IGNORE_DELETE_OPTION.into(), v.into());
425 }
426 f
427 }),
428 None => None,
430 },
431 };
432
433 let definition = context.normalized_sql().to_owned();
434 let mut plan_root = if is_iceberg_engine_internal {
435 Planner::new_for_iceberg_table_engine_sink(context.into()).plan_query(bound)?
436 } else {
437 Planner::new_for_stream(context.into()).plan_query(bound)?
438 };
439 if let Some(col_names) = &col_names {
440 plan_root.set_out_names(col_names.clone())?;
441 };
442
443 let without_snapshot = matches!(
444 resolved_with_options.remove(SINK_SNAPSHOT_OPTION),
445 Some(flag) if flag.eq_ignore_ascii_case("false")
446 );
447
448 if since_timestamp_epoch.is_some() && !without_snapshot {
449 return Err(ErrorCode::BindError(format!(
450 "`{SINK_SINCE_TIMESTAMP_OPTION}` requires `snapshot = false`"
451 ))
452 .into());
453 }
454
455 let target_table_catalog = stmt
456 .into_table_name
457 .as_ref()
458 .map(|table_name| fetch_table_catalog_for_alter(session, table_name).map(|t| t.0))
459 .transpose()?;
460
461 if let Some(target_table_catalog) = &target_table_catalog {
462 if let Some(col_names) = col_names {
463 let target_table_columns = target_table_catalog
464 .columns()
465 .iter()
466 .map(|c| c.name())
467 .collect::<BTreeSet<_>>();
468 for c in col_names {
469 if !target_table_columns.contains(c.as_str()) {
470 return Err(RwError::from(ErrorCode::BindError(format!(
471 "Column {} not found in table {}",
472 c,
473 target_table_catalog.name()
474 ))));
475 }
476 }
477 }
478 if target_table_catalog
479 .columns()
480 .iter()
481 .any(|col| !col.nullable())
482 {
483 notice_to_user(format!(
484 "The target table `{}` contains NOT NULL columns. Rows written by the sink that violate those constraints will be ignored silently.",
485 target_table_catalog.name(),
486 ));
487 }
488 }
489
490 let sink_plan = plan_root.gen_sink_plan(
491 sink_table_name,
492 definition,
493 resolved_with_options,
494 emit_on_window_close,
495 db_name.to_owned(),
496 sink_from_table_name,
497 format_desc,
498 without_snapshot,
499 since_timestamp_epoch.is_some(),
500 is_iceberg_engine_internal,
501 target_table_catalog.clone(),
502 partition_info,
503 user_specified_columns,
504 auto_refresh_schema_from_table,
505 )?;
506
507 let sink_desc = sink_plan.sink_desc().clone();
508
509 let mut sink_plan: PlanRef = sink_plan.into_stream_plan()?;
510 sink_plan = ensure_sync_log_store_fragment_root(sink_plan);
511
512 let ctx = sink_plan.ctx();
513 let explain_trace = ctx.is_explain_trace();
514 if explain_trace {
515 ctx.trace("Create Sink:");
516 ctx.trace(sink_plan.explain_to_string());
517 }
518 tracing::trace!("sink_plan: {:?}", sink_plan.explain_to_string());
519
520 let dependencies =
523 RelationCollectorVisitor::collect_with(dependent_relations, sink_plan.clone())
524 .into_iter()
525 .chain(dependent_udfs.iter().copied().map_into())
526 .chain(
527 dependent_secrets
528 .iter()
529 .copied()
530 .map(|id| id.as_object_id()),
531 )
532 .collect();
533
534 let sink_catalog = sink_desc.into_catalog(
535 sink_schema_id,
536 sink_database_id,
537 session.user_id(),
538 connector_conn_ref,
539 );
540
541 if let Some(table_catalog) = &target_table_catalog {
542 for column in sink_catalog.full_columns() {
543 if !column.can_dml() {
544 unreachable!(
545 "cannot derive generated columns or the `_rw_timestamp` system column in a sink catalog, but found one"
546 );
547 }
548 }
549
550 let table_columns_without_rw_timestamp = table_catalog.columns_without_rw_timestamp();
551 let exprs = derive_default_column_project_for_sink(
552 &sink_catalog,
553 sink_plan.schema(),
554 &table_columns_without_rw_timestamp,
555 user_specified_columns,
556 )?;
557
558 let logical_project = generic::Project::new(exprs, sink_plan);
559
560 sink_plan = StreamProject::new(logical_project).into();
561
562 let exprs = LogicalSource::derive_output_exprs_from_generated_columns(
563 &table_columns_without_rw_timestamp,
564 )?;
565
566 if let Some(exprs) = exprs {
567 let logical_project = generic::Project::new(exprs, sink_plan);
568 sink_plan = StreamProject::new(logical_project).into();
569 }
570 };
571
572 Ok(SinkPlanContext {
573 query,
574 sink_plan,
575 sink_catalog,
576 target_table_catalog,
577 dependencies,
578 since_timestamp_epoch,
579 })
580}
581
582pub async fn get_partition_compute_info(
587 with_options: &WithOptionsSecResolved,
588) -> Result<Option<PartitionComputeInfo>> {
589 let (options, secret_refs) = with_options.clone().into_parts();
590 let Some(connector) = options.get(UPSTREAM_SOURCE_KEY).cloned() else {
591 return Ok(None);
592 };
593 let properties = LocalSecretManager::global().fill_secrets(options, secret_refs)?;
594 match connector.as_str() {
595 ICEBERG_SINK => {
596 let iceberg_config = IcebergConfig::from_btreemap(properties)?;
597 get_partition_compute_info_for_iceberg(&iceberg_config).await
598 }
599 _ => Ok(None),
600 }
601}
602
603async fn get_partition_compute_info_for_iceberg(
604 _iceberg_config: &IcebergConfig,
605) -> Result<Option<PartitionComputeInfo>> {
606 if _iceberg_config.create_table_if_not_exists {
608 return Ok(None);
609 }
610 let table = _iceberg_config.load_table().await?;
611 let partition_spec = table.metadata().default_partition_spec();
612 if partition_spec.is_unpartitioned() {
613 return Ok(None);
614 }
615
616 let has_sparse_partition = partition_spec.fields().iter().any(|f| match f.transform {
621 Transform::Identity | Transform::Truncate(_) | Transform::Bucket(_) => true,
623 Transform::Year
625 | Transform::Month
626 | Transform::Day
627 | Transform::Hour
628 | Transform::Void
629 | Transform::Unknown => false,
630 });
631 if !has_sparse_partition {
632 return Ok(None);
633 }
634
635 let arrow_type = type_to_arrow_type(&iceberg::spec::Type::Struct(
636 table.metadata().default_partition_type().clone(),
637 ))
638 .map_err(|_| {
639 RwError::from(ErrorCode::SinkError(
640 "Failed to convert the Iceberg partition type to an Arrow type".into(),
641 ))
642 })?;
643 let ArrowDataType::Struct(struct_fields) = arrow_type else {
644 return Err(RwError::from(ErrorCode::SinkError(
645 "The Iceberg partition type must be a struct type".into(),
646 )));
647 };
648
649 let schema = table.metadata().current_schema();
650 let partition_fields = partition_spec
651 .fields()
652 .iter()
653 .map(|f| {
654 let source_f =
655 schema
656 .field_by_id(f.source_id)
657 .ok_or(RwError::from(ErrorCode::SinkError(
658 "Failed to look up the Iceberg partition field".into(),
659 )))?;
660 Ok((source_f.name.clone(), f.transform))
661 })
662 .collect::<Result<Vec<_>>>()?;
663
664 Ok(Some(PartitionComputeInfo::Iceberg(IcebergPartitionInfo {
665 partition_type: IcebergArrowConvert.struct_from_fields(&struct_fields)?,
666 partition_fields,
667 })))
668}
669
670pub async fn handle_create_sink(
671 mut handle_args: HandlerArgs,
672 stmt: CreateSinkStatement,
673 is_iceberg_engine_internal: bool,
674) -> Result<RwPgResponse> {
675 let session = handle_args.session.clone();
676
677 session.check_cluster_limits().await?;
678
679 let mode = if stmt.or_replace {
680 prepare_replace_sink(&mut handle_args, &stmt)?
681 } else {
682 let if_not_exists = stmt.if_not_exists;
683 if let Either::Right(resp) = session.check_relation_name_duplicated(
684 stmt.sink_name.clone(),
685 StatementType::CREATE_SINK,
686 if_not_exists,
687 )? {
688 return Ok(resp);
689 }
690
691 if stmt.sink_name.base_name().starts_with(ICEBERG_SINK_PREFIX) {
692 return Err(RwError::from(ErrorCode::InvalidInputSyntax(format!(
693 "Sink name cannot start with reserved prefix '{}'",
694 ICEBERG_SINK_PREFIX
695 ))));
696 }
697
698 SinkCreateMode::Create { if_not_exists }
699 };
700
701 create_sink_or_replace(handle_args, stmt, is_iceberg_engine_internal, mode).await
702}
703
704enum SinkCreateMode {
705 Create { if_not_exists: bool },
706 Replace { original_sink: Arc<SinkCatalog> },
707}
708
709impl SinkCreateMode {
710 fn statement_name(&self) -> &'static str {
711 match self {
712 SinkCreateMode::Create { .. } => "CREATE SINK",
713 SinkCreateMode::Replace { .. } => "REPLACE SINK",
714 }
715 }
716}
717
718async fn create_sink_or_replace(
719 mut handle_args: HandlerArgs,
720 stmt: CreateSinkStatement,
721 is_iceberg_engine_internal: bool,
722 mode: SinkCreateMode,
723) -> Result<RwPgResponse> {
724 let session = handle_args.session.clone();
725
726 let resource_type =
727 resolve_streaming_job_resource_type(session.as_ref(), &mut handle_args.with_options)?;
728
729 let (sink, graph, dependencies, since_timestamp_epoch) = {
730 let backfill_order_strategy = handle_args.with_options.backfill_order_strategy();
731 let SinkPlanContext {
732 query,
733 sink_plan: plan,
734 sink_catalog: mut sink,
735 target_table_catalog,
736 dependencies,
737 since_timestamp_epoch,
738 } = gen_sink_plan(handle_args, stmt, None, is_iceberg_engine_internal).await?;
739
740 let has_order_by = !query.order_by.is_empty();
741 if has_order_by {
742 plan.ctx().warn_to_user(
743 r#"The ORDER BY clause in the CREATE SINK statement has no effect at all."#
744 .to_owned(),
745 );
746 }
747
748 match &mode {
749 SinkCreateMode::Create { .. } => {
750 if let Some(table_catalog) = &target_table_catalog {
751 sink.original_target_columns = table_catalog.columns_without_rw_timestamp();
752 }
753 }
754 SinkCreateMode::Replace { original_sink } => {
755 if target_table_catalog.is_some() {
756 return Err(ErrorCode::NotSupported(
757 "REPLACE SINK INTO TABLE is not supported yet".to_owned(),
758 "replace ordinary sinks first".to_owned(),
759 )
760 .into());
761 }
762
763 sink.schema_id = original_sink.schema_id;
764 sink.database_id = original_sink.database_id;
765 sink.name = original_sink.name.clone();
766 sink.owner = original_sink.owner;
767 }
768 }
769
770 let backfill_order =
771 plan_backfill_order(session.as_ref(), backfill_order_strategy, plan.clone())?;
772 let graph =
773 build_graph_with_strategy(plan, Some(GraphJobType::Sink), Some(backfill_order))?;
774
775 (sink, graph, dependencies, since_timestamp_epoch)
776 };
777
778 let statement_name = mode.statement_name();
779 let catalog_writer = session.catalog_writer()?;
780 match mode {
781 SinkCreateMode::Create { if_not_exists } => {
782 let _job_guard = session.env().creating_streaming_job_tracker().guard(
783 CreatingStreamingJobInfo::new(
784 session.session_id(),
785 sink.database_id,
786 sink.schema_id,
787 sink.name.clone(),
788 ),
789 );
790
791 execute_with_long_running_notification(
792 catalog_writer.create_sink(
793 sink.to_proto(),
794 graph,
795 dependencies,
796 resource_type,
797 if_not_exists,
798 since_timestamp_epoch,
799 ),
800 &session,
801 statement_name,
802 LongRunningNotificationAction::MonitorBackfillJob,
803 )
804 .await?;
805 }
806 SinkCreateMode::Replace { original_sink } => {
807 let original_sink_id = original_sink.id;
808 execute_with_long_running_notification(
809 catalog_writer.replace_sink(
810 original_sink_id,
811 sink.to_proto(),
812 graph,
813 dependencies,
814 resource_type,
815 ),
816 &session,
817 statement_name,
818 LongRunningNotificationAction::DiagnoseBarrierLatency,
819 )
820 .await?;
821
822 tracing::info!(
823 old_sink_id = %original_sink_id,
824 sink_name = %sink.name,
825 "replace sink plan submitted"
826 );
827 }
828 }
829
830 Ok(PgResponse::empty_result(StatementType::CREATE_SINK))
831}
832
833fn prepare_replace_sink(
834 handle_args: &mut HandlerArgs,
835 stmt: &CreateSinkStatement,
836) -> Result<SinkCreateMode> {
837 let session = handle_args.session.clone();
838 if stmt.if_not_exists {
839 return Err(ErrorCode::InvalidInputSyntax(
840 "REPLACE SINK does not support IF NOT EXISTS".to_owned(),
841 )
842 .into());
843 }
844 if !matches!(&stmt.sink_from, CreateSink::From(_)) {
845 return Err(ErrorCode::NotSupported(
846 "REPLACE SINK currently only supports REPLACE SINK ... FROM table_or_mv".to_owned(),
847 "use REPLACE SINK name FROM existing_relation ...".to_owned(),
848 )
849 .into());
850 }
851 if stmt.into_table_name.is_some() {
852 return Err(ErrorCode::NotSupported(
853 "REPLACE SINK INTO TABLE is not supported yet".to_owned(),
854 "replace ordinary sinks first".to_owned(),
855 )
856 .into());
857 }
858 if handle_args
859 .with_options
860 .get(AUTO_SCHEMA_CHANGE_KEY)
861 .is_some_and(|value| value.eq_ignore_ascii_case("true"))
862 {
863 return Err(ErrorCode::NotSupported(
864 "REPLACE SINK with auto schema change is not supported yet".to_owned(),
865 "disable auto schema change for this replacement".to_owned(),
866 )
867 .into());
868 }
869 if handle_args
870 .with_options
871 .contains_key(SINK_SINCE_TIMESTAMP_OPTION)
872 {
873 return Err(ErrorCode::NotSupported(
874 "REPLACE SINK with since_timestamp is not supported yet".to_owned(),
875 "create a new sink with since_timestamp instead".to_owned(),
876 )
877 .into());
878 }
879 match handle_args.with_options.get(SINK_SNAPSHOT_OPTION) {
880 Some(value) if !value.eq_ignore_ascii_case("false") => {
881 return Err(ErrorCode::InvalidInputSyntax(
882 "REPLACE SINK must not enable snapshot backfill".to_owned(),
883 )
884 .into());
885 }
886 Some(_) => {}
887 None => {
888 handle_args
889 .with_options
890 .insert(SINK_SNAPSHOT_OPTION.to_owned(), "false".to_owned());
891 }
892 }
893
894 let db_name = session.database();
895 let (sink_schema_name, sink_table_name) =
896 Binder::resolve_schema_qualified_name(&db_name, &stmt.sink_name)?;
897 let original_sink = {
898 let search_path = session.config().search_path();
899 let user_name = session.user_name();
900 let schema_path = SchemaPath::new(sink_schema_name.as_deref(), &search_path, &user_name);
901 let reader = session.env().catalog_reader().read_guard();
902 let (sink, schema_name) =
903 reader.get_created_sink_by_name(&db_name, schema_path, &sink_table_name)?;
904 session.check_privilege_for_drop_alter(schema_name, &**sink)?;
905 if sink.target_table.is_some() {
906 return Err(ErrorCode::NotSupported(
907 "REPLACE SINK INTO TABLE is not supported yet".to_owned(),
908 "replace ordinary sinks first".to_owned(),
909 )
910 .into());
911 }
912 if sink.auto_refresh_schema_from_table.is_some() {
913 return Err(ErrorCode::NotSupported(
914 "REPLACE SINK with auto schema change is not supported yet".to_owned(),
915 "drop and recreate this auto schema change sink".to_owned(),
916 )
917 .into());
918 }
919 if sink_is_exactly_once(&sink.properties)? {
920 return Err(ErrorCode::NotSupported(
921 "REPLACE SINK does not support exactly-once sinks yet".to_owned(),
922 "set is_exactly_once=false or recreate the sink manually".to_owned(),
923 )
924 .into());
925 }
926 sink.clone()
927 };
928
929 Ok(SinkCreateMode::Replace { original_sink })
930}
931
932pub fn fetch_incoming_sinks(
933 session: &Arc<SessionImpl>,
934 table: &TableCatalog,
935) -> Result<Vec<Arc<SinkCatalog>>> {
936 let reader = session.env().catalog_reader().read_guard();
937 let schema = reader.get_schema_by_id(table.database_id, table.schema_id)?;
938 let Some(incoming_sinks) = schema.table_incoming_sinks(table.id) else {
939 return Ok(vec![]);
940 };
941 let mut sinks = vec![];
942 for sink_id in incoming_sinks {
943 sinks.push(
944 schema
945 .get_sink_by_id(*sink_id)
946 .expect("should exist")
947 .clone(),
948 );
949 }
950 Ok(sinks)
951}
952
953fn derive_sink_to_table_expr(
954 sink_schema: &Schema,
955 idx: usize,
956 target_type: &DataType,
957) -> Result<ExprImpl> {
958 let input_type = &sink_schema.fields()[idx].data_type;
959
960 if !target_type.equals_datatype(input_type) {
961 bail!(
962 "column type mismatch: {:?} vs {:?}, column name: {:?}",
963 target_type,
964 input_type,
965 sink_schema.fields()[idx].name
966 );
967 } else {
968 Ok(ExprImpl::InputRef(Box::new(InputRef::new(
969 idx,
970 input_type.clone(),
971 ))))
972 }
973}
974
975pub(crate) fn derive_default_column_project_for_sink(
976 sink: &SinkCatalog,
977 sink_schema: &Schema,
978 columns: &[ColumnCatalog],
979 user_specified_columns: bool,
980) -> Result<Vec<ExprImpl>> {
981 assert_eq!(sink.full_schema().len(), sink_schema.len());
982
983 let default_column_exprs = TableCatalog::default_column_exprs(columns);
984
985 let mut exprs = vec![];
986
987 let sink_visible_col_idxes = sink
988 .full_columns()
989 .iter()
990 .positions(|c| !c.is_hidden())
991 .collect_vec();
992 let sink_visible_col_idxes_by_name = sink
993 .full_columns()
994 .iter()
995 .enumerate()
996 .filter(|(_, c)| !c.is_hidden())
997 .map(|(i, c)| (c.name(), i))
998 .collect::<BTreeMap<_, _>>();
999
1000 for (idx, column) in columns.iter().enumerate() {
1001 if !column.can_dml() {
1002 continue;
1003 }
1004
1005 let default_col_expr =
1006 || -> ExprImpl { rewrite_now_to_proctime(default_column_exprs[idx].clone()) };
1007
1008 let sink_col_expr = |sink_col_idx: usize| -> Result<ExprImpl> {
1009 derive_sink_to_table_expr(sink_schema, sink_col_idx, column.data_type())
1010 };
1011
1012 if user_specified_columns {
1016 if let Some(idx) = sink_visible_col_idxes_by_name.get(column.name()) {
1017 exprs.push(sink_col_expr(*idx)?);
1018 } else {
1019 exprs.push(default_col_expr());
1020 }
1021 } else {
1022 if idx < sink_visible_col_idxes.len() {
1023 exprs.push(sink_col_expr(sink_visible_col_idxes[idx])?);
1024 } else {
1025 exprs.push(default_col_expr());
1026 };
1027 }
1028 }
1029 Ok(exprs)
1030}
1031
1032fn bind_sink_format_desc(
1036 session: &SessionImpl,
1037 value: FormatEncodeOptions,
1038) -> Result<SinkFormatDesc> {
1039 use risingwave_connector::sink::catalog::{SinkEncode, SinkFormat};
1040 use risingwave_connector::sink::encoder::TimestamptzHandlingMode;
1041 use risingwave_sqlparser::ast::{Encode as E, Format as F};
1042
1043 let format = match value.format {
1044 F::Plain => SinkFormat::AppendOnly,
1045 F::Upsert => SinkFormat::Upsert,
1046 F::Debezium => SinkFormat::Debezium,
1047 f @ (F::Native | F::DebeziumMongo | F::Maxwell | F::Canal | F::None) => {
1048 return Err(ErrorCode::BindError(format!("sink format unsupported: {f}")).into());
1049 }
1050 };
1051 let encode = match value.row_encode {
1052 E::Json => SinkEncode::Json,
1053 E::Protobuf => SinkEncode::Protobuf,
1054 E::Avro => SinkEncode::Avro,
1055 E::Template => SinkEncode::Template,
1056 E::Parquet => SinkEncode::Parquet,
1057 E::Bytes => SinkEncode::Bytes,
1058 e @ (E::Native | E::Csv | E::None | E::Text) => {
1059 return Err(ErrorCode::BindError(format!("sink encode unsupported: {e}")).into());
1060 }
1061 };
1062
1063 let mut key_encode = None;
1064 if let Some(encode) = value.key_encode {
1065 match encode {
1066 E::Text => key_encode = Some(SinkEncode::Text),
1067 E::Bytes => key_encode = Some(SinkEncode::Bytes),
1068 _ => {
1069 return Err(ErrorCode::BindError(format!(
1070 "sink key encode unsupported: {encode}, only TEXT and BYTES supported"
1071 ))
1072 .into());
1073 }
1074 }
1075 }
1076
1077 let (props, connection_type_flag, schema_registry_conn_ref) =
1078 resolve_connection_ref_and_secret_ref(
1079 WithOptions::try_from(value.row_options.as_slice())?,
1080 session,
1081 Some(TelemetryDatabaseObject::Sink),
1082 )?;
1083 ensure_connection_type_allowed(
1084 connection_type_flag,
1085 &SINK_ALLOWED_CONNECTION_SCHEMA_REGISTRY,
1086 )?;
1087 let (mut options, secret_refs) = props.into_parts();
1088
1089 options
1090 .entry(TimestamptzHandlingMode::OPTION_KEY.to_owned())
1091 .or_insert(TimestamptzHandlingMode::FRONTEND_DEFAULT.to_owned());
1092
1093 Ok(SinkFormatDesc {
1094 format,
1095 encode,
1096 options,
1097 secret_refs,
1098 key_encode,
1099 connection_id: schema_registry_conn_ref,
1100 })
1101}
1102
1103static CONNECTORS_COMPATIBLE_FORMATS: LazyLock<HashMap<String, HashMap<Format, Vec<Encode>>>> =
1104 LazyLock::new(|| {
1105 use risingwave_connector::sink::Sink as _;
1106 use risingwave_connector::sink::file_sink::azblob::AzblobSink;
1107 use risingwave_connector::sink::file_sink::fs::FsSink;
1108 use risingwave_connector::sink::file_sink::gcs::GcsSink;
1109 use risingwave_connector::sink::file_sink::opendal_sink::FileSink;
1110 use risingwave_connector::sink::file_sink::s3::S3Sink;
1111 use risingwave_connector::sink::file_sink::webhdfs::WebhdfsSink;
1112 use risingwave_connector::sink::google_pubsub::GooglePubSubSink;
1113 use risingwave_connector::sink::kafka::KafkaSink;
1114 use risingwave_connector::sink::kinesis::KinesisSink;
1115 use risingwave_connector::sink::mqtt::MqttSink;
1116 use risingwave_connector::sink::pulsar::PulsarSink;
1117 use risingwave_connector::sink::redis::RedisSink;
1118
1119 convert_args!(hashmap!(
1120 GooglePubSubSink::SINK_NAME => hashmap!(
1121 Format::Plain => vec![Encode::Json],
1122 ),
1123 KafkaSink::SINK_NAME => hashmap!(
1124 Format::Plain => vec![Encode::Json, Encode::Avro, Encode::Protobuf, Encode::Bytes],
1125 Format::Upsert => vec![Encode::Json, Encode::Avro, Encode::Protobuf],
1126 Format::Debezium => vec![Encode::Json],
1127 ),
1128 FileSink::<S3Sink>::SINK_NAME => hashmap!(
1129 Format::Plain => vec![Encode::Parquet, Encode::Json],
1130 ),
1131 FileSink::<SnowflakeSink>::SINK_NAME => hashmap!(
1132 Format::Plain => vec![Encode::Parquet, Encode::Json],
1133 ),
1134 FileSink::<GcsSink>::SINK_NAME => hashmap!(
1135 Format::Plain => vec![Encode::Parquet, Encode::Json],
1136 ),
1137 FileSink::<AzblobSink>::SINK_NAME => hashmap!(
1138 Format::Plain => vec![Encode::Parquet, Encode::Json],
1139 ),
1140 FileSink::<WebhdfsSink>::SINK_NAME => hashmap!(
1141 Format::Plain => vec![Encode::Parquet, Encode::Json],
1142 ),
1143 FileSink::<FsSink>::SINK_NAME => hashmap!(
1144 Format::Plain => vec![Encode::Parquet, Encode::Json],
1145 ),
1146 KinesisSink::SINK_NAME => hashmap!(
1147 Format::Plain => vec![Encode::Json],
1148 Format::Upsert => vec![Encode::Json],
1149 Format::Debezium => vec![Encode::Json],
1150 ),
1151 MqttSink::SINK_NAME => hashmap!(
1152 Format::Plain => vec![Encode::Json, Encode::Protobuf],
1153 ),
1154 PulsarSink::SINK_NAME => hashmap!(
1155 Format::Plain => vec![Encode::Json],
1156 Format::Upsert => vec![Encode::Json],
1157 Format::Debezium => vec![Encode::Json],
1158 ),
1159 RedisSink::SINK_NAME => hashmap!(
1160 Format::Plain => vec![Encode::Json, Encode::Template],
1161 Format::Upsert => vec![Encode::Json, Encode::Template],
1162 ),
1163 ))
1164 });
1165
1166pub fn validate_compatibility(connector: &str, format_desc: &FormatEncodeOptions) -> Result<()> {
1167 let compatible_formats = CONNECTORS_COMPATIBLE_FORMATS
1168 .get(connector)
1169 .ok_or_else(|| {
1170 ErrorCode::BindError(format!(
1171 "connector {} is not supported by FORMAT ... ENCODE ... syntax",
1172 connector
1173 ))
1174 })?;
1175 let compatible_encodes = compatible_formats.get(&format_desc.format).ok_or_else(|| {
1176 ErrorCode::BindError(format!(
1177 "connector {} does not support format {:?}",
1178 connector, format_desc.format
1179 ))
1180 })?;
1181 if !compatible_encodes.contains(&format_desc.row_encode) {
1182 return Err(ErrorCode::BindError(format!(
1183 "connector {} does not support format {:?} with encode {:?}",
1184 connector, format_desc.format, format_desc.row_encode
1185 ))
1186 .into());
1187 }
1188
1189 if let Some(encode) = &format_desc.key_encode
1191 && connector != KAFKA_SINK
1192 && matches!(encode, Encode::Bytes)
1193 {
1194 return Err(ErrorCode::BindError(format!(
1195 "key encode bytes only works with kafka connector, but found {}",
1196 connector
1197 ))
1198 .into());
1199 }
1200
1201 Ok(())
1202}
1203
1204#[cfg(test)]
1205pub mod tests {
1206 use std::collections::BTreeMap;
1207
1208 use risingwave_common::catalog::{CreateType, DEFAULT_DATABASE_NAME, DEFAULT_SCHEMA_NAME};
1209 use risingwave_common::config::FrontendConfig;
1210 use risingwave_connector::sink::Sink;
1211 use risingwave_connector::sink::snowflake_redshift::redshift::RedshiftSink;
1212 use risingwave_connector::sink::snowflake_redshift::snowflake::SnowflakeV2Sink;
1213 use risingwave_connector::{
1214 SINK_CREATE_TABLE_IF_NOT_EXISTS_KEY, SINK_INTERMEDIATE_TABLE_NAME, SINK_TARGET_TABLE_NAME,
1215 };
1216
1217 use crate::WithOptionsSecResolved;
1218 use crate::catalog::root_catalog::SchemaPath;
1219 use crate::handler::create_sink::maybe_fill_intermediate_table_name;
1220 use crate::test_utils::{LocalFrontend, PROTO_FILE_DATA, create_proto_file};
1221
1222 #[test]
1223 fn test_fill_intermediate_table_name_for_upsert_auto_create_only() {
1224 for (connector, sink_type, create_table, should_fill) in [
1225 (SnowflakeV2Sink::SINK_NAME, "append-only", false, false),
1226 (SnowflakeV2Sink::SINK_NAME, "append-only", true, false),
1227 (SnowflakeV2Sink::SINK_NAME, "upsert", false, false),
1228 (SnowflakeV2Sink::SINK_NAME, "upsert", true, true),
1229 (RedshiftSink::SINK_NAME, "upsert", true, true),
1230 ("jdbc", "upsert", true, false),
1231 ] {
1232 let mut options = WithOptionsSecResolved::without_secrets(BTreeMap::from([
1233 ("type".to_owned(), sink_type.to_owned()),
1234 (
1235 SINK_CREATE_TABLE_IF_NOT_EXISTS_KEY.to_owned(),
1236 create_table.to_string(),
1237 ),
1238 (SINK_TARGET_TABLE_NAME.to_owned(), "target".to_owned()),
1239 ]));
1240
1241 maybe_fill_intermediate_table_name(&mut options, connector, "sink").unwrap();
1242
1243 assert_eq!(
1244 options.contains_key(SINK_INTERMEDIATE_TABLE_NAME),
1245 should_fill,
1246 "connector={connector}, type={sink_type}, create_table={create_table}"
1247 );
1248 }
1249 }
1250
1251 #[test]
1252 fn test_fill_intermediate_table_name_preserves_user_value() {
1253 let mut options = WithOptionsSecResolved::without_secrets(BTreeMap::from([
1254 ("type".to_owned(), "upsert".to_owned()),
1255 (
1256 SINK_CREATE_TABLE_IF_NOT_EXISTS_KEY.to_owned(),
1257 "true".to_owned(),
1258 ),
1259 (SINK_TARGET_TABLE_NAME.to_owned(), "target".to_owned()),
1260 (
1261 SINK_INTERMEDIATE_TABLE_NAME.to_owned(),
1262 "custom_cdc".to_owned(),
1263 ),
1264 ]));
1265
1266 maybe_fill_intermediate_table_name(&mut options, SnowflakeV2Sink::SINK_NAME, "sink")
1267 .unwrap();
1268
1269 assert_eq!(
1270 options.get(SINK_INTERMEDIATE_TABLE_NAME).unwrap(),
1271 "custom_cdc"
1272 );
1273 }
1274
1275 #[test]
1276 fn test_sink_replace_requires_exactly_once_state_defaults() {
1277 let properties = BTreeMap::from([("connector".to_owned(), "iceberg".to_owned())]);
1278 assert!(super::sink_is_exactly_once(&properties).unwrap());
1279
1280 let properties = BTreeMap::from([
1281 ("connector".to_owned(), "iceberg".to_owned()),
1282 ("is_exactly_once".to_owned(), "false".to_owned()),
1283 ]);
1284 assert!(!super::sink_is_exactly_once(&properties).unwrap());
1285
1286 let properties = BTreeMap::from([
1287 ("connector".to_owned(), "jdbc".to_owned()),
1288 ("is_exactly_once".to_owned(), "true".to_owned()),
1289 ]);
1290 assert!(!super::sink_is_exactly_once(&properties).unwrap());
1291 }
1292
1293 #[tokio::test]
1294 async fn test_create_sink_handler() {
1295 let proto_file = create_proto_file(PROTO_FILE_DATA);
1296 let sql = format!(
1297 r#"CREATE SOURCE t1
1298 WITH (connector = 'kafka', kafka.topic = 'abc', kafka.brokers = 'localhost:1001')
1299 FORMAT PLAIN ENCODE PROTOBUF (message = '.test.TestRecord', schema.location = 'file://{}')"#,
1300 proto_file.path().to_str().unwrap()
1301 );
1302 let frontend = LocalFrontend::new(Default::default()).await;
1303 frontend.run_sql(sql).await.unwrap();
1304
1305 let sql = "create materialized view mv1 as select t1.country from t1;";
1306 frontend.run_sql(sql).await.unwrap();
1307
1308 let sql = r#"CREATE SINK snk1 FROM mv1
1309 WITH (connector = 'jdbc', mysql.endpoint = '127.0.0.1:3306', mysql.table =
1310 '<table_name>', mysql.database = '<database_name>', mysql.user = '<user_name>',
1311 mysql.password = '<password>', type = 'append-only', force_append_only = 'true');"#.to_owned();
1312 frontend.run_sql(sql).await.unwrap();
1313
1314 let session = frontend.session_ref();
1315 let catalog_reader = session.env().catalog_reader().read_guard();
1316 let schema_path = SchemaPath::Name(DEFAULT_SCHEMA_NAME);
1317
1318 let (source, _) = catalog_reader
1320 .get_source_by_name(DEFAULT_DATABASE_NAME, schema_path, "t1")
1321 .unwrap();
1322 assert_eq!(source.name, "t1");
1323
1324 let (table, schema_name) = catalog_reader
1326 .get_created_table_by_name(DEFAULT_DATABASE_NAME, schema_path, "mv1")
1327 .unwrap();
1328 assert_eq!(table.name(), "mv1");
1329 let schema_name = schema_name.to_owned();
1330
1331 let (sink, _) = catalog_reader
1333 .get_created_sink_by_name(
1334 DEFAULT_DATABASE_NAME,
1335 SchemaPath::Name(&schema_name),
1336 "snk1",
1337 )
1338 .unwrap();
1339 assert_eq!(sink.name, "snk1");
1340 drop(catalog_reader);
1341
1342 let sql = r#"REPLACE SINK snk1 FROM mv1
1343 WITH (connector = 'jdbc', mysql.endpoint = '127.0.0.1:3306', mysql.table =
1344 '<table_name>', mysql.database = '<database_name>', mysql.user = '<user_name>',
1345 mysql.password = '<password>', type = 'append-only', force_append_only = 'true');"#.to_owned();
1346 frontend.run_sql(sql).await.unwrap();
1347
1348 let catalog_reader = session.env().catalog_reader().read_guard();
1349 let (sink, _) = catalog_reader
1350 .get_created_sink_by_name(
1351 DEFAULT_DATABASE_NAME,
1352 SchemaPath::Name(&schema_name),
1353 "snk1",
1354 )
1355 .unwrap();
1356 assert_eq!(sink.name, "snk1");
1357 assert_eq!(sink.create_type, CreateType::Foreground);
1360 }
1361
1362 #[tokio::test]
1363 async fn test_create_fs_sink_requires_frontend_config() {
1364 let frontend = LocalFrontend::with_frontend_config(
1365 Default::default(),
1366 FrontendConfig {
1367 unsafe_enable_local_fs_connector: false,
1368 ..Default::default()
1369 },
1370 )
1371 .await;
1372 frontend.run_sql("CREATE TABLE t(v int);").await.unwrap();
1373 frontend
1374 .run_sql("CREATE MATERIALIZED VIEW mv AS SELECT * FROM t;")
1375 .await
1376 .unwrap();
1377
1378 let err = frontend
1379 .run_sql(
1380 r#"CREATE SINK local_sink FROM mv
1381 WITH (
1382 connector = 'fs',
1383 fs.path = '/tmp/rw-local-sink',
1384 type = 'append-only',
1385 force_append_only = 'true'
1386 ) FORMAT PLAIN ENCODE JSON (force_append_only = 'true');"#
1387 .to_owned(),
1388 )
1389 .await
1390 .unwrap_err();
1391
1392 assert!(
1393 err.to_string()
1394 .contains("frontend.unsafe_enable_local_fs_connector = true"),
1395 "{err:?}"
1396 );
1397 }
1398}