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