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
Original file line number Diff line number Diff line change
Expand Up @@ -1368,11 +1368,11 @@ class AstBuilder extends DataTypeAstBuilder
withSchemaEvolution)
}

protected def buildAutoCdcIntoCommand(ctx: AutoCdcCommandContext): AutoCdcIntoCommand =
protected def buildAutoCdcInto(ctx: AutoCdcCommandContext): AutoCdcInto =
withOrigin(ctx) {
val target = UnresolvedIdentifier(visitMultipartIdentifier(ctx.target))
val params = parseAutoCdcParams(ctx.autoCdcParameters())
AutoCdcIntoCommand(
AutoCdcInto(
targetTable = target,
source = params.source,
keys = params.keys,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,19 +18,24 @@
package org.apache.spark.sql.catalyst.plans.logical

import org.apache.spark.sql.catalyst.analysis.UnresolvedAttribute
import org.apache.spark.sql.catalyst.expressions.Expression
import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression}
import org.apache.spark.sql.catalyst.trees.BinaryLike

/**
* Logical plan node for an AUTO CDC INTO command, used by Spark Declarative Pipelines.
* 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.
*
* This node serves as a parse-time placeholder for a pipeline CDC definition and cannot be
* executed directly. It will be interpreted by the pipeline submodule once execution support
* is added (SPARK-57402). The [[targetTable]] and [[source]] relations are exposed as the node's
* children (left and right respectively) so the analyzer resolves them through the normal plan
* resolution path.
* 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
* plain [[LogicalPlan]] rather than a [[Command]] so that it is never eagerly executed or planned
* on its own -- it is interpreted by the pipeline submodule during a pipeline execution. Unlike a
* [[ParsedStatement]], it is not rewritten into another plan during analysis; the pipeline
* submodule consumes it as-is, so it relies on the default resolution semantics (resolved once its
* children and expressions are resolved) rather than being forced unresolved. The [[targetTable]]
* and [[source]] relations are exposed as the node's children (left and right respectively) so the
* analyzer resolves them through the normal plan resolution path.
*
* @param targetTable The target table to apply changes into, as an `UnresolvedIdentifier`.
* Exposed as the node's left child.
Expand All @@ -49,19 +54,23 @@ import org.apache.spark.sql.catalyst.expressions.Expression
* except these). [[None]] when no COLUMNS clause was specified. Mutually
* exclusive with [[includeColumns]].
*/
case class AutoCdcIntoCommand(
case class AutoCdcInto(
targetTable: LogicalPlan,
source: LogicalPlan,
keys: Seq[UnresolvedAttribute],
deleteCondition: Option[Expression],
sequenceByExpr: Expression,
includeColumns: Option[Seq[UnresolvedAttribute]],
excludeColumns: Option[Seq[UnresolvedAttribute]]
) extends BinaryCommand {
) extends LogicalPlan with BinaryLike[LogicalPlan] {
override def left: LogicalPlan = targetTable
override def right: LogicalPlan = source

// This node is a parse-time placeholder that is never executed or projected from directly; the
// pipeline submodule reads its fields instead. It therefore produces no output columns.
override def output: Seq[Attribute] = Nil

override protected def withNewChildrenInternal(
newLeft: LogicalPlan, newRight: LogicalPlan): AutoCdcIntoCommand =
newLeft: LogicalPlan, newRight: LogicalPlan): AutoCdcInto =
copy(targetTable = newLeft, source = newRight)
}
Original file line number Diff line number Diff line change
Expand Up @@ -803,7 +803,7 @@ case class CreateStreamingTableAsSelect(

/**
* Command parsed from `CREATE STREAMING TABLE ...` SQL syntax. This command serves as a logical
* representation of the matching SQL syntac and cannot be executed. It is instead interpreted by
* representation of the matching SQL syntax and cannot be executed. It is instead interpreted by
* the pipeline submodule during a pipeline execution.
*
* Differs from [[CreateStreamingTableAsSelect]] in that the AS [subquery] clause is not provided
Expand All @@ -826,9 +826,9 @@ case class CreateStreamingTable(

/**
* Command parsed from `CREATE STREAMING TABLE <name> FLOW AUTO CDC ...` SQL syntax.
* This command serves as a parse-time placeholder for a pipeline CDC definition and cannot be
* executed directly. It will be interpreted by the pipeline submodule once execution support
* is added (SPARK-57402).
* This command serves as a parse-time placeholder for a pipeline CDC definition and cannot
* be executed directly. It is instead interpreted by the pipeline submodule during a pipeline
* execution.
*
* The target of the CDC operation is the streaming table itself (given by [[name]]).
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1576,7 +1576,7 @@ class SparkSqlAstBuilder extends AstBuilder {
val flowHeaderCtx = ctx.createPipelineFlowHeader()
val ident = withIdentClause(flowHeaderCtx.flowName, UnresolvedIdentifier(_))
val commentOpt = Option(flowHeaderCtx.commentSpec()).map(visitCommentSpec)
val autoCdcInto = buildAutoCdcIntoCommand(ctx.autoCdcCommand())
val autoCdcInto = buildAutoCdcInto(ctx.autoCdcCommand())
CreateFlowCommand(
name = ident,
flowOperation = autoCdcInto,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1001,7 +1001,8 @@ abstract class SparkStrategies extends QueryPlanner[SparkPlan] {
case cmv: CreateMaterializedViewAsSelect =>
throw QueryCompilationErrors.unsupportedCreatePipelineDatasetQueryExecutionError(
pipelineDatasetType = "MATERIALIZED VIEW")
case cst: CreateStreamingTableAsSelect =>
case _: CreateStreamingTableAsSelect | _: CreateStreamingTable |
_: CreateStreamingTableAutoCdc =>
throw QueryCompilationErrors.unsupportedCreatePipelineDatasetQueryExecutionError(
pipelineDatasetType = "STREAMING TABLE")
case _ => Nil
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2541,6 +2541,36 @@ abstract class DDLSuite extends QueryTest with DDLSuiteBase {
)
}

test("CREATE STREAMING TABLE without subquery cannot be directly executed") {
checkError(
exception = intercept[SparkUnsupportedOperationException] {
sql("CREATE STREAMING TABLE table1")
},
condition = "UNSUPPORTED_FEATURE.CREATE_PIPELINE_DATASET_QUERY_EXECUTION",
sqlState = "0A000",
parameters = Map("pipelineDatasetType" -> "STREAMING TABLE")
)
}

test("CREATE STREAMING TABLE FLOW AUTO CDC cannot be directly executed") {
withTable("cdc_src") {
sql("CREATE TABLE cdc_src AS SELECT 1 AS id, 1 AS ts")
checkError(
exception = intercept[SparkUnsupportedOperationException] {
sql(
"""CREATE STREAMING TABLE table1
|FLOW AUTO CDC
|FROM STREAM(cdc_src)
|KEYS (id)
|SEQUENCE BY ts""".stripMargin)
},
condition = "UNSUPPORTED_FEATURE.CREATE_PIPELINE_DATASET_QUERY_EXECUTION",
sqlState = "0A000",
parameters = Map("pipelineDatasetType" -> "STREAMING TABLE")
)
}
}

test(s"CREATE FLOW statement cannot be directly executed") {
sql("CREATE TABLE table1 AS SELECT 1")
sql("CREATE TABLE table2 AS SELECT 2")
Expand All @@ -2553,6 +2583,24 @@ abstract class DDLSuite extends QueryTest with DDLSuiteBase {
parameters = Map.empty
)
}

test("CREATE FLOW AS AUTO CDC cannot be directly executed") {
withTable("cdc_src") {
sql("CREATE TABLE cdc_src AS SELECT 1 AS id, 1 AS ts")
checkError(
exception = intercept[SparkUnsupportedOperationException] {
sql(
"""CREATE FLOW f AS AUTO CDC INTO target
|FROM STREAM(cdc_src)
|KEYS (id)
|SEQUENCE BY ts""".stripMargin)
},
condition = "UNSUPPORTED_FEATURE.CREATE_FLOW_QUERY_EXECUTION",
sqlState = "0A000",
parameters = Map.empty
)
}
}
}

object FakeLocalFsFileSystem {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import org.apache.spark.sql.catalyst.analysis.{
UnresolvedRelation}
import org.apache.spark.sql.catalyst.parser.ParseException
import org.apache.spark.sql.catalyst.plans.logical.{
AutoCdcIntoCommand,
AutoCdcInto,
CreateFlowCommand,
CreateStreamingTableAutoCdc,
LogicalPlan,
Expand Down Expand Up @@ -71,7 +71,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
assert(cmd.name.asInstanceOf[UnresolvedIdentifier].nameParts == Seq("myflow"))
assert(cmd.comment.isEmpty)

val cdc = cmd.flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = cmd.flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.targetTable.asInstanceOf[UnresolvedIdentifier].nameParts == Seq("target"))
val source = streamSource(cdc.source)
assert(source.multipartIdentifier == Seq("source"))
Expand All @@ -89,7 +89,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|KEYS (id)
|SEQUENCE BY ts""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
val source = streamSource(cdc.source)
assert(source.multipartIdentifier == Seq("mycat", "myschema", "source"))
}
Expand Down Expand Up @@ -124,7 +124,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|KEYS (k)
|SEQUENCE BY ts""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.targetTable.asInstanceOf[UnresolvedIdentifier].nameParts ==
Seq("myschema", "mytable"))
}
Expand All @@ -136,7 +136,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|KEYS (k)
|SEQUENCE BY ts""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.targetTable.asInstanceOf[UnresolvedIdentifier].nameParts ==
Seq("mycat", "myschema", "mytable"))
}
Expand All @@ -149,7 +149,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|APPLY AS DELETE WHEN op = 'DELETE'
|SEQUENCE BY ts""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.deleteCondition.isDefined)
assert(cdc.deleteCondition.get.sql.contains("op"))
}
Expand All @@ -162,7 +162,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|SEQUENCE BY ts
|COLUMNS (id, name, value)""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.includeColumns.get.map(_.name) == Seq("id", "name", "value"))
assert(cdc.excludeColumns.isEmpty)
}
Expand All @@ -175,7 +175,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|SEQUENCE BY ts
|COLUMNS * EXCEPT (op, ts)""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.includeColumns.isEmpty)
assert(cdc.excludeColumns.get.map(_.name) == Seq("op", "ts"))
}
Expand All @@ -189,7 +189,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|SEQUENCE BY timestamp
|COLUMNS (key1, key2, key3, timestamp)""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.keys.map(_.name) == Seq("key1", "key2"))
assert(cdc.deleteCondition.isDefined)
assert(cdc.sequenceByExpr == UnresolvedAttribute("timestamp"))
Expand Down Expand Up @@ -523,7 +523,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|KEYS (id)
|SEQUENCE BY ts""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
val source = streamSource(cdc.source)
assert(source.multipartIdentifier == Seq("source"))
}
Expand All @@ -535,7 +535,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {
|KEYS (id)
|SEQUENCE BY ts""".stripMargin)

val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.source.isStreaming)
val alias = cdc.source.asInstanceOf[SubqueryAlias]
assert(alias.alias == "s")
Expand All @@ -550,7 +550,7 @@ class AutoCdcParserSuite extends CommandSuiteBase with AnalysisTest {

// The subquery wraps the STREAM read in Project/Filter/SubqueryAlias nodes; isStreaming
// propagates up through them, so the whole source is recognized as streaming.
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcIntoCommand]
val cdc = plan.asInstanceOf[CreateFlowCommand].flowOperation.asInstanceOf[AutoCdcInto]
assert(cdc.source.isStreaming)
assert(cdc.source.isInstanceOf[SubqueryAlias])
}
Expand Down
Loading