From 7a67574d136c84635a5e031afe59a1bb4d98ff26 Mon Sep 17 00:00:00 2001 From: andreas-neumann_data Date: Wed, 22 Jul 2026 18:14:03 +0000 Subject: [PATCH 1/7] [SPARK-58261][SQL][SDP] Support SCD Type 2 syntax in SQL AUTO CDC ### What changes were proposed in this pull request? Add SQL syntax for selecting the SCD type and history-tracking columns of an AUTO CDC flow, mirroring the recently added Python/Connect surface. Both SQL AUTO CDC forms -- `CREATE STREAMING TABLE ... FLOW AUTO CDC ...` and `CREATE FLOW ... AS AUTO CDC INTO ...` -- now accept: - `STORED AS SCD TYPE 1 | 2` (defaults to SCD Type 1 when omitted), and - `TRACK HISTORY ON ()` / `TRACK HISTORY ON * EXCEPT ()` under SCD Type 2. Specifically: - Grammar: add multi-word lexer tokens `SCD_TYPE_1`/`SCD_TYPE_2` and `HISTORY`/`TRACK` keywords; add `autoCdcStoredAsClause` and `autoCdcTrackHistoryClause`, appended in fixed order after the existing `COLUMNS` clause. `HISTORY`/`TRACK` are registered as non-reserved. - AST/plans: thread `storedAsScdType` (Int, default 1) and the track-history column lists through `AutoCdcParams`, `AutoCdcInto`, and `CreateStreamingTableAutoCdc`. - Registration: `SqlGraphRegistrationContext.buildChangeArgs` maps the SCD type onto `ScdType`, builds `trackHistorySelection`, and rejects TRACK HISTORY under SCD1 as well as specifying both track-history lists. The SCD type and track-history columns feed the existing `ChangeArgs` model (`storedAsScdType` / `trackHistorySelection`), the same contract the engine's SCD2 execution work consumes. ### Why are the changes needed? SQL AUTO CDC previously only supported SCD Type 1 and rejected `STORED AS SCD TYPE 2` / `TRACK HISTORY` at parse time. SCD Type 2 (history tracking) is a common CDC requirement and is already modeled in the engine's `ChangeArgs` / `ScdType`, so the SQL surface was the missing piece. ### Does this PR introduce any user-facing change? Yes. SQL AUTO CDC statements accept `STORED AS SCD TYPE 1|2` and, under SCD Type 2, `TRACK HISTORY ON ...`. Previously these clauses failed to parse. This is a change within the unreleased master branch only. ### How was this patch tested? New parser tests in `AutoCdcParserSuite` (SCD type 1/2, TRACK HISTORY include/except, all clauses combined, clause-ordering and invalid-type negative cases) and registration tests in `SqlPipelineSuite` (SCD2 and track-history map onto `ChangeArgs`; TRACK HISTORY without SCD2 is rejected). Ran `AutoCdcParserSuite`, `SQLKeywordSuite`, and the AUTO CDC `SqlPipelineSuite` cases. Co-authored-by: Isaac --- docs/sql-ref-ansi-compliance.md | 2 + .../spark/sql/catalyst/parser/SqlBaseLexer.g4 | 4 + .../sql/catalyst/parser/SqlBaseParser.g4 | 20 ++ .../sql/catalyst/parser/AstBuilder.scala | 31 ++- .../catalyst/plans/logical/AutoCdcInto.scala | 18 +- .../catalyst/plans/logical/v2Commands.scala | 15 +- .../analysis/UnsupportedOperationsSuite.scala | 5 +- .../spark/sql/execution/SparkSqlParser.scala | 5 +- .../command/v2/AutoCdcParserSuite.scala | 182 ++++++++++++++---- .../graph/SqlGraphRegistrationContext.scala | 58 +++++- .../pipelines/graph/SqlPipelineSuite.scala | 77 +++++++- 11 files changed, 370 insertions(+), 47 deletions(-) diff --git a/docs/sql-ref-ansi-compliance.md b/docs/sql-ref-ansi-compliance.md index 367abce6cb8d3..cb2acf5ee29a3 100644 --- a/docs/sql-ref-ansi-compliance.md +++ b/docs/sql-ref-ansi-compliance.md @@ -583,6 +583,7 @@ Below is a list of all the keywords in Spark SQL. |GROUPING|non-reserved|non-reserved|reserved| |HANDLER|non-reserved|non-reserved|non-reserved| |HAVING|reserved|non-reserved|reserved| +|HISTORY|non-reserved|non-reserved|non-reserved| |HOUR|non-reserved|non-reserved|non-reserved| |HOURS|non-reserved|non-reserved|non-reserved| |IDENTIFIER|non-reserved|non-reserved|non-reserved| @@ -803,6 +804,7 @@ Below is a list of all the keywords in Spark SQL. |TINYINT|non-reserved|non-reserved|non-reserved| |TO|reserved|non-reserved|reserved| |TOUCH|non-reserved|non-reserved|non-reserved| +|TRACK|non-reserved|non-reserved|non-reserved| |TRAILING|reserved|non-reserved|reserved| |TRANSACTION|non-reserved|non-reserved|non-reserved| |TRANSACTIONS|non-reserved|non-reserved|non-reserved| diff --git a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 index 64a99e82053b1..dbd45012350ba 100644 --- a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 +++ b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 @@ -300,6 +300,7 @@ GROUPING: 'GROUPING'; HANDLER: 'HANDLER'; HAVING: 'HAVING'; BINARY_HEX: 'X'; +HISTORY: 'HISTORY'; HOUR: 'HOUR'; HOURS: 'HOURS'; IDENTIFIER_KW: 'IDENTIFIER'; @@ -460,6 +461,8 @@ SECOND: 'SECOND'; SECONDS: 'SECONDS'; SCHEMA: 'SCHEMA'; SCHEMAS: 'SCHEMAS'; +SCD_TYPE_1: 'SCD TYPE 1'; +SCD_TYPE_2: 'SCD TYPE 2'; SECURITY: 'SECURITY'; SELECT: 'SELECT'; SEMI: 'SEMI'; @@ -519,6 +522,7 @@ TINYINT: 'TINYINT'; TO: 'TO'; EXECUTE: 'EXECUTE'; TOUCH: 'TOUCH'; +TRACK: 'TRACK'; TRAILING: 'TRAILING'; TRANSACTION: 'TRANSACTION'; TRANSACTIONS: 'TRANSACTIONS'; diff --git a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 index 29094679e0fb6..8ba223b29f1ed 100644 --- a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 +++ b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 @@ -774,6 +774,8 @@ autoCdcParameters autoCdcDeleteClause? autoCdcSequenceByClause autoCdcColumnsClause? + autoCdcStoredAsClause? + autoCdcTrackHistoryClause? ; autoCdcDeleteClause @@ -790,6 +792,16 @@ autoCdcColumnsClause ASTERISK EXCEPT LEFT_PAREN exceptCols=identifierSeq RIGHT_PAREN) ; +autoCdcStoredAsClause + : STORED AS (SCD_TYPE_1 | SCD_TYPE_2) + ; + +autoCdcTrackHistoryClause + : TRACK HISTORY ON ( + LEFT_PAREN trackCols=identifierSeq RIGHT_PAREN | + ASTERISK EXCEPT LEFT_PAREN nonTrackCols=identifierSeq RIGHT_PAREN) + ; + identifierReference : IDENTIFIER_KW LEFT_PAREN expression RIGHT_PAREN | multipartIdentifier @@ -2160,6 +2172,7 @@ ansiNonReserved | GLOBAL | GROUPING | HANDLER + | HISTORY | HOUR | HOURS | IDENTIFIER_KW @@ -2294,6 +2307,8 @@ ansiNonReserved | ROWS | SCHEMA | SCHEMAS + | SCD_TYPE_1 + | SCD_TYPE_2 | SECOND | SECONDS | SECURITY @@ -2346,6 +2361,7 @@ ansiNonReserved | TIMESTAMPDIFF | TINYINT | TOUCH + | TRACK | TRANSACTION | TRANSACTIONS | TRANSFORM @@ -2586,6 +2602,7 @@ nonReserved | GROUPING | HANDLER | HAVING + | HISTORY | HOUR | HOURS | IDENTIFIER_KW @@ -2737,6 +2754,8 @@ nonReserved | ROWS | SCHEMA | SCHEMAS + | SCD_TYPE_1 + | SCD_TYPE_2 | SECOND | SECONDS | SECURITY @@ -2795,6 +2814,7 @@ nonReserved | TINYINT | TO | TOUCH + | TRACK | TRAILING | TRANSACTION | TRANSACTIONS diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala index 529aad74956a1..994a6df400863 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala @@ -1400,7 +1400,10 @@ class AstBuilder extends DataTypeAstBuilder deleteCondition = params.deleteCondition, sequenceByExpr = params.sequencing, includeColumns = params.includeColumns, - excludeColumns = params.excludeColumns) + excludeColumns = params.excludeColumns, + storedAsScdType = params.storedAsScdType, + trackHistoryColumns = params.trackHistoryColumns, + trackHistoryExceptColumns = params.trackHistoryExceptColumns) } protected def parseAutoCdcParams(params: AutoCdcParametersContext): AutoCdcParams = @@ -1430,13 +1433,32 @@ class AstBuilder extends DataTypeAstBuilder visitIdentifierSeq(c.exceptCols).map(UnresolvedAttribute.quoted) } + // STORED AS SCD TYPE 1|2. Absent clause defaults to SCD Type 1. + val storedAsScdType = Option(params.autoCdcStoredAsClause()) match { + case Some(c) if c.SCD_TYPE_2() != null => 2 + case _ => 1 + } + + val trackHistoryClause = Option(params.autoCdcTrackHistoryClause()) + val trackHistoryColumns = trackHistoryClause.collect { + case c if c.trackCols != null => + visitIdentifierSeq(c.trackCols).map(UnresolvedAttribute.quoted) + } + val trackHistoryExceptColumns = trackHistoryClause.collect { + case c if c.nonTrackCols != null => + visitIdentifierSeq(c.nonTrackCols).map(UnresolvedAttribute.quoted) + } + AutoCdcParams( source = source, keys = keys, deleteCondition = deleteCondition, sequencing = sequencing, includeColumns = includeColumns, - excludeColumns = excludeColumns) + excludeColumns = excludeColumns, + storedAsScdType = storedAsScdType, + trackHistoryColumns = trackHistoryColumns, + trackHistoryExceptColumns = trackHistoryExceptColumns) } /** @@ -8038,4 +8060,7 @@ case class AutoCdcParams( deleteCondition: Option[Expression], sequencing: Expression, includeColumns: Option[Seq[UnresolvedAttribute]], - excludeColumns: Option[Seq[UnresolvedAttribute]]) + excludeColumns: Option[Seq[UnresolvedAttribute]], + storedAsScdType: Int, + trackHistoryColumns: Option[Seq[UnresolvedAttribute]], + trackHistoryExceptColumns: Option[Seq[UnresolvedAttribute]]) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/AutoCdcInto.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/AutoCdcInto.scala index cfea98579c630..d78291b7dea3e 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/AutoCdcInto.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/AutoCdcInto.scala @@ -25,7 +25,8 @@ import org.apache.spark.sql.catalyst.trees.BinaryLike * Logical plan node for an AUTO CDC INTO operation, used by Spark Declarative Pipelines. * * This represents a CDC (Change Data Capture) operation that applies an ordered change event - * stream from [[source]] into [[targetTable]] using SCD Type 1 (upsert) semantics. + * stream from [[source]] into [[targetTable]] using SCD Type 1 (upsert) or SCD Type 2 + * (history-tracking) semantics, as selected by [[storedAsScdType]]. * * This node only ever appears as the flow operation of a [[CreateFlowCommand]] (parsed from * `CREATE FLOW ... AS AUTO CDC INTO ...`); there is no standalone `AUTO CDC INTO` syntax. It is a @@ -53,6 +54,16 @@ import org.apache.spark.sql.catalyst.trees.BinaryLike * @param excludeColumns Source columns to exclude from the target table (i.e., all columns * except these). [[None]] when no COLUMNS clause was specified. Mutually * exclusive with [[includeColumns]]. + * @param storedAsScdType The SCD type of the target table, from `STORED AS SCD TYPE `. 1 for + * SCD Type 1 (upsert) and 2 for SCD Type 2 (history-tracking). Defaults to + * 1 when no `STORED AS` clause is specified. + * @param trackHistoryColumns SCD2-only. An explicit list of columns whose value change opens a + * new history record, from `TRACK HISTORY ON (...)`. [[None]] when no TRACK + * HISTORY clause was specified. Mutually exclusive with + * [[trackHistoryExceptColumns]]. + * @param trackHistoryExceptColumns SCD2-only. Columns excluded from history tracking, from + * `TRACK HISTORY ON * EXCEPT (...)`. [[None]] when no TRACK HISTORY clause + * was specified. Mutually exclusive with [[trackHistoryColumns]]. */ case class AutoCdcInto( targetTable: LogicalPlan, @@ -61,7 +72,10 @@ case class AutoCdcInto( deleteCondition: Option[Expression], sequenceByExpr: Expression, includeColumns: Option[Seq[UnresolvedAttribute]], - excludeColumns: Option[Seq[UnresolvedAttribute]] + excludeColumns: Option[Seq[UnresolvedAttribute]], + storedAsScdType: Int, + trackHistoryColumns: Option[Seq[UnresolvedAttribute]], + trackHistoryExceptColumns: Option[Seq[UnresolvedAttribute]] ) extends LogicalPlan with BinaryLike[LogicalPlan] { override def left: LogicalPlan = targetTable override def right: LogicalPlan = source diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala index 468b43d4915fc..e82a65c8f99e7 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/v2Commands.scala @@ -852,6 +852,16 @@ case class CreateStreamingTable( * clause was specified. Mutually exclusive with [[excludeColumns]]. * @param excludeColumns Source columns to exclude. [[None]] when no COLUMNS clause was specified. * Mutually exclusive with [[includeColumns]]. + * @param storedAsScdType The SCD type of the target table, from `STORED AS SCD TYPE `. 1 for + * SCD Type 1 (upsert) and 2 for SCD Type 2 (history-tracking). Defaults to + * 1 when no `STORED AS` clause is specified. + * @param trackHistoryColumns SCD2-only. An explicit list of columns whose value change opens a + * new history record, from `TRACK HISTORY ON (...)`. [[None]] when no TRACK + * HISTORY clause was specified. Mutually exclusive with + * [[trackHistoryExceptColumns]]. + * @param trackHistoryExceptColumns SCD2-only. Columns excluded from history tracking, from + * `TRACK HISTORY ON * EXCEPT (...)`. [[None]] when no TRACK HISTORY clause + * was specified. Mutually exclusive with [[trackHistoryColumns]]. */ case class CreateStreamingTableAutoCdc( name: LogicalPlan, @@ -864,7 +874,10 @@ case class CreateStreamingTableAutoCdc( deleteCondition: Option[Expression], sequenceByExpr: Expression, includeColumns: Option[Seq[UnresolvedAttribute]], - excludeColumns: Option[Seq[UnresolvedAttribute]] + excludeColumns: Option[Seq[UnresolvedAttribute]], + storedAsScdType: Int, + trackHistoryColumns: Option[Seq[UnresolvedAttribute]], + trackHistoryExceptColumns: Option[Seq[UnresolvedAttribute]] ) extends BinaryCommand with CreatePipelineDataset { override def left: LogicalPlan = name override def right: LogicalPlan = source diff --git a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationsSuite.scala b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationsSuite.scala index 785ca2f6f5978..293523b86f998 100644 --- a/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationsSuite.scala +++ b/sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationsSuite.scala @@ -103,7 +103,10 @@ class UnsupportedOperationsSuite extends SparkFunSuite with SQLHelper { deleteCondition = None, sequenceByExpr = attribute, includeColumns = None, - excludeColumns = None)) + excludeColumns = None, + storedAsScdType = 1, + trackHistoryColumns = None, + trackHistoryExceptColumns = None)) /* ======================================================================================= diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala index 57fadd3ecb29e..9acd0ba6def57 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/SparkSqlParser.scala @@ -1747,7 +1747,10 @@ class SparkSqlAstBuilder extends AstBuilder { deleteCondition = params.deleteCondition, sequenceByExpr = params.sequencing, includeColumns = params.includeColumns, - excludeColumns = params.excludeColumns + excludeColumns = params.excludeColumns, + storedAsScdType = params.storedAsScdType, + trackHistoryColumns = params.trackHistoryColumns, + trackHistoryExceptColumns = params.trackHistoryExceptColumns ) } else { Option(ctx.query) match { diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala index c9d28edaf6a6f..32b6a749f0611 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala @@ -39,9 +39,10 @@ import org.apache.spark.sql.types.IntegerType * 1. CREATE FLOW [COMMENT ...] AS AUTO CDC INTO ... * 2. CREATE STREAMING TABLE FLOW AUTO CDC ... * - * Snapshot CDC, SCD Type 2, IGNORE NULL UPDATES, and APPLY AS TRUNCATE WHEN are not - * supported and should fail to parse. The standalone AUTO CDC INTO form (without CREATE FLOW - * or CREATE STREAMING TABLE) is also not supported. + * Both forms support `STORED AS SCD TYPE 1|2` and, under SCD Type 2, `TRACK HISTORY ON ...`. + * Snapshot CDC, IGNORE NULL UPDATES, and APPLY AS TRUNCATE WHEN are not supported and should + * fail to parse. The standalone AUTO CDC INTO form (without CREATE FLOW or CREATE STREAMING + * TABLE) is also not supported. */ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { protected lazy val parser = new SparkSqlParser() @@ -80,6 +81,10 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { assert(cdc.sequenceByExpr == UnresolvedAttribute("timestamp")) assert(cdc.includeColumns.isEmpty) assert(cdc.excludeColumns.isEmpty) + // No STORED AS clause defaults to SCD Type 1 with no history tracking. + assert(cdc.storedAsScdType == 1) + assert(cdc.trackHistoryColumns.isEmpty) + assert(cdc.trackHistoryExceptColumns.isEmpty) } test("CREATE FLOW AS AUTO CDC INTO - multipart source name") { @@ -196,6 +201,135 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { assert(cdc.includeColumns.get.map(_.name) == Seq("key1", "key2", "key3", "timestamp")) } + test("CREATE FLOW AS AUTO CDC INTO - STORED AS SCD TYPE 1") { + val plan = parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 1""".stripMargin) + + val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto] + assert(cdc.storedAsScdType == 1) + assert(cdc.trackHistoryColumns.isEmpty) + assert(cdc.trackHistoryExceptColumns.isEmpty) + } + + test("CREATE FLOW AS AUTO CDC INTO - STORED AS SCD TYPE 2") { + val plan = parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 2""".stripMargin) + + val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto] + assert(cdc.storedAsScdType == 2) + assert(cdc.trackHistoryColumns.isEmpty) + assert(cdc.trackHistoryExceptColumns.isEmpty) + } + + test("CREATE FLOW AS AUTO CDC INTO - STORED AS SCD TYPE 2 with TRACK HISTORY ON") { + val plan = parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON (val1, val2)""".stripMargin) + + val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto] + assert(cdc.storedAsScdType == 2) + assert(cdc.trackHistoryColumns.get.map(_.name) == Seq("val1", "val2")) + assert(cdc.trackHistoryExceptColumns.isEmpty) + } + + test("CREATE FLOW AS AUTO CDC INTO - STORED AS SCD TYPE 2 with TRACK HISTORY ON * EXCEPT") { + val plan = parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON * EXCEPT (op, ts)""".stripMargin) + + val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto] + assert(cdc.storedAsScdType == 2) + assert(cdc.trackHistoryColumns.isEmpty) + assert(cdc.trackHistoryExceptColumns.get.map(_.name) == Seq("op", "ts")) + } + + test("CREATE FLOW AS AUTO CDC INTO - all clauses combined including SCD2 and TRACK HISTORY") { + val plan = parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (key1, key2) + |APPLY AS DELETE WHEN key3 = 3 + |SEQUENCE BY timestamp + |COLUMNS (key1, key2, key3, timestamp) + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON (key3)""".stripMargin) + + val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto] + assert(cdc.deleteCondition.isDefined) + assert(cdc.includeColumns.get.map(_.name) == Seq("key1", "key2", "key3", "timestamp")) + assert(cdc.storedAsScdType == 2) + assert(cdc.trackHistoryColumns.get.map(_.name) == Seq("key3")) + } + + test("CREATE STREAMING TABLE FLOW AUTO CDC - STORED AS SCD TYPE 2 with TRACK HISTORY ON") { + val plan = parser.parsePlan( + """CREATE STREAMING TABLE st FLOW AUTO CDC + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON (val1, val2)""".stripMargin) + + val cst = plan.asInstanceOf[CreateStreamingTableAutoCdc] + assert(cst.storedAsScdType == 2) + assert(cst.trackHistoryColumns.get.map(_.name) == Seq("val1", "val2")) + assert(cst.trackHistoryExceptColumns.isEmpty) + } + + test("CREATE STREAMING TABLE FLOW AUTO CDC - STORED AS SCD TYPE 2 with TRACK HISTORY EXCEPT") { + val plan = parser.parsePlan( + """CREATE STREAMING TABLE st FLOW AUTO CDC + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON * EXCEPT (op)""".stripMargin) + + val cst = plan.asInstanceOf[CreateStreamingTableAutoCdc] + assert(cst.storedAsScdType == 2) + assert(cst.trackHistoryColumns.isEmpty) + assert(cst.trackHistoryExceptColumns.get.map(_.name) == Seq("op")) + } + + test("AUTO CDC - TRACK HISTORY before STORED AS is not allowed") { + intercept[ParseException] { + parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |TRACK HISTORY ON (val) + |STORED AS SCD TYPE 2""".stripMargin) + } + } + + test("AUTO CDC - STORED AS SCD TYPE 3 is not allowed") { + intercept[ParseException] { + parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 3""".stripMargin) + } + } + // --------------------------------------------------------------------------- // CREATE STREAMING TABLE ... FLOW AUTO CDC // --------------------------------------------------------------------------- @@ -810,35 +944,17 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { ) } - test("STORED AS SCD TYPE 2 is not supported") { - checkError( - intercept[ParseException] { - parser.parsePlan( - """CREATE FLOW f AS AUTO CDC INTO target - |FROM STREAM(source) - |KEYS (id) - |SEQUENCE BY ts - |STORED AS SCD TYPE 2""".stripMargin) - }, - condition = "PARSE_SYNTAX_ERROR", - sqlState = "42601", - parameters = Map("error" -> "'STORED'", "hint" -> "") - ) - } - - test("TRACK HISTORY ON is not supported") { - checkError( - intercept[ParseException] { - parser.parsePlan( - """CREATE FLOW f AS AUTO CDC INTO target - |FROM STREAM(source) - |KEYS (id) - |SEQUENCE BY ts - |TRACK HISTORY ON value1, value2""".stripMargin) - }, - condition = "PARSE_SYNTAX_ERROR", - sqlState = "42601", - parameters = Map("error" -> "'TRACK'", "hint" -> "") - ) + test("TRACK HISTORY ON without parentheses is not allowed") { + // The grammar requires a parenthesized column list or `* EXCEPT (...)`, so a bare column + // list after TRACK HISTORY ON fails to parse. + intercept[ParseException] { + parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON value1, value2""".stripMargin) + } } } diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/SqlGraphRegistrationContext.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/SqlGraphRegistrationContext.scala index 36a1f17209a23..a8dd94040b91c 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/SqlGraphRegistrationContext.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/SqlGraphRegistrationContext.scala @@ -275,9 +275,10 @@ class SqlGraphRegistrationContext( * into the [[ChangeArgs]] consumed by an [[AutoCdcFlow]]. Shared by the two SQL AUTO CDC entry * points: `CREATE STREAMING TABLE ... FLOW AUTO CDC ...` and `CREATE FLOW ... AS AUTO CDC INTO`. * - * SQL AUTO CDC syntax only supports SCD Type 1, so [[ChangeArgs.storedAsScdType]] is always - * [[ScdType.Type1]]. [[includeColumns]] and [[excludeColumns]] are mutually exclusive at the - * grammar level; the guard here is defensive. + * `STORED AS SCD TYPE ` selects [[ScdType]] (defaulting to [[ScdType.Type1]] when absent). + * `TRACK HISTORY ON ...` populates the SCD2-only [[ChangeArgs.trackHistorySelection]] and is + * rejected under SCD1. [[includeColumns]]/[[excludeColumns]] and the two track-history lists are + * each mutually exclusive at the grammar level; the guards here are defensive. */ private def buildChangeArgs( keys: Seq[UnresolvedAttribute], @@ -285,6 +286,9 @@ class SqlGraphRegistrationContext( deleteCondition: Option[Expression], includeColumns: Option[Seq[UnresolvedAttribute]], excludeColumns: Option[Seq[UnresolvedAttribute]], + storedAsScdType: Int, + trackHistoryColumns: Option[Seq[UnresolvedAttribute]], + trackHistoryExceptColumns: Option[Seq[UnresolvedAttribute]], queryOrigin: QueryOrigin): ChangeArgs = { val columnSelection: Option[ColumnSelection] = (includeColumns, excludeColumns) match { case (Some(_), Some(_)) => @@ -300,12 +304,50 @@ class SqlGraphRegistrationContext( None } + val scdType: ScdType = storedAsScdType match { + case 1 => ScdType.Type1 + case 2 => ScdType.Type2 + case other => + throw SqlGraphElementRegistrationException( + msg = s"Unsupported AUTO CDC SCD type: $other. Only SCD TYPE 1 and SCD TYPE 2 are " + + "supported.", + queryOrigin = queryOrigin + ) + } + + // TRACK HISTORY is only meaningful under SCD2; reject it for SCD1 before interpreting the + // clause. The grammar allows the clause regardless of the STORED AS type, so this is where + // the SCD2 requirement is enforced. + if (scdType != ScdType.Type2 && + (trackHistoryColumns.isDefined || trackHistoryExceptColumns.isDefined)) { + throw SqlGraphElementRegistrationException( + msg = "AUTO CDC TRACK HISTORY requires STORED AS SCD TYPE 2.", + queryOrigin = queryOrigin + ) + } + + val trackHistorySelection: Option[ColumnSelection] = + (trackHistoryColumns, trackHistoryExceptColumns) match { + case (Some(_), Some(_)) => + throw SqlGraphElementRegistrationException( + msg = "AUTO CDC cannot specify both TRACK HISTORY ON and TRACK HISTORY ON * EXCEPT.", + queryOrigin = queryOrigin + ) + case (Some(tracked), None) => + Option(ColumnSelection.IncludeColumns(tracked.map(toUnqualifiedColumnName))) + case (None, Some(nonTracked)) => + Option(ColumnSelection.ExcludeColumns(nonTracked.map(toUnqualifiedColumnName))) + case (None, None) => + None + } + ChangeArgs( keys = keys.map(toUnqualifiedColumnName), sequencing = Column(sequenceByExpr), - storedAsScdType = ScdType.Type1, + storedAsScdType = scdType, deleteCondition = deleteCondition.map(Column(_)), - columnSelection = columnSelection + columnSelection = columnSelection, + trackHistorySelection = trackHistorySelection ) } @@ -374,6 +416,9 @@ class SqlGraphRegistrationContext( deleteCondition = cst.deleteCondition, includeColumns = cst.includeColumns, excludeColumns = cst.excludeColumns, + storedAsScdType = cst.storedAsScdType, + trackHistoryColumns = cst.trackHistoryColumns, + trackHistoryExceptColumns = cst.trackHistoryExceptColumns, queryOrigin = queryOrigin ) ) @@ -588,6 +633,9 @@ class SqlGraphRegistrationContext( deleteCondition = a.deleteCondition, includeColumns = a.includeColumns, excludeColumns = a.excludeColumns, + storedAsScdType = a.storedAsScdType, + trackHistoryColumns = a.trackHistoryColumns, + trackHistoryExceptColumns = a.trackHistoryExceptColumns, queryOrigin = queryOrigin ) ) diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala index 3297709174c5a..b751c9ff8f890 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala @@ -1120,7 +1120,7 @@ class SqlPipelineSuite extends PipelineTest with SharedSparkSession { // Two SQL forms register an [[AutoCdcFlow]] into the dataflow graph: // 1. CREATE STREAMING TABLE FLOW AUTO CDC FROM KEYS (...) SEQUENCE BY // 2. CREATE FLOW AS AUTO CDC INTO FROM KEYS (...) SEQUENCE BY - // SQL AUTO CDC only supports SCD Type 1. + // Both forms support STORED AS SCD TYPE 1|2 and, under SCD2, TRACK HISTORY ON .... // =========================================================================== /** Returns the single unresolved [[AutoCdcFlow]] registered for the given flow identifier. */ @@ -1215,6 +1215,81 @@ class SqlPipelineSuite extends PipelineTest with SharedSparkSession { } } + test("AUTO CDC STORED AS SCD TYPE 2 maps onto ChangeArgs") { + val graph = unresolvedDataflowGraphFromSql( + sqlText = s""" + |CREATE STREAMING TABLE st + |FLOW AUTO CDC + |FROM STREAM $externalTable1Ident + |KEYS (id) + |SEQUENCE BY id + |STORED AS SCD TYPE 2 + |""".stripMargin + ) + + val flow = autoCdcFlowFor(graph, "st") + assert(flow.changeArgs.storedAsScdType == ScdType.Type2) + assert(flow.changeArgs.trackHistorySelection.isEmpty) + } + + test("AUTO CDC STORED AS SCD TYPE 2 with TRACK HISTORY ON maps onto ChangeArgs") { + val graph = unresolvedDataflowGraphFromSql( + sqlText = s""" + |CREATE STREAMING TABLE target; + |CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM $externalTable1Ident + |KEYS (id) + |SEQUENCE BY id + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON (id) + |""".stripMargin + ) + + val flow = autoCdcFlowFor(graph, "f") + assert(flow.changeArgs.storedAsScdType == ScdType.Type2) + flow.changeArgs.trackHistorySelection match { + case Some(ColumnSelection.IncludeColumns(cols)) => assert(cols.map(_.name) == Seq("id")) + case other => fail(s"Expected IncludeColumns(id), got $other") + } + } + + test("AUTO CDC STORED AS SCD TYPE 2 with TRACK HISTORY ON * EXCEPT maps onto ChangeArgs") { + val graph = unresolvedDataflowGraphFromSql( + sqlText = s""" + |CREATE STREAMING TABLE st + |FLOW AUTO CDC + |FROM STREAM $externalTable1Ident + |KEYS (id) + |SEQUENCE BY id + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON * EXCEPT (id) + |""".stripMargin + ) + + val flow = autoCdcFlowFor(graph, "st") + assert(flow.changeArgs.storedAsScdType == ScdType.Type2) + flow.changeArgs.trackHistorySelection match { + case Some(ColumnSelection.ExcludeColumns(cols)) => assert(cols.map(_.name) == Seq("id")) + case other => fail(s"Expected ExcludeColumns(id), got $other") + } + } + + test("AUTO CDC TRACK HISTORY without SCD TYPE 2 is rejected") { + val ex = intercept[SqlGraphElementRegistrationException] { + unresolvedDataflowGraphFromSql( + sqlText = s""" + |CREATE STREAMING TABLE st + |FLOW AUTO CDC + |FROM STREAM $externalTable1Ident + |KEYS (id) + |SEQUENCE BY id + |TRACK HISTORY ON (id) + |""".stripMargin + ) + } + assert(ex.getMessage.contains("TRACK HISTORY requires STORED AS SCD TYPE 2")) + } + test("Explicit column list is rejected for CREATE STREAMING TABLE FLOW AUTO CDC") { val ex = intercept[SqlGraphElementRegistrationException] { unresolvedDataflowGraphFromSql( From e1db1d26359b1b1db4218e05f8ae84fd59024daf Mon Sep 17 00:00:00 2001 From: andreas-neumann_data Date: Wed, 22 Jul 2026 21:43:28 +0000 Subject: [PATCH 2/7] Regenerate keyword golden files for HISTORY and TRACK Adding the HISTORY and TRACK non-reserved keywords for the SCD2 AUTO CDC syntax changes the output of SQL_KEYWORDS(), which the keywords.sql / keywords-enforced.sql / nonansi/keywords.sql golden-file tests assert against (run by both SQLQueryTestSuite and ThriftServerQueryTestSuite). Regenerate the .sql.out files via SPARK_GENERATE_GOLDEN_FILES=1. The multi-word SCD_TYPE_1 / SCD_TYPE_2 tokens contain spaces and digits and are excluded by the keyword regex, so only HISTORY and TRACK are added. Co-authored-by: Isaac --- .../test/resources/sql-tests/results/keywords-enforced.sql.out | 2 ++ sql/core/src/test/resources/sql-tests/results/keywords.sql.out | 2 ++ .../test/resources/sql-tests/results/nonansi/keywords.sql.out | 2 ++ 3 files changed, 6 insertions(+) diff --git a/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out b/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out index 989bfb7251117..054df304b9d50 100644 --- a/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out +++ b/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out @@ -176,6 +176,7 @@ GROUP true GROUPING false HANDLER false HAVING true +HISTORY false HOUR false HOURS false IDENTIFIED false @@ -392,6 +393,7 @@ TIMESTAMP_NTZ false TINYINT false TO true TOUCH false +TRACK false TRAILING true TRANSACTION false TRANSACTIONS false diff --git a/sql/core/src/test/resources/sql-tests/results/keywords.sql.out b/sql/core/src/test/resources/sql-tests/results/keywords.sql.out index ee127e584f172..6f0cf9bf40938 100644 --- a/sql/core/src/test/resources/sql-tests/results/keywords.sql.out +++ b/sql/core/src/test/resources/sql-tests/results/keywords.sql.out @@ -176,6 +176,7 @@ GROUP false GROUPING false HANDLER false HAVING false +HISTORY false HOUR false HOURS false IDENTIFIED false @@ -392,6 +393,7 @@ TIMESTAMP_NTZ false TINYINT false TO false TOUCH false +TRACK false TRAILING false TRANSACTION false TRANSACTIONS false diff --git a/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out b/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out index ee127e584f172..6f0cf9bf40938 100644 --- a/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out +++ b/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out @@ -176,6 +176,7 @@ GROUP false GROUPING false HANDLER false HAVING false +HISTORY false HOUR false HOURS false IDENTIFIED false @@ -392,6 +393,7 @@ TIMESTAMP_NTZ false TINYINT false TO false TOUCH false +TRACK false TRAILING false TRANSACTION false TRANSACTIONS false From bf2610f55d1cd1defac8a2b5abf34f856fb8d5d8 Mon Sep 17 00:00:00 2001 From: andreas-neumann_data Date: Thu, 23 Jul 2026 00:01:23 +0000 Subject: [PATCH 3/7] Add HISTORY and TRACK to getSQLKeywords expected list The HISTORY and TRACK non-reserved keywords added for the SCD2 AUTO CDC syntax appear in DatabaseMetaData.getSQLKeywords(), which SparkConnectDatabaseMetaDataSuite asserts against with a hardcoded expected string. Insert both keywords in their sorted positions (after HANDLER and TOUCH respectively). Co-authored-by: Isaac --- .../connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala b/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala index babac5dd94863..2e25c692b98c2 100644 --- a/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala +++ b/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala @@ -210,7 +210,7 @@ class SparkConnectDatabaseMetaDataSuite extends ConnectFunSuite with RemoteSpark val metadata = conn.getMetaData // scalastyle:off line.size.limit // CURRENT_PATH and SYSTEM are excluded: getSQLKeywords drops SQL:2003 reserved words (see companion). - assert(metadata.getSQLKeywords === "ADD,AFTER,AGGREGATE,ALIGN,ALWAYS,ANALYZE,ANTI,ANY_VALUE,APPLY,APPROX,ARCHIVE,ASC,ASOF,AUTO,BERNOULLI,BIN,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BUCKET,BUCKETS,BYTE,CACHE,CASCADE,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CLEAR,CLUSTER,CLUSTERED,CODEGEN,COLLATION,COLLATIONS,COLLECTION,COLUMNS,COMMENT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONTAINS,CONTINUE,COST,CURRENT_DATABASE,CURRENT_SCHEMA,DATA,DATABASE,DATABASES,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAYOFYEAR,DAYS,DBPROPERTIES,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELIMITED,DESC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTRIBUTE,DIV,DO,ELSEIF,ENFORCED,ESCAPED,EVOLUTION,EXACT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,FIELDS,FILEFORMAT,FIRST,FLOW,FOLLOWING,FORMAT,FORMATTED,FOUND,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,HANDLER,HOURS,IDENTIFIED,IDENTIFIER,IF,IGNORE,ILIKE,IMMEDIATE,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INPATH,INPUT,INPUTFORMAT,INVOKER,ITEMS,ITERATE,JSON,KEY,KEYS,LAST,LAZY,LEAVE,LEVEL,LIMIT,LINES,LIST,LOAD,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MEASURE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTES,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NEAREST,NORELY,NULLS,OFFSET,OPTION,OPTIONS,OUTPUTFORMAT,OVERWRITE,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,PRECEDING,PRINCIPALS,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,REDUCE,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,ROLE,ROLES,SCHEMA,SCHEMAS,SECONDS,SECURITY,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SORT,SORTED,SOURCE,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SYNC,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLES,TARGET,TBLPROPERTIES,TERMINATED,TIMEDIFF,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TOUCH,TRANSACTION,TRANSACTIONS,TRANSFORM,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNLOCK,UNPIVOT,UNSET,UNTIL,USE,VAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHILE,WIDTH,X,YEARS,ZONE") + assert(metadata.getSQLKeywords === "ADD,AFTER,AGGREGATE,ALIGN,ALWAYS,ANALYZE,ANTI,ANY_VALUE,APPLY,APPROX,ARCHIVE,ASC,ASOF,AUTO,BERNOULLI,BIN,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BUCKET,BUCKETS,BYTE,CACHE,CASCADE,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CLEAR,CLUSTER,CLUSTERED,CODEGEN,COLLATION,COLLATIONS,COLLECTION,COLUMNS,COMMENT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONTAINS,CONTINUE,COST,CURRENT_DATABASE,CURRENT_SCHEMA,DATA,DATABASE,DATABASES,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAYOFYEAR,DAYS,DBPROPERTIES,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELIMITED,DESC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTRIBUTE,DIV,DO,ELSEIF,ENFORCED,ESCAPED,EVOLUTION,EXACT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,FIELDS,FILEFORMAT,FIRST,FLOW,FOLLOWING,FORMAT,FORMATTED,FOUND,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,HANDLER,HISTORY,HOURS,IDENTIFIED,IDENTIFIER,IF,IGNORE,ILIKE,IMMEDIATE,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INPATH,INPUT,INPUTFORMAT,INVOKER,ITEMS,ITERATE,JSON,KEY,KEYS,LAST,LAZY,LEAVE,LEVEL,LIMIT,LINES,LIST,LOAD,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MEASURE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTES,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NEAREST,NORELY,NULLS,OFFSET,OPTION,OPTIONS,OUTPUTFORMAT,OVERWRITE,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,PRECEDING,PRINCIPALS,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,REDUCE,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,ROLE,ROLES,SCHEMA,SCHEMAS,SECONDS,SECURITY,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SORT,SORTED,SOURCE,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SYNC,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLES,TARGET,TBLPROPERTIES,TERMINATED,TIMEDIFF,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TOUCH,TRACK,TRANSACTION,TRANSACTIONS,TRANSFORM,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNLOCK,UNPIVOT,UNSET,UNTIL,USE,VAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHILE,WIDTH,X,YEARS,ZONE") // scalastyle:on line.size.limit } } From c56c6f97dabf73e1442e74fe79057e9dce82f330 Mon Sep 17 00:00:00 2001 From: andreas-neumann_data Date: Thu, 23 Jul 2026 00:37:18 +0000 Subject: [PATCH 4/7] Lex STORED AS SCD TYPE n as separate tokens instead of multi-word literals Address review: replace the multi-word lexer literals SCD_TYPE_1 ('SCD TYPE 1') and SCD_TYPE_2 ('SCD TYPE 2') with a single SCD keyword token, and parse the clause as `STORED AS SCD TYPE type=INTEGER_VALUE`, validating the number (1 or 2) in the AstBuilder with a clear error. Benefits: - Fixes the whitespace brittleness the multi-word literals imposed (SPARK-58270): arbitrary whitespace between SCD, TYPE and the number is now accepted. - Removes the SCD_TYPE_1/SCD_TYPE_2 entries from the non-reserved lists, which had no effect (multi-word literals can never stand in as identifiers, and the digit in the token name excluded them from SQLKeywordSuite / docs / golden files). - SCD TYPE 3 (and other unsupported numbers) now fail with a clear "Unsupported SCD type" message rather than a generic parse error. SCD is a plain [A-Z_]+ keyword, so it flows into the keyword lists: added to the docs table, keywords*.sql.out golden files (regenerated), and the JDBC getSQLKeywords expected string. `scd` remains usable as an identifier (non-reserved). Co-authored-by: Isaac --- docs/sql-ref-ansi-compliance.md | 1 + .../spark/sql/catalyst/parser/SqlBaseLexer.g4 | 3 +-- .../spark/sql/catalyst/parser/SqlBaseParser.g4 | 8 +++----- .../spark/sql/catalyst/parser/AstBuilder.scala | 14 +++++++++++--- .../SparkConnectDatabaseMetaDataSuite.scala | 2 +- .../sql-tests/results/keywords-enforced.sql.out | 1 + .../sql-tests/results/keywords.sql.out | 1 + .../sql-tests/results/nonansi/keywords.sql.out | 1 + .../command/v2/AutoCdcParserSuite.scala | 17 ++++++++++++++++- 9 files changed, 36 insertions(+), 12 deletions(-) diff --git a/docs/sql-ref-ansi-compliance.md b/docs/sql-ref-ansi-compliance.md index cb2acf5ee29a3..f4ac58a4e3f02 100644 --- a/docs/sql-ref-ansi-compliance.md +++ b/docs/sql-ref-ansi-compliance.md @@ -742,6 +742,7 @@ Below is a list of all the keywords in Spark SQL. |ROLLUP|non-reserved|non-reserved|reserved| |ROW|non-reserved|non-reserved|reserved| |ROWS|non-reserved|non-reserved|reserved| +|SCD|non-reserved|non-reserved|non-reserved| |SCHEMA|non-reserved|non-reserved|non-reserved| |SCHEMAS|non-reserved|non-reserved|non-reserved| |SECOND|non-reserved|non-reserved|non-reserved| diff --git a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 index dbd45012350ba..1445e4015beb6 100644 --- a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 +++ b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseLexer.g4 @@ -459,10 +459,9 @@ ROW: 'ROW'; ROWS: 'ROWS'; SECOND: 'SECOND'; SECONDS: 'SECONDS'; +SCD: 'SCD'; SCHEMA: 'SCHEMA'; SCHEMAS: 'SCHEMAS'; -SCD_TYPE_1: 'SCD TYPE 1'; -SCD_TYPE_2: 'SCD TYPE 2'; SECURITY: 'SECURITY'; SELECT: 'SELECT'; SEMI: 'SEMI'; diff --git a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 index 8ba223b29f1ed..441011f9bb80c 100644 --- a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 +++ b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 @@ -793,7 +793,7 @@ autoCdcColumnsClause ; autoCdcStoredAsClause - : STORED AS (SCD_TYPE_1 | SCD_TYPE_2) + : STORED AS SCD TYPE type=INTEGER_VALUE ; autoCdcTrackHistoryClause @@ -2305,10 +2305,9 @@ ansiNonReserved | ROLLUP | ROW | ROWS + | SCD | SCHEMA | SCHEMAS - | SCD_TYPE_1 - | SCD_TYPE_2 | SECOND | SECONDS | SECURITY @@ -2752,10 +2751,9 @@ nonReserved | ROLLUP | ROW | ROWS + | SCD | SCHEMA | SCHEMAS - | SCD_TYPE_1 - | SCD_TYPE_2 | SECOND | SECONDS | SECURITY diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala index 994a6df400863..2e77eae904994 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala @@ -1433,10 +1433,18 @@ class AstBuilder extends DataTypeAstBuilder visitIdentifierSeq(c.exceptCols).map(UnresolvedAttribute.quoted) } - // STORED AS SCD TYPE 1|2. Absent clause defaults to SCD Type 1. + // STORED AS SCD TYPE . Absent clause defaults to SCD Type 1. Only 1 and 2 are + // supported; reject anything else with a clear error rather than a generic parse failure. val storedAsScdType = Option(params.autoCdcStoredAsClause()) match { - case Some(c) if c.SCD_TYPE_2() != null => 2 - case _ => 1 + case Some(c) => + val scdType = c.`type`.getText.toInt + if (scdType != 1 && scdType != 2) { + operationNotAllowed( + s"Unsupported SCD type: $scdType. AUTO CDC only supports STORED AS SCD TYPE 1 " + + "or SCD TYPE 2.", c) + } + scdType + case None => 1 } val trackHistoryClause = Option(params.autoCdcTrackHistoryClause()) diff --git a/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala b/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala index 2e25c692b98c2..3d3260c750f6f 100644 --- a/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala +++ b/sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala @@ -210,7 +210,7 @@ class SparkConnectDatabaseMetaDataSuite extends ConnectFunSuite with RemoteSpark val metadata = conn.getMetaData // scalastyle:off line.size.limit // CURRENT_PATH and SYSTEM are excluded: getSQLKeywords drops SQL:2003 reserved words (see companion). - assert(metadata.getSQLKeywords === "ADD,AFTER,AGGREGATE,ALIGN,ALWAYS,ANALYZE,ANTI,ANY_VALUE,APPLY,APPROX,ARCHIVE,ASC,ASOF,AUTO,BERNOULLI,BIN,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BUCKET,BUCKETS,BYTE,CACHE,CASCADE,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CLEAR,CLUSTER,CLUSTERED,CODEGEN,COLLATION,COLLATIONS,COLLECTION,COLUMNS,COMMENT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONTAINS,CONTINUE,COST,CURRENT_DATABASE,CURRENT_SCHEMA,DATA,DATABASE,DATABASES,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAYOFYEAR,DAYS,DBPROPERTIES,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELIMITED,DESC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTRIBUTE,DIV,DO,ELSEIF,ENFORCED,ESCAPED,EVOLUTION,EXACT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,FIELDS,FILEFORMAT,FIRST,FLOW,FOLLOWING,FORMAT,FORMATTED,FOUND,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,HANDLER,HISTORY,HOURS,IDENTIFIED,IDENTIFIER,IF,IGNORE,ILIKE,IMMEDIATE,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INPATH,INPUT,INPUTFORMAT,INVOKER,ITEMS,ITERATE,JSON,KEY,KEYS,LAST,LAZY,LEAVE,LEVEL,LIMIT,LINES,LIST,LOAD,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MEASURE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTES,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NEAREST,NORELY,NULLS,OFFSET,OPTION,OPTIONS,OUTPUTFORMAT,OVERWRITE,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,PRECEDING,PRINCIPALS,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,REDUCE,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,ROLE,ROLES,SCHEMA,SCHEMAS,SECONDS,SECURITY,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SORT,SORTED,SOURCE,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SYNC,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLES,TARGET,TBLPROPERTIES,TERMINATED,TIMEDIFF,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TOUCH,TRACK,TRANSACTION,TRANSACTIONS,TRANSFORM,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNLOCK,UNPIVOT,UNSET,UNTIL,USE,VAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHILE,WIDTH,X,YEARS,ZONE") + assert(metadata.getSQLKeywords === "ADD,AFTER,AGGREGATE,ALIGN,ALWAYS,ANALYZE,ANTI,ANY_VALUE,APPLY,APPROX,ARCHIVE,ASC,ASOF,AUTO,BERNOULLI,BIN,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BUCKET,BUCKETS,BYTE,CACHE,CASCADE,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CLEAR,CLUSTER,CLUSTERED,CODEGEN,COLLATION,COLLATIONS,COLLECTION,COLUMNS,COMMENT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONTAINS,CONTINUE,COST,CURRENT_DATABASE,CURRENT_SCHEMA,DATA,DATABASE,DATABASES,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAYOFYEAR,DAYS,DBPROPERTIES,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELIMITED,DESC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTRIBUTE,DIV,DO,ELSEIF,ENFORCED,ESCAPED,EVOLUTION,EXACT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,FIELDS,FILEFORMAT,FIRST,FLOW,FOLLOWING,FORMAT,FORMATTED,FOUND,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,HANDLER,HISTORY,HOURS,IDENTIFIED,IDENTIFIER,IF,IGNORE,ILIKE,IMMEDIATE,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INPATH,INPUT,INPUTFORMAT,INVOKER,ITEMS,ITERATE,JSON,KEY,KEYS,LAST,LAZY,LEAVE,LEVEL,LIMIT,LINES,LIST,LOAD,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MEASURE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTES,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NEAREST,NORELY,NULLS,OFFSET,OPTION,OPTIONS,OUTPUTFORMAT,OVERWRITE,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,PRECEDING,PRINCIPALS,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,REDUCE,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,ROLE,ROLES,SCD,SCHEMA,SCHEMAS,SECONDS,SECURITY,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SORT,SORTED,SOURCE,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SYNC,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLES,TARGET,TBLPROPERTIES,TERMINATED,TIMEDIFF,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TOUCH,TRACK,TRANSACTION,TRANSACTIONS,TRANSFORM,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNLOCK,UNPIVOT,UNSET,UNTIL,USE,VAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHILE,WIDTH,X,YEARS,ZONE") // scalastyle:on line.size.limit } } diff --git a/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out b/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out index 054df304b9d50..7a6608885e3c0 100644 --- a/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out +++ b/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out @@ -333,6 +333,7 @@ ROLLBACK false ROLLUP false ROW false ROWS false +SCD false SCHEMA false SCHEMAS false SECOND false diff --git a/sql/core/src/test/resources/sql-tests/results/keywords.sql.out b/sql/core/src/test/resources/sql-tests/results/keywords.sql.out index 6f0cf9bf40938..db8bfa0e205c9 100644 --- a/sql/core/src/test/resources/sql-tests/results/keywords.sql.out +++ b/sql/core/src/test/resources/sql-tests/results/keywords.sql.out @@ -333,6 +333,7 @@ ROLLBACK false ROLLUP false ROW false ROWS false +SCD false SCHEMA false SCHEMAS false SECOND false diff --git a/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out b/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out index 6f0cf9bf40938..db8bfa0e205c9 100644 --- a/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out +++ b/sql/core/src/test/resources/sql-tests/results/nonansi/keywords.sql.out @@ -333,6 +333,7 @@ ROLLBACK false ROLLUP false ROW false ROWS false +SCD false SCHEMA false SCHEMAS false SECOND false diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala index 32b6a749f0611..ace2f688ca157 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala @@ -320,7 +320,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { } test("AUTO CDC - STORED AS SCD TYPE 3 is not allowed") { - intercept[ParseException] { + val ex = intercept[ParseException] { parser.parsePlan( """CREATE FLOW f AS AUTO CDC INTO target |FROM STREAM(source) @@ -328,6 +328,21 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { |SEQUENCE BY ts |STORED AS SCD TYPE 3""".stripMargin) } + assert(ex.getMessage.contains("Unsupported SCD type: 3")) + } + + test("AUTO CDC - STORED AS SCD TYPE tolerates extra whitespace between words") { + // SCD, TYPE and the number are separate tokens, so any whitespace between them is accepted. + val plan = parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE + | 2""".stripMargin) + + val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto] + assert(cdc.storedAsScdType == 2) } // --------------------------------------------------------------------------- From 65f0dbb88994621980853db1e696a08280535815 Mon Sep 17 00:00:00 2001 From: andreas-neumann_data Date: Thu, 23 Jul 2026 00:45:15 +0000 Subject: [PATCH 5/7] address review: test explicit SCD1 + TRACK HISTORY, and keywords as identifiers Two test additions suggested in review: - AUTO CDC with an explicit STORED AS SCD TYPE 1 plus TRACK HISTORY is rejected (SqlPipelineSuite); previously only the implicit-default-SCD1 case was covered. - The new non-reserved keywords history, track and scd still parse as ordinary table/column identifiers (AutoCdcParserSuite), guarding their non-reserved classification against regressions. Co-authored-by: Isaac --- .../command/v2/AutoCdcParserSuite.scala | 22 +++++++++++++++++++ .../pipelines/graph/SqlPipelineSuite.scala | 18 +++++++++++++++ 2 files changed, 40 insertions(+) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala index ace2f688ca157..120da0d102313 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala @@ -972,4 +972,26 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { |TRACK HISTORY ON value1, value2""".stripMargin) } } + + test("new keywords history, track, scd remain usable as identifiers") { + // history / track / scd are non-reserved, so they must still parse as ordinary table and + // column identifiers. This guards their non-reserved classification against regressions. + val plan = parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO track + |FROM STREAM(scd) + |KEYS (history) + |SEQUENCE BY track + |COLUMNS (history, track, scd) + |STORED AS SCD TYPE 2 + |TRACK HISTORY ON (history, scd)""".stripMargin) + + val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto] + assert(cdc.targetTable.asInstanceOf[UnresolvedIdentifier].nameParts == Seq("track")) + assert(streamSource(cdc.source).multipartIdentifier == Seq("scd")) + assert(cdc.keys.map(_.name) == Seq("history")) + assert(cdc.sequenceByExpr == UnresolvedAttribute("track")) + assert(cdc.includeColumns.get.map(_.name) == Seq("history", "track", "scd")) + assert(cdc.storedAsScdType == 2) + assert(cdc.trackHistoryColumns.get.map(_.name) == Seq("history", "scd")) + } } diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala index b751c9ff8f890..cec8db6ec5288 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/graph/SqlPipelineSuite.scala @@ -1290,6 +1290,24 @@ class SqlPipelineSuite extends PipelineTest with SharedSparkSession { assert(ex.getMessage.contains("TRACK HISTORY requires STORED AS SCD TYPE 2")) } + test("AUTO CDC TRACK HISTORY with explicit STORED AS SCD TYPE 1 is rejected") { + // The implicit-default-SCD1 case is covered above; this pins the explicit SCD TYPE 1 path. + val ex = intercept[SqlGraphElementRegistrationException] { + unresolvedDataflowGraphFromSql( + sqlText = s""" + |CREATE STREAMING TABLE st + |FLOW AUTO CDC + |FROM STREAM $externalTable1Ident + |KEYS (id) + |SEQUENCE BY id + |STORED AS SCD TYPE 1 + |TRACK HISTORY ON (id) + |""".stripMargin + ) + } + assert(ex.getMessage.contains("TRACK HISTORY requires STORED AS SCD TYPE 2")) + } + test("Explicit column list is rejected for CREATE STREAMING TABLE FLOW AUTO CDC") { val ex = intercept[SqlGraphElementRegistrationException] { unresolvedDataflowGraphFromSql( From 33fce78328510230c5925417a173322886da5b2a Mon Sep 17 00:00:00 2001 From: andreas-neumann_data Date: Thu, 23 Jul 2026 01:49:30 +0000 Subject: [PATCH 6/7] Add HISTORY, TRACK, SCD to ThriftServerWithSparkContextSuite keyword list The SPARK-43119 "Get SQL Keywords" test asserts against a hardcoded keyword string (a fourth keyword-list surface, alongside the docs table, the keywords*.sql.out golden files, and the JDBC getSQLKeywords string). Insert the three new non-reserved keywords in their sorted positions. Co-authored-by: Isaac --- .../hive/thriftserver/ThriftServerWithSparkContextSuite.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sql/hive-thriftserver/src/test/scala/org/apache/spark/sql/hive/thriftserver/ThriftServerWithSparkContextSuite.scala b/sql/hive-thriftserver/src/test/scala/org/apache/spark/sql/hive/thriftserver/ThriftServerWithSparkContextSuite.scala index ff9e203f5915d..fbc150f79ab43 100644 --- a/sql/hive-thriftserver/src/test/scala/org/apache/spark/sql/hive/thriftserver/ThriftServerWithSparkContextSuite.scala +++ b/sql/hive-thriftserver/src/test/scala/org/apache/spark/sql/hive/thriftserver/ThriftServerWithSparkContextSuite.scala @@ -214,7 +214,7 @@ trait ThriftServerWithSparkContextSuite extends SharedThriftServer { val sessionHandle = client.openSession(user, "") val infoValue = client.getInfo(sessionHandle, GetInfoType.CLI_ODBC_KEYWORDS) // scalastyle:off line.size.limit - assert(infoValue.getStringValue == "ADD,AFTER,AGGREGATE,ALIGN,ALL,ALTER,ALWAYS,ANALYZE,AND,ANTI,ANY,ANY_VALUE,APPLY,APPROX,ARCHIVE,ARRAY,AS,ASC,ASENSITIVE,ASOF,AT,ATOMIC,AUTHORIZATION,AUTO,BEGIN,BERNOULLI,BETWEEN,BIGINT,BIN,BINARY,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BOOLEAN,BOTH,BUCKET,BUCKETS,BY,BYTE,CACHE,CALL,CALLED,CASCADE,CASE,CAST,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CHAR,CHARACTER,CHECK,CLEAR,CLOSE,CLUSTER,CLUSTERED,CODEGEN,COLLATE,COLLATION,COLLATIONS,COLLECTION,COLUMN,COLUMNS,COMMENT,COMMIT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONDITION,CONSTRAINT,CONTAINS,CONTINUE,COST,CREATE,CROSS,CUBE,CURRENT,CURRENT_DATABASE,CURRENT_DATE,CURRENT_PATH,CURRENT_SCHEMA,CURRENT_TIME,CURRENT_TIMESTAMP,CURRENT_USER,CURSOR,DATA,DATABASE,DATABASES,DATE,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAY,DAYOFYEAR,DAYS,DBPROPERTIES,DEC,DECIMAL,DECLARE,DEFAULT,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELETE,DELIMITED,DESC,DESCRIBE,DETERMINISTIC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTINCT,DISTRIBUTE,DIV,DO,DOUBLE,DROP,ELSE,ELSEIF,END,ENFORCED,ESCAPE,ESCAPED,EVOLUTION,EXACT,EXCEPT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXECUTE,EXISTS,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,EXTERNAL,EXTRACT,FALSE,FETCH,FIELDS,FILEFORMAT,FILTER,FIRST,FLOAT,FLOW,FOLLOWING,FOR,FOREIGN,FORMAT,FORMATTED,FOUND,FROM,FULL,FUNCTION,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,GLOBAL,GRANT,GROUP,GROUPING,HANDLER,HAVING,HOUR,HOURS,IDENTIFIED,IDENTIFIER,IDENTITY,IF,IGNORE,ILIKE,IMMEDIATE,IMPORT,IN,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INNER,INPATH,INPUT,INPUTFORMAT,INSENSITIVE,INSERT,INT,INTEGER,INTERSECT,INTERVAL,INTO,INVOKER,IS,ITEMS,ITERATE,JOIN,JSON,KEY,KEYS,LANGUAGE,LAST,LATERAL,LAZY,LEADING,LEAVE,LEFT,LEVEL,LIKE,LIMIT,LINES,LIST,LOAD,LOCAL,LOCALTIME,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MAX,MEASURE,MERGE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTE,MINUTES,MODIFIES,MONTH,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NATURAL,NEAREST,NEXT,NO,NONE,NORELY,NOT,NULL,NULLS,NUMERIC,OF,OFFSET,ON,ONLY,OPEN,OPTION,OPTIONS,OR,ORDER,OUT,OUTER,OUTPUTFORMAT,OVER,OVERLAPS,OVERLAY,OVERWRITE,PARTITION,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,POSITION,PRECEDING,PRIMARY,PRINCIPALS,PROCEDURE,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RANGE,READ,READS,REAL,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,RECURSIVE,REDUCE,REFERENCES,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,RETURN,RETURNS,REVOKE,RIGHT,ROLE,ROLES,ROLLBACK,ROLLUP,ROW,ROWS,SCHEMA,SCHEMAS,SECOND,SECONDS,SECURITY,SELECT,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SESSION_USER,SET,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SMALLINT,SOME,SORT,SORTED,SOURCE,SPECIFIC,SQL,SQLEXCEPTION,SQLSTATE,START,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SUBSTRING,SYNC,SYSTEM,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLE,TABLES,TABLESAMPLE,TARGET,TBLPROPERTIES,TERMINATED,THEN,TIME,TIMEDIFF,TIMESTAMP,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TO,TOUCH,TRAILING,TRANSACTION,TRANSACTIONS,TRANSFORM,TRIM,TRUE,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNION,UNIQUE,UNKNOWN,UNLOCK,UNPIVOT,UNSET,UNTIL,UPDATE,USE,USER,USING,VALUE,VALUES,VAR,VARCHAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHEN,WHERE,WHILE,WIDTH,WINDOW,WITH,WITHIN,WITHOUT,X,YEAR,YEARS,ZONE") + assert(infoValue.getStringValue == "ADD,AFTER,AGGREGATE,ALIGN,ALL,ALTER,ALWAYS,ANALYZE,AND,ANTI,ANY,ANY_VALUE,APPLY,APPROX,ARCHIVE,ARRAY,AS,ASC,ASENSITIVE,ASOF,AT,ATOMIC,AUTHORIZATION,AUTO,BEGIN,BERNOULLI,BETWEEN,BIGINT,BIN,BINARY,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BOOLEAN,BOTH,BUCKET,BUCKETS,BY,BYTE,CACHE,CALL,CALLED,CASCADE,CASE,CAST,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CHAR,CHARACTER,CHECK,CLEAR,CLOSE,CLUSTER,CLUSTERED,CODEGEN,COLLATE,COLLATION,COLLATIONS,COLLECTION,COLUMN,COLUMNS,COMMENT,COMMIT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONDITION,CONSTRAINT,CONTAINS,CONTINUE,COST,CREATE,CROSS,CUBE,CURRENT,CURRENT_DATABASE,CURRENT_DATE,CURRENT_PATH,CURRENT_SCHEMA,CURRENT_TIME,CURRENT_TIMESTAMP,CURRENT_USER,CURSOR,DATA,DATABASE,DATABASES,DATE,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAY,DAYOFYEAR,DAYS,DBPROPERTIES,DEC,DECIMAL,DECLARE,DEFAULT,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELETE,DELIMITED,DESC,DESCRIBE,DETERMINISTIC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTINCT,DISTRIBUTE,DIV,DO,DOUBLE,DROP,ELSE,ELSEIF,END,ENFORCED,ESCAPE,ESCAPED,EVOLUTION,EXACT,EXCEPT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXECUTE,EXISTS,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,EXTERNAL,EXTRACT,FALSE,FETCH,FIELDS,FILEFORMAT,FILTER,FIRST,FLOAT,FLOW,FOLLOWING,FOR,FOREIGN,FORMAT,FORMATTED,FOUND,FROM,FULL,FUNCTION,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,GLOBAL,GRANT,GROUP,GROUPING,HANDLER,HAVING,HISTORY,HOUR,HOURS,IDENTIFIED,IDENTIFIER,IDENTITY,IF,IGNORE,ILIKE,IMMEDIATE,IMPORT,IN,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INNER,INPATH,INPUT,INPUTFORMAT,INSENSITIVE,INSERT,INT,INTEGER,INTERSECT,INTERVAL,INTO,INVOKER,IS,ITEMS,ITERATE,JOIN,JSON,KEY,KEYS,LANGUAGE,LAST,LATERAL,LAZY,LEADING,LEAVE,LEFT,LEVEL,LIKE,LIMIT,LINES,LIST,LOAD,LOCAL,LOCALTIME,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MAX,MEASURE,MERGE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTE,MINUTES,MODIFIES,MONTH,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NATURAL,NEAREST,NEXT,NO,NONE,NORELY,NOT,NULL,NULLS,NUMERIC,OF,OFFSET,ON,ONLY,OPEN,OPTION,OPTIONS,OR,ORDER,OUT,OUTER,OUTPUTFORMAT,OVER,OVERLAPS,OVERLAY,OVERWRITE,PARTITION,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,POSITION,PRECEDING,PRIMARY,PRINCIPALS,PROCEDURE,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RANGE,READ,READS,REAL,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,RECURSIVE,REDUCE,REFERENCES,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,RETURN,RETURNS,REVOKE,RIGHT,ROLE,ROLES,ROLLBACK,ROLLUP,ROW,ROWS,SCD,SCHEMA,SCHEMAS,SECOND,SECONDS,SECURITY,SELECT,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SESSION_USER,SET,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SMALLINT,SOME,SORT,SORTED,SOURCE,SPECIFIC,SQL,SQLEXCEPTION,SQLSTATE,START,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SUBSTRING,SYNC,SYSTEM,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLE,TABLES,TABLESAMPLE,TARGET,TBLPROPERTIES,TERMINATED,THEN,TIME,TIMEDIFF,TIMESTAMP,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TO,TOUCH,TRACK,TRAILING,TRANSACTION,TRANSACTIONS,TRANSFORM,TRIM,TRUE,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNION,UNIQUE,UNKNOWN,UNLOCK,UNPIVOT,UNSET,UNTIL,UPDATE,USE,USER,USING,VALUE,VALUES,VAR,VARCHAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHEN,WHERE,WHILE,WIDTH,WINDOW,WITH,WITHIN,WITHOUT,X,YEAR,YEARS,ZONE") // scalastyle:on line.size.limit } } From 2b85dd37d91444cd355f57d021602c99fbb69ff4 Mon Sep 17 00:00:00 2001 From: andreas-neumann_data Date: Thu, 23 Jul 2026 03:01:53 +0000 Subject: [PATCH 7/7] address review: guard SCD type overflow and rename grammar label Two review nits on the STORED AS SCD TYPE handling: - INTEGER_VALUE is DIGIT+, so an oversized literal (e.g. SCD TYPE 99999999999999999999) made `.toInt` throw an uncaught NumberFormatException, surfacing as an internal error instead of the intended "Unsupported SCD type" message. Match on the token text ("1"/"2") and route everything else through operationNotAllowed, closing the overflow hole. - Rename the grammar label `type=INTEGER_VALUE` to `scdType=INTEGER_VALUE` to avoid the Scala-keyword collision that forced a backtick escape in the AstBuilder. Add a parser test for the overflowing-number case. Co-authored-by: Isaac --- .../spark/sql/catalyst/parser/SqlBaseParser.g4 | 2 +- .../spark/sql/catalyst/parser/AstBuilder.scala | 17 ++++++++++------- .../command/v2/AutoCdcParserSuite.scala | 14 ++++++++++++++ 3 files changed, 25 insertions(+), 8 deletions(-) diff --git a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 index 441011f9bb80c..cdeb34c6a0ca6 100644 --- a/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 +++ b/sql/api/src/main/antlr4/org/apache/spark/sql/catalyst/parser/SqlBaseParser.g4 @@ -793,7 +793,7 @@ autoCdcColumnsClause ; autoCdcStoredAsClause - : STORED AS SCD TYPE type=INTEGER_VALUE + : STORED AS SCD TYPE scdType=INTEGER_VALUE ; autoCdcTrackHistoryClause diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala index 2e77eae904994..26d88c07af5c8 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala @@ -1434,16 +1434,19 @@ class AstBuilder extends DataTypeAstBuilder } // STORED AS SCD TYPE . Absent clause defaults to SCD Type 1. Only 1 and 2 are - // supported; reject anything else with a clear error rather than a generic parse failure. + // supported; reject anything else (including oversized numeric literals) with a clear + // error rather than a generic parse failure or NumberFormatException. Match on the token + // text rather than parsing to Int so an overflowing literal cannot throw. val storedAsScdType = Option(params.autoCdcStoredAsClause()) match { case Some(c) => - val scdType = c.`type`.getText.toInt - if (scdType != 1 && scdType != 2) { - operationNotAllowed( - s"Unsupported SCD type: $scdType. AUTO CDC only supports STORED AS SCD TYPE 1 " + - "or SCD TYPE 2.", c) + c.scdType.getText match { + case "1" => 1 + case "2" => 2 + case other => + operationNotAllowed( + s"Unsupported SCD type: $other. AUTO CDC only supports STORED AS SCD TYPE 1 " + + "or SCD TYPE 2.", c) } - scdType case None => 1 } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala index 120da0d102313..34642497eb57a 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/command/v2/AutoCdcParserSuite.scala @@ -331,6 +331,20 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest { assert(ex.getMessage.contains("Unsupported SCD type: 3")) } + test("AUTO CDC - STORED AS SCD TYPE with an overflowing number is rejected cleanly") { + // The number is INTEGER_VALUE (DIGIT+), so an oversized literal must route through the + // clear "Unsupported SCD type" error rather than throwing a NumberFormatException. + val ex = intercept[ParseException] { + parser.parsePlan( + """CREATE FLOW f AS AUTO CDC INTO target + |FROM STREAM(source) + |KEYS (id) + |SEQUENCE BY ts + |STORED AS SCD TYPE 99999999999999999999""".stripMargin) + } + assert(ex.getMessage.contains("Unsupported SCD type: 99999999999999999999")) + } + test("AUTO CDC - STORED AS SCD TYPE tolerates extra whitespace between words") { // SCD, TYPE and the number are separate tokens, so any whitespace between them is accepted. val plan = parser.parsePlan(