From 3f989bb2c7ddf6e2f3c64dd9e9c7fd7f4420a9b3 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Sat, 18 Jul 2026 00:28:43 +0000 Subject: [PATCH 1/2] [SPARK-57402][SQL][FOLLOWUP] Allow pipeline definitions in batch operation checks Skip the generic batch streaming-source rejection for streaming table and flow definitions, which are interpreted by Spark Declarative Pipelines instead of executed as batch queries. Add direct checker coverage for both definition forms. --- .../UnsupportedOperationChecker.scala | 9 +++++++ .../analysis/UnsupportedOperationsSuite.scala | 25 +++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala index 83ad97fdc4faf..b60002cb1da6b 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala @@ -37,7 +37,16 @@ import org.apache.spark.sql.streaming.{GroupStateTimeout, OutputMode} */ object UnsupportedOperationChecker extends Logging { + private def isPipelineDefinition(plan: LogicalPlan): Boolean = plan match { + case _: CreateStreamingTableAsSelect | _: CreateStreamingTableAutoCdc | + _: CreateFlowCommand => true + case _ => false + } + def checkForBatch(plan: LogicalPlan): Unit = { + if (isPipelineDefinition(plan)) { + return + } plan.foreachUp { case p if p.isStreaming => throwError("Queries with streaming sources must be executed with writeStream.start(), or " + 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 906bb260b0ed0..92843b5d56107 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 @@ -65,6 +65,31 @@ class UnsupportedOperationsSuite extends SparkFunSuite with SQLHelper { streamRelation.select($"`count(*)`"), Seq("with streaming source", "start")) + private val emptyTableSpec = TableSpec( + properties = Map.empty, + provider = None, + options = Map.empty, + location = None, + comment = None, + collation = None, + serde = None, + external = false) + + assertSupportedInBatchPlan( + "streaming table definition with streaming source", + CreateStreamingTableAsSelect( + name = batchRelation, + columns = Nil, + partitioning = Nil, + tableSpec = emptyTableSpec, + query = streamRelation, + originalText = "", + ifNotExists = false)) + + assertSupportedInBatchPlan( + "flow definition with streaming source", + CreateFlowCommand(batchRelation, streamRelation, comment = None)) + /* ======================================================================================= From f8ecc01e3cf5b7591bd3bcde21356b93afa501a4 Mon Sep 17 00:00:00 2001 From: Wenchen Fan Date: Mon, 20 Jul 2026 15:18:01 +0000 Subject: [PATCH 2/2] =?UTF-8?q?=E2=9C=85=20test(sql):=20Cover=20AUTO=20CDC?= =?UTF-8?q?=20batch=20validation?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../analysis/UnsupportedOperationChecker.scala | 3 +++ .../analysis/UnsupportedOperationsSuite.scala | 14 ++++++++++++++ 2 files changed, 17 insertions(+) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala index b60002cb1da6b..eddcee169e377 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/analysis/UnsupportedOperationChecker.scala @@ -37,6 +37,9 @@ import org.apache.spark.sql.streaming.{GroupStateTimeout, OutputMode} */ object UnsupportedOperationChecker extends Logging { + // Pipeline definition commands are always root plans and must reach SparkStrategies.Pipelines + // for their dedicated errors. Materialized views are intentionally excluded because their + // queries are batch queries and must continue to reject streaming sources here. private def isPipelineDefinition(plan: LogicalPlan): Boolean = plan match { case _: CreateStreamingTableAsSelect | _: CreateStreamingTableAutoCdc | _: CreateFlowCommand => true 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 92843b5d56107..785ca2f6f5978 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 @@ -90,6 +90,20 @@ class UnsupportedOperationsSuite extends SparkFunSuite with SQLHelper { "flow definition with streaming source", CreateFlowCommand(batchRelation, streamRelation, comment = None)) + assertSupportedInBatchPlan( + "AUTO CDC streaming table definition with streaming source", + CreateStreamingTableAutoCdc( + name = batchRelation, + columns = Nil, + partitioning = Nil, + tableSpec = emptyTableSpec, + ifNotExists = false, + source = streamRelation, + keys = Nil, + deleteCondition = None, + sequenceByExpr = attribute, + includeColumns = None, + excludeColumns = None)) /* =======================================================================================