Skip to main content

risingwave_frontend/handler/
cdc.rs

1// Copyright 2026 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 anyhow::{Context, anyhow};
16use fancy_regex::Regex;
17use itertools::Itertools;
18use risingwave_common::catalog::{CdcKeyComparison, ColumnCatalog};
19use risingwave_connector::WithOptionsSecResolved;
20use risingwave_connector::source::UPSTREAM_SOURCE_KEY;
21use risingwave_connector::source::cdc::external::{
22    DATABASE_NAME_KEY, ExternalTableConfig, ExternalTableImpl, SCHEMA_NAME_KEY, SchemaTableName,
23    TABLE_NAME_KEY,
24};
25use risingwave_connector::source::cdc::{
26    MYSQL_CDC_CONNECTOR, POSTGRES_CDC_CONNECTOR, SQL_SERVER_CDC_CONNECTOR,
27};
28use risingwave_sqlparser::ast::{ColumnDef, ColumnOption, SourceWatermark, TableConstraint};
29use thiserror_ext::AsReport;
30
31use crate::error::{ErrorCode, Result, RwError};
32use crate::handler::create_source::reject_variant_columns;
33use crate::handler::create_table::{bind_sql_columns, bind_sql_pk_names, bind_table_constraints};
34
35/// Derive connector properties and normalize `external_table_name` for CDC tables.
36///
37/// Returns (`connector_properties`, `normalized_external_table_name`) where:
38/// - For SQL Server: Normalizes 'db.schema.table' (3 parts) to 'schema.table' (2 parts),
39///   because users can optionally include database name for verification, but it needs to be
40///   stripped to match the format returned by Debezium's `extract_table_name()`.
41/// - For MySQL/Postgres: Returns the original `external_table_name` unchanged.
42pub(crate) fn derive_with_options_for_cdc_table(
43    source_with_properties: &WithOptionsSecResolved,
44    external_table_name: String,
45) -> Result<(WithOptionsSecResolved, String)> {
46    // we should remove the prefix from `full_table_name`
47    let source_database_name: &str = source_with_properties
48        .get("database.name")
49        .ok_or_else(|| anyhow!("The source with properties does not contain 'database.name'"))?
50        .as_str();
51    let mut with_options = source_with_properties.clone();
52    if let Some(connector) = source_with_properties.get(UPSTREAM_SOURCE_KEY) {
53        match connector.as_str() {
54            MYSQL_CDC_CONNECTOR => {
55                // MySQL doesn't allow '.' in database name and table name, so we can split the
56                // external table name by '.' to get the table name
57                let (db_name, table_name) = external_table_name.split_once('.').ok_or_else(|| {
58                    anyhow!("The upstream table name must contain database name prefix, e.g. 'database.table'")
59                })?;
60                // We allow multiple database names in the source definition
61                if !source_database_name
62                    .split(',')
63                    .map(|s| s.trim())
64                    .any(|name| name == db_name)
65                {
66                    return Err(anyhow!(
67                        "The database name `{}` in the FROM clause is not included in the database name `{}` in source definition",
68                        db_name,
69                        source_database_name
70                    ).into());
71                }
72                with_options.insert(DATABASE_NAME_KEY.into(), db_name.into());
73                with_options.insert(TABLE_NAME_KEY.into(), table_name.into());
74                // Return original external_table_name unchanged for MySQL
75                return Ok((with_options, external_table_name));
76            }
77            POSTGRES_CDC_CONNECTOR => {
78                let (schema_name, table_name) =
79                    parse_postgres_cdc_external_table_name(&external_table_name)?;
80
81                // insert 'schema.name' into connect properties
82                with_options.insert(SCHEMA_NAME_KEY.into(), schema_name);
83                with_options.insert(TABLE_NAME_KEY.into(), table_name);
84                // Return original external_table_name unchanged for Postgres
85                return Ok((with_options, external_table_name));
86            }
87            SQL_SERVER_CDC_CONNECTOR => {
88                // SQL Server external table name must be in one of two formats:
89                // 1. 'schemaName.tableName' (2 parts) - database is already specified in source
90                // 2. 'databaseName.schemaName.tableName' (3 parts) - for explicit verification
91                //
92                // We do NOT allow single table name (e.g., 't') because:
93                // - Unlike database name (already in source), schema name is NOT pre-specified
94                // - User must explicitly provide schema (even if it's 'dbo')
95                let parts: Vec<&str> = external_table_name.split('.').collect();
96                let (schema_name, table_name) = match parts.len() {
97                    3 => {
98                        // Format: database.schema.table
99                        // Verify that the database name matches the one in source definition
100                        let db_name = parts[0];
101                        let schema_name = parts[1];
102                        let table_name = parts[2];
103
104                        if db_name != source_database_name {
105                            return Err(anyhow!(
106                                "The database name '{}' in FROM clause does not match the database name '{}' specified in source definition. \
107                                 You can either use 'schema.table' format (recommended) or ensure the database name matches.",
108                                db_name,
109                                source_database_name
110                            ).into());
111                        }
112                        (schema_name, table_name)
113                    }
114                    2 => {
115                        // Format: schema.table (recommended)
116                        // Database name is taken from source definition
117                        let schema_name = parts[0];
118                        let table_name = parts[1];
119                        (schema_name, table_name)
120                    }
121                    1 => {
122                        // Format: table only
123                        // Reject with clear error message
124                        return Err(anyhow!(
125                            "Invalid table name format '{}'. For SQL Server CDC, you must specify the schema name. \
126                             Use 'schema.table' format (e.g., 'dbo.{}') or 'database.schema.table' format (e.g., '{}.dbo.{}').",
127                            external_table_name,
128                            external_table_name,
129                            source_database_name,
130                            external_table_name
131                        ).into());
132                    }
133                    _ => {
134                        // Invalid format (4+ parts or empty)
135                        return Err(anyhow!(
136                            "Invalid table name format '{}'. Expected 'schema.table' or 'database.schema.table'.",
137                            external_table_name
138                        ).into());
139                    }
140                };
141
142                // Insert schema and table names into connector properties
143                with_options.insert(SCHEMA_NAME_KEY.into(), schema_name.into());
144                with_options.insert(TABLE_NAME_KEY.into(), table_name.into());
145
146                // Normalize external_table_name to 'schema.table' format
147                // This ensures consistency with extract_table_name() in message.rs
148                let normalized_external_table_name = format!("{}.{}", schema_name, table_name);
149                return Ok((with_options, normalized_external_table_name));
150            }
151            _ => {
152                return Err(RwError::from(anyhow!(
153                    "connector {} is not supported for cdc table",
154                    connector
155                )));
156            }
157        };
158    }
159    unreachable!("All valid CDC connectors should have returned by now")
160}
161
162/// Parse the schema/table name from the CDC `TABLE` clause.
163///
164/// Column names do not need the same parsing here: wildcard schema derivation reads
165/// them from PostgreSQL catalogs after the exact table has been identified.
166fn parse_postgres_cdc_external_table_name(external_table_name: &str) -> Result<(String, String)> {
167    let mut parts = vec![];
168    let mut current = String::new();
169    let mut chars = external_table_name.chars().peekable();
170    let mut in_quote = false;
171    let mut just_closed_quote = false;
172
173    while let Some(ch) = chars.next() {
174        if in_quote {
175            if ch == '"' {
176                if chars.peek() == Some(&'"') {
177                    current.push('"');
178                    chars.next();
179                } else {
180                    in_quote = false;
181                    just_closed_quote = true;
182                }
183            } else {
184                current.push(ch);
185            }
186        } else {
187            match ch {
188                '.' => {
189                    if current.is_empty() {
190                        return Err(anyhow!(
191                            "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
192                            external_table_name
193                        )
194                        .into());
195                    }
196                    parts.push(std::mem::take(&mut current));
197                    just_closed_quote = false;
198                }
199                '"' if current.is_empty() => {
200                    in_quote = true;
201                }
202                '"' => {
203                    return Err(anyhow!(
204                        "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
205                        external_table_name
206                    )
207                    .into());
208                }
209                _ if just_closed_quote => {
210                    return Err(anyhow!(
211                        "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
212                        external_table_name
213                    )
214                    .into());
215                }
216                _ => current.push(ch),
217            }
218        }
219    }
220
221    if in_quote || current.is_empty() {
222        return Err(anyhow!(
223            "Invalid Postgres CDC table name '{}'. Expected 'schema.table'.",
224            external_table_name
225        )
226        .into());
227    }
228    parts.push(current);
229
230    if let [schema_name, table_name] = parts.as_slice() {
231        Ok((schema_name.clone(), table_name.clone()))
232    } else {
233        Err(
234            anyhow!("The upstream table name must contain schema name prefix, e.g. 'public.table'")
235                .into(),
236        )
237    }
238}
239
240/// Reject a CDC table when a primary-key column is filtered out of Debezium change-event
241/// values via `debezium.column.exclude.list` or `debezium.column.include.list`.
242///
243/// Debezium's column filters only apply to the change-event **value** payload. Message keys are
244/// always built from the upstream PRIMARY KEY and are not affected. If a PK column is filtered out
245/// of the value, RisingWave reads NULL for that PK column from the payload, causing silent data
246/// corruption: UPDATE turns into a fresh INSERT (PK mismatch with the original row) and DELETE
247/// silently no-ops.
248///
249/// Debezium entries are regex patterns matched against the fully qualified column name
250/// `<namespace>.<table>.<column>`, where namespace is `schema` for Postgres / SQL Server and
251/// `database` for MySQL.
252pub(crate) fn reject_pk_filtered_by_debezium_column_filter(
253    pk_names: &[String],
254    cdc_with_options: &WithOptionsSecResolved,
255) -> Result<()> {
256    const EXCLUDE_KEY: &str = "debezium.column.exclude.list";
257    const INCLUDE_KEY: &str = "debezium.column.include.list";
258
259    let st = SchemaTableName::from_properties(cdc_with_options.as_plaintext());
260    reject_pk_filtered_by_debezium_column_filter_inner(
261        pk_names,
262        &st,
263        cdc_with_options.get(EXCLUDE_KEY).map(String::as_str),
264        cdc_with_options.get(INCLUDE_KEY).map(String::as_str),
265    )
266}
267
268fn reject_pk_filtered_by_debezium_column_filter_inner(
269    pk_names: &[String],
270    st: &SchemaTableName,
271    exclude_list: Option<&str>,
272    include_list: Option<&str>,
273) -> Result<()> {
274    const EXCLUDE_KEY: &str = "debezium.column.exclude.list";
275    const INCLUDE_KEY: &str = "debezium.column.include.list";
276
277    let pk_full_names = pk_names
278        .iter()
279        .map(|pk| (pk, format!("{}.{}.{}", st.schema_name, st.table_name, pk)))
280        .collect_vec();
281
282    if let Some(exclude_list) = exclude_list {
283        let patterns = compile_debezium_column_filter_patterns(EXCLUDE_KEY, exclude_list)?;
284        for (pk, pk_full_name) in &pk_full_names {
285            for (pattern, regex) in &patterns {
286                if regex.is_match(pk_full_name).map_err(|err| {
287                    ErrorCode::InvalidInputSyntax(format!(
288                        "failed to evaluate Debezium column filter pattern `{pattern}` in `{EXCLUDE_KEY}`: {}",
289                        err.as_report()
290                    ))
291                })? {
292                    return Err(ErrorCode::InvalidInputSyntax(format!(
293                        "primary key column `{pk}` is excluded by `{EXCLUDE_KEY}` pattern \
294                         `{pattern}`. Excluding a PK column causes silent data corruption: \
295                         Debezium keeps the PK in the message key but drops it from the payload, \
296                         so RisingWave cannot match UPDATE/DELETE events against the original row."
297                    ))
298                    .into());
299                }
300            }
301        }
302    }
303
304    if let Some(include_list) = include_list {
305        let patterns = compile_debezium_column_filter_patterns(INCLUDE_KEY, include_list)?;
306        for (pk, pk_full_name) in &pk_full_names {
307            let mut included = false;
308            for (_, regex) in &patterns {
309                if regex.is_match(pk_full_name).map_err(|err| {
310                    ErrorCode::InvalidInputSyntax(format!(
311                        "failed to evaluate Debezium column filter pattern in `{INCLUDE_KEY}`: {}",
312                        err.as_report()
313                    ))
314                })? {
315                    included = true;
316                    break;
317                }
318            }
319            if !included {
320                return Err(ErrorCode::InvalidInputSyntax(format!(
321                    "primary key column `{pk}` is not included by `{INCLUDE_KEY}`. Omitting a PK \
322                     column causes silent data corruption: Debezium keeps the PK in the message key \
323                     but drops it from the payload, so RisingWave cannot match UPDATE/DELETE events \
324                     against the original row."
325                ))
326                .into());
327            }
328        }
329    }
330
331    Ok(())
332}
333
334fn compile_debezium_column_filter_patterns(
335    key: &str,
336    filter_list: &str,
337) -> Result<Vec<(String, Regex)>> {
338    filter_list
339        .split(',')
340        .map(str::trim)
341        .filter(|pattern| !pattern.is_empty())
342        .map(|pattern| {
343            let anchored_pattern = format!("(?i:^(?:{pattern})$)");
344            let regex = Regex::new(&anchored_pattern).map_err(|err| {
345                ErrorCode::InvalidInputSyntax(format!(
346                    "invalid Debezium column filter pattern `{pattern}` in `{key}`: {}",
347                    err.as_report()
348                ))
349            })?;
350            Ok((pattern.to_owned(), regex))
351        })
352        .collect()
353}
354
355// For both table from cdc source and table with cdc connector
356pub(crate) fn not_null_check_for_cdc_table(
357    wildcard_idx: &Option<usize>,
358    column_defs: &Vec<ColumnDef>,
359) -> Result<()> {
360    if !wildcard_idx.is_some()
361        && column_defs.iter().any(|col| {
362            col.options
363                .iter()
364                .any(|opt| matches!(opt.option, ColumnOption::NotNull))
365        })
366    {
367        return Err(ErrorCode::NotSupported(
368            "CDC table with NOT NULL constraint is not supported".to_owned(),
369            "Please remove the NOT NULL constraint for columns".to_owned(),
370        )
371        .into());
372    }
373    Ok(())
374}
375
376// Only for table from cdc source
377pub(crate) fn sanity_check_for_table_on_cdc_source(
378    append_only: bool,
379    column_defs: &Vec<ColumnDef>,
380    wildcard_idx: &Option<usize>,
381    constraints: &Vec<TableConstraint>,
382    source_watermarks: &Vec<SourceWatermark>,
383) -> Result<()> {
384    // wildcard cannot be used with column definitions
385    if wildcard_idx.is_some() && !column_defs.is_empty() {
386        return Err(ErrorCode::NotSupported(
387            "wildcard(*) and column definitions cannot be used together".to_owned(),
388            "Remove the wildcard or column definitions".to_owned(),
389        )
390        .into());
391    }
392
393    // cdc table must have primary key constraint or primary key column
394    if !wildcard_idx.is_some()
395        && !constraints.iter().any(|c| {
396            matches!(
397                c,
398                TableConstraint::Unique {
399                    is_primary: true,
400                    ..
401                }
402            )
403        })
404        && !column_defs.iter().any(|col| {
405            col.options
406                .iter()
407                .any(|opt| matches!(opt.option, ColumnOption::Unique { is_primary: true }))
408        })
409    {
410        return Err(ErrorCode::NotSupported(
411            "CDC table without primary key constraint is not supported".to_owned(),
412            "Please define a primary key".to_owned(),
413        )
414        .into());
415    }
416
417    if append_only {
418        return Err(ErrorCode::NotSupported(
419            "append only modifier on the table created from a CDC source".into(),
420            "Remove the APPEND ONLY clause".into(),
421        )
422        .into());
423    }
424
425    if !source_watermarks.is_empty()
426        && source_watermarks
427            .iter()
428            .any(|watermark| !watermark.with_ttl)
429    {
430        return Err(ErrorCode::NotSupported(
431            "non-TTL watermark defined on the table created from a CDC source".into(),
432            "Use `WATERMARK ... WITH TTL` instead.".into(),
433        )
434        .into());
435    }
436
437    Ok(())
438}
439
440/// Derive the schema of a CDC table from its upstream external table.
441pub(crate) async fn bind_cdc_table_schema_externally(
442    cdc_with_options: WithOptionsSecResolved,
443) -> Result<(Vec<ColumnCatalog>, Vec<String>, Vec<CdcKeyComparison>)> {
444    let (options, secret_refs) = cdc_with_options.into_parts();
445    let config = ExternalTableConfig::try_from_btreemap(options, secret_refs)
446        .context("failed to extract external table config")?;
447
448    let table = ExternalTableImpl::connect(config)
449        .await
450        .context("failed to auto derive table schema")?;
451
452    let pk_names = table.pk_names().clone();
453    let pk_comparisons = table.pk_column_comparisons(&pk_names)?;
454    Ok((
455        table
456            .column_descs()
457            .iter()
458            .cloned()
459            .map(|column_desc| ColumnCatalog {
460                column_desc,
461                is_hidden: false,
462            })
463            .collect(),
464        pk_names,
465        pk_comparisons,
466    ))
467}
468
469pub(crate) async fn bind_cdc_pk_comparisons_externally(
470    cdc_with_options: WithOptionsSecResolved,
471    pk_names: &[String],
472) -> Result<Vec<CdcKeyComparison>> {
473    // Replacement plans also use this path, so a successful schema change persists the current
474    // upstream comparison semantics in the new stream graph.
475    let (options, secret_refs) = cdc_with_options.into_parts();
476    let config = ExternalTableConfig::try_from_btreemap(options, secret_refs)
477        .context("failed to extract external table config")?;
478    Ok(ExternalTableImpl::discover_pk_column_comparisons(&config, pk_names).await?)
479}
480
481/// Derive the schema of a CDC table from explicit SQL columns and constraints.
482pub(crate) fn bind_cdc_table_schema(
483    column_defs: &Vec<ColumnDef>,
484    constraints: &Vec<TableConstraint>,
485    is_for_replace_plan: bool,
486) -> Result<(Vec<ColumnCatalog>, Vec<String>)> {
487    let columns = bind_sql_columns(column_defs, is_for_replace_plan)?;
488    // CDC parsers cannot produce variant values, and this path has no FORMAT/ENCODE gate.
489    reject_variant_columns(&columns, "on a table created from a CDC source")?;
490
491    let pk_names = bind_sql_pk_names(column_defs, bind_table_constraints(constraints)?)?;
492    Ok((columns, pk_names))
493}
494
495#[cfg(test)]
496mod tests {
497    use super::*;
498
499    fn test_schema_table_name() -> SchemaTableName {
500        SchemaTableName {
501            schema_name: "public".to_owned(),
502            table_name: "orders".to_owned(),
503        }
504    }
505
506    fn pk_names() -> Vec<String> {
507        vec!["plan_id".to_owned(), "site_id".to_owned()]
508    }
509
510    #[test]
511    fn test_debezium_filter_rejects_literal_excluded_pk() {
512        let err = reject_pk_filtered_by_debezium_column_filter_inner(
513            &pk_names(),
514            &test_schema_table_name(),
515            Some("public.orders.site_id"),
516            None,
517        )
518        .unwrap_err();
519
520        assert!(err.to_report_string().contains("site_id"));
521        assert!(
522            err.to_report_string()
523                .contains("debezium.column.exclude.list")
524        );
525    }
526
527    #[test]
528    fn test_debezium_filter_rejects_include_list_missing_pk() {
529        let err = reject_pk_filtered_by_debezium_column_filter_inner(
530            &pk_names(),
531            &test_schema_table_name(),
532            None,
533            Some("public.orders.plan_id,public.orders.payload"),
534        )
535        .unwrap_err();
536
537        assert!(err.to_report_string().contains("site_id"));
538        assert!(
539            err.to_report_string()
540                .contains("debezium.column.include.list")
541        );
542    }
543
544    #[test]
545    fn test_debezium_filter_accepts_include_list_covering_all_pks() {
546        reject_pk_filtered_by_debezium_column_filter_inner(
547            &pk_names(),
548            &test_schema_table_name(),
549            None,
550            Some("public.orders.plan_id,public.orders.site_id,public.orders.payload"),
551        )
552        .unwrap();
553    }
554
555    #[test]
556    fn test_debezium_filter_matches_regex_patterns() {
557        reject_pk_filtered_by_debezium_column_filter_inner(
558            &pk_names(),
559            &test_schema_table_name(),
560            None,
561            Some(r"public[.]orders[.](plan_id|site_id),public[.]orders[.]payload"),
562        )
563        .unwrap();
564
565        let err = reject_pk_filtered_by_debezium_column_filter_inner(
566            &pk_names(),
567            &test_schema_table_name(),
568            Some(r".*[.]orders[.]site_id"),
569            None,
570        )
571        .unwrap_err();
572
573        assert!(err.to_report_string().contains("site_id"));
574    }
575
576    #[test]
577    fn test_debezium_filter_matches_patterns_case_insensitively() {
578        let err = reject_pk_filtered_by_debezium_column_filter_inner(
579            &pk_names(),
580            &test_schema_table_name(),
581            Some("PUBLIC.ORDERS.SITE_ID"),
582            None,
583        )
584        .unwrap_err();
585
586        assert!(err.to_report_string().contains("site_id"));
587
588        reject_pk_filtered_by_debezium_column_filter_inner(
589            &pk_names(),
590            &test_schema_table_name(),
591            None,
592            Some("Public.Orders.Plan_ID,Public.Orders.Site_ID"),
593        )
594        .unwrap();
595    }
596
597    #[test]
598    fn test_parse_postgres_cdc_external_table_name() {
599        for (input, expected) in [
600            ("public.Note", ("public", "Note")),
601            ("public.\"Note\"", ("public", "Note")),
602            (
603                "\"Mixed.Schema\".\"Note.Table\"",
604                ("Mixed.Schema", "Note.Table"),
605            ),
606            ("public.\"Note\"\"Archive\"", ("public", "Note\"Archive")),
607        ] {
608            assert_eq!(
609                parse_postgres_cdc_external_table_name(input).unwrap(),
610                (expected.0.to_owned(), expected.1.to_owned()),
611                "input: {input}"
612            );
613        }
614
615        for input in [
616            "Note",
617            "public.",
618            ".Note",
619            "public.\"Note",
620            "public.\"Note\"Archive",
621            "public.Note.Archive",
622        ] {
623            assert!(
624                parse_postgres_cdc_external_table_name(input).is_err(),
625                "input should be rejected: {input}"
626            );
627        }
628    }
629}