Skip to main content

risingwave_frontend/handler/
create_sink.rs

1// Copyright 2022 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::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
113// used to store result of `gen_sink_plan`
114pub 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    // These are frontend-level streaming job options. They must not be passed to connector
166    // property validation.
167    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 not using connection, we don't need to check connector match connection type
209    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    // Used for debezium's table name
260    let sink_from_table_name;
261    // `true` means that sink statement has the form: `CREATE SINK s1 FROM ...`
262    // `false` means that sink statement has the form: `CREATE SINK s1 AS <query>`
263    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        // If column names not specified, use the name in the bound query, which is equal with the plan root's original field name.
402        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        // Case A: new syntax `format ... encode ...`
412        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            // Case B: old syntax `type = '...'`
418            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            // Case C: no format + encode required
429            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    // TODO(rc): To be consistent with UDF dependency check, we should collect relation dependencies
521    // during binding instead of visiting the optimized plan.
522    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
582// This function is used to return partition compute info for a sink. More details refer in `PartitionComputeInfo`.
583// Return:
584// `Some(PartitionComputeInfo)` if the sink need to compute partition.
585// `None` if the sink does not need to compute partition.
586pub 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    // TODO: check table if exists
607    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    // Separate the partition spec into two parts: sparse partition and range partition.
617    // Sparse partition means that the data distribution is more sparse at a given time.
618    // Range partition means that the data distribution is likely same at a given time.
619    // Only compute the partition and shuffle by them for the sparse partition.
620    let has_sparse_partition = partition_spec.fields().iter().any(|f| match f.transform {
621        // Sparse partition
622        Transform::Identity | Transform::Truncate(_) | Transform::Bucket(_) => true,
623        // Range partition
624        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 users specified the columns to be inserted e.g. `CREATE SINK s INTO t(a, b)`, the expressions of `Project` will be generated accordingly.
1013        // The missing columns will be filled with default value (`null` if not explicitly defined).
1014        // Otherwise, e.g. `CREATE SINK s INTO t`, the columns will be matched by their order in `select` query and the target table.
1015        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
1032/// Transforms the (format, encode, options) from sqlparser AST into an internal struct `SinkFormatDesc`.
1033/// This is an analogy to (part of) [`crate::handler::create_source::bind_columns_from_source`]
1034/// which transforms sqlparser AST `SourceSchemaV2` into `StreamSourceInfo`.
1035fn 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    // only allow Kafka connector work with `bytes` as key encode
1190    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        // Check source exists.
1319        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        // Check table exists.
1325        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        // Check sink exists.
1332        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        // Frontend leaves the replacement job foreground for the meta foreground wait path.
1358        // Meta switches it to background when marking the job Creating during cutover.
1359        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}