Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions docs/sql-ref-ansi-compliance.md
Original file line number Diff line number Diff line change
Expand Up @@ -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|
Expand Down Expand Up @@ -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|
Expand Down Expand Up @@ -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|
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,7 @@ GROUPING: 'GROUPING';
HANDLER: 'HANDLER';
HAVING: 'HAVING';
BINARY_HEX: 'X';
HISTORY: 'HISTORY';
HOUR: 'HOUR';
HOURS: 'HOURS';
IDENTIFIER_KW: 'IDENTIFIER';
Expand Down Expand Up @@ -458,6 +459,7 @@ ROW: 'ROW';
ROWS: 'ROWS';
SECOND: 'SECOND';
SECONDS: 'SECONDS';
SCD: 'SCD';
SCHEMA: 'SCHEMA';
SCHEMAS: 'SCHEMAS';
SECURITY: 'SECURITY';
Expand Down Expand Up @@ -519,6 +521,7 @@ TINYINT: 'TINYINT';
TO: 'TO';
EXECUTE: 'EXECUTE';
TOUCH: 'TOUCH';
TRACK: 'TRACK';
TRAILING: 'TRAILING';
TRANSACTION: 'TRANSACTION';
TRANSACTIONS: 'TRANSACTIONS';
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -774,6 +774,8 @@ autoCdcParameters
autoCdcDeleteClause?
autoCdcSequenceByClause
autoCdcColumnsClause?
autoCdcStoredAsClause?
autoCdcTrackHistoryClause?
;

autoCdcDeleteClause
Expand All @@ -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
Expand Down Expand Up @@ -2160,6 +2172,7 @@ ansiNonReserved
| GLOBAL
| GROUPING
| HANDLER
| HISTORY
| HOUR
| HOURS
| IDENTIFIER_KW
Expand Down Expand Up @@ -2292,6 +2305,7 @@ ansiNonReserved
| ROLLUP
| ROW
| ROWS
| SCD
| SCHEMA
| SCHEMAS
| SECOND
Expand Down Expand Up @@ -2346,6 +2360,7 @@ ansiNonReserved
| TIMESTAMPDIFF
| TINYINT
| TOUCH
| TRACK
| TRANSACTION
| TRANSACTIONS
| TRANSFORM
Expand Down Expand Up @@ -2586,6 +2601,7 @@ nonReserved
| GROUPING
| HANDLER
| HAVING
| HISTORY
| HOUR
| HOURS
| IDENTIFIER_KW
Expand Down Expand Up @@ -2735,6 +2751,7 @@ nonReserved
| ROLLUP
| ROW
| ROWS
| SCD
| SCHEMA
| SCHEMAS
| SECOND
Expand Down Expand Up @@ -2795,6 +2812,7 @@ nonReserved
| TINYINT
| TO
| TOUCH
| TRACK
| TRAILING
| TRANSACTION
| TRANSACTIONS
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down Expand Up @@ -1430,13 +1433,43 @@ class AstBuilder extends DataTypeAstBuilder
visitIdentifierSeq(c.exceptCols).map(UnresolvedAttribute.quoted)
}

// STORED AS SCD TYPE <n>. 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)
}

/**
Expand Down Expand Up @@ -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]])
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 <n>`. 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,
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 <n>`. 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,
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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))

/*
=======================================================================================
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,7 @@ GROUP true
GROUPING false
HANDLER false
HAVING true
HISTORY false
HOUR false
HOURS false
IDENTIFIED false
Expand Down Expand Up @@ -332,6 +333,7 @@ ROLLBACK false
ROLLUP false
ROW false
ROWS false
SCD false
SCHEMA false
SCHEMAS false
SECOND false
Expand Down Expand Up @@ -392,6 +394,7 @@ TIMESTAMP_NTZ false
TINYINT false
TO true
TOUCH false
TRACK false
TRAILING true
TRANSACTION false
TRANSACTIONS false
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,7 @@ GROUP false
GROUPING false
HANDLER false
HAVING false
HISTORY false
HOUR false
HOURS false
IDENTIFIED false
Expand Down Expand Up @@ -332,6 +333,7 @@ ROLLBACK false
ROLLUP false
ROW false
ROWS false
SCD false
SCHEMA false
SCHEMAS false
SECOND false
Expand Down Expand Up @@ -392,6 +394,7 @@ TIMESTAMP_NTZ false
TINYINT false
TO false
TOUCH false
TRACK false
TRAILING false
TRANSACTION false
TRANSACTIONS false
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,7 @@ GROUP false
GROUPING false
HANDLER false
HAVING false
HISTORY false
HOUR false
HOURS false
IDENTIFIED false
Expand Down Expand Up @@ -332,6 +333,7 @@ ROLLBACK false
ROLLUP false
ROW false
ROWS false
SCD false
SCHEMA false
SCHEMAS false
SECOND false
Expand Down Expand Up @@ -392,6 +394,7 @@ TIMESTAMP_NTZ false
TINYINT false
TO false
TOUCH false
TRACK false
TRAILING false
TRANSACTION false
TRANSACTIONS false
Expand Down
Loading