Skip to content

Commit 2cd31c0

Browse files
committed
Guard Delta CDF DV remove scans
1 parent e62f121 commit 2cd31c0

2 files changed

Lines changed: 23 additions & 4 deletions

File tree

gluten-delta/src/main/scala/org/apache/gluten/execution/DeltaScanTransformer.scala

Lines changed: 18 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ import org.apache.spark.sql.catalyst.TableIdentifier
2323
import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression}
2424
import org.apache.spark.sql.catalyst.plans.QueryPlan
2525
import org.apache.spark.sql.connector.read.streaming.SparkDataStream
26+
import org.apache.spark.sql.delta.files.TahoeRemoveFileIndex
2627
import org.apache.spark.sql.execution.FileSourceScanExec
2728
import org.apache.spark.sql.execution.datasources.HadoopFsRelation
2829
import org.apache.spark.sql.types.StructType
@@ -57,16 +58,27 @@ case class DeltaScanTransformer(
5758

5859
override protected def doValidateInternal(): ValidationResult = {
5960
if (
60-
requiredSchema.fields.exists(
61-
_.name == "__delta_internal_is_row_deleted") || requiredSchema.fields.exists(
62-
_.name == "__delta_internal_row_index")
61+
requiredSchemaIncludesDeletionVectorColumns ||
62+
cdfRemoveFilesHaveDeletionVectors
6363
) {
64-
return ValidationResult.failed(s"Deletion vector is not supported in native.")
64+
return ValidationResult.failed(DeltaScanTransformer.DELETION_VECTOR_UNSUPPORTED)
6565
}
6666

6767
super.doValidateInternal()
6868
}
6969

70+
private def requiredSchemaIncludesDeletionVectorColumns: Boolean = {
71+
requiredSchema.fields.exists(
72+
_.name == "__delta_internal_is_row_deleted") || requiredSchema.fields.exists(
73+
_.name == "__delta_internal_row_index")
74+
}
75+
76+
private def cdfRemoveFilesHaveDeletionVectors: Boolean = relation.location match {
77+
case index: TahoeRemoveFileIndex =>
78+
index.filesByVersion.exists(_.actions.exists(_.deletionVector != null))
79+
case _ => false
80+
}
81+
7082
override def doCanonicalize(): DeltaScanTransformer = {
7183
DeltaScanTransformer(
7284
relation,
@@ -91,6 +103,8 @@ case class DeltaScanTransformer(
91103

92104
object DeltaScanTransformer {
93105

106+
val DELETION_VECTOR_UNSUPPORTED = "Deletion vector is not supported in native."
107+
94108
def apply(scanExec: FileSourceScanExec): DeltaScanTransformer = {
95109
new DeltaScanTransformer(
96110
scanExec.relation,

gluten-delta/src/test/scala/org/apache/gluten/execution/DeltaSuite.scala

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -262,12 +262,17 @@ abstract class DeltaSuite extends WholeStageTransformerSuite {
262262
|""".stripMargin)
263263

264264
// Native DV scan support is separate; this regression locks down CDF correctness.
265+
import org.apache.spark.sql.execution.GlutenImplicits._
265266
val df = runAndCompare(
266267
s"""
267268
|select id, name, _change_type
268269
|from table_changes('delta_cdf_dv', 0)
269270
|order by id, name, _change_type
270271
|""".stripMargin)
272+
assert(
273+
df.fallbackSummary.fallbackNodeToReason
274+
.flatMap(_.values)
275+
.exists(_.contains(DeltaScanTransformer.DELETION_VECTOR_UNSUPPORTED)))
271276
checkAnswer(
272277
df,
273278
Seq(

0 commit comments

Comments
 (0)