diff --git a/docs/sql-ref-ansi-compliance.md b/docs/sql-ref-ansi-compliance.md index 367abce6cb8d3..f4ac58a4e3f02 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| @@ -741,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| @@ -803,6 +805,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..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 @@ -300,6 +300,7 @@ GROUPING: 'GROUPING'; HANDLER: 'HANDLER'; HAVING: 'HAVING'; BINARY_HEX: 'X'; +HISTORY: 'HISTORY'; HOUR: 'HOUR'; HOURS: 'HOURS'; IDENTIFIER_KW: 'IDENTIFIER'; @@ -458,6 +459,7 @@ ROW: 'ROW'; ROWS: 'ROWS'; SECOND: 'SECOND'; SECONDS: 'SECONDS'; +SCD: 'SCD'; SCHEMA: 'SCHEMA'; SCHEMAS: 'SCHEMAS'; SECURITY: 'SECURITY'; @@ -519,6 +521,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..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 @@ -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 scdType=INTEGER_VALUE + ; + +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 @@ -2292,6 +2305,7 @@ ansiNonReserved | ROLLUP | ROW | ROWS + | SCD | SCHEMA | SCHEMAS | SECOND @@ -2346,6 +2360,7 @@ ansiNonReserved | TIMESTAMPDIFF | TINYINT | TOUCH + | TRACK | TRANSACTION | TRANSACTIONS | TRANSFORM @@ -2586,6 +2601,7 @@ nonReserved | GROUPING | HANDLER | HAVING + | HISTORY | HOUR | HOURS | IDENTIFIER_KW @@ -2735,6 +2751,7 @@ nonReserved | ROLLUP | ROW | ROWS + | SCD | SCHEMA | SCHEMAS | SECOND @@ -2795,6 +2812,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..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 @@ -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,43 @@ class AstBuilder extends DataTypeAstBuilder visitIdentifierSeq(c.exceptCols).map(UnresolvedAttribute.quoted) } + // STORED AS SCD TYPE . Absent clause defaults to SCD Type 1. Only 1 and 2 are + // 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) => + 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) + } + case None => 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 +8071,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/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..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,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,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/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/resources/sql-tests/results/keywords-enforced.sql.out b/sql/core/src/test/resources/sql-tests/results/keywords-enforced.sql.out index 989bfb7251117..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 @@ -176,6 +176,7 @@ GROUP true GROUPING false HANDLER false HAVING true +HISTORY false HOUR false HOURS false IDENTIFIED false @@ -332,6 +333,7 @@ ROLLBACK false ROLLUP false ROW false ROWS false +SCD false SCHEMA false SCHEMAS false SECOND false @@ -392,6 +394,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..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 @@ -176,6 +176,7 @@ GROUP false GROUPING false HANDLER false HAVING false +HISTORY false HOUR false HOURS false IDENTIFIED false @@ -332,6 +333,7 @@ ROLLBACK false ROLLUP false ROW false ROWS false +SCD false SCHEMA false SCHEMAS false SECOND false @@ -392,6 +394,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..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 @@ -176,6 +176,7 @@ GROUP false GROUPING false HANDLER false HAVING false +HISTORY false HOUR false HOURS false IDENTIFIED false @@ -332,6 +333,7 @@ ROLLBACK false ROLLUP false ROW false ROWS false +SCD false SCHEMA false SCHEMAS false SECOND false @@ -392,6 +394,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/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..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 @@ -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,164 @@ 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") { + 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 3""".stripMargin) + } + 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( + """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) + } + // --------------------------------------------------------------------------- // CREATE STREAMING TABLE ... FLOW AUTO CDC // --------------------------------------------------------------------------- @@ -810,35 +973,39 @@ 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 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) + } } - 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("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/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 } } 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..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 @@ -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,99 @@ 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("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(