Skip to content

Commit 6517820

Browse files
committed
[VL][Delta] Add CDF offload guard
1 parent d740b0d commit 6517820

5 files changed

Lines changed: 72 additions & 13 deletions

File tree

backends-velox/src-delta/main/scala/org/apache/gluten/component/VeloxDeltaComponent.scala

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
package org.apache.gluten.component
1818

1919
import org.apache.gluten.backendsapi.velox.VeloxBackend
20-
import org.apache.gluten.config.GlutenConfig
20+
import org.apache.gluten.config.{GlutenConfig, VeloxDeltaConfig}
2121
import org.apache.gluten.extension.{DeltaCDFScanStrategy, DeltaPostTransformRules, OffloadDeltaFilter, OffloadDeltaProject, OffloadDeltaScan}
2222
import org.apache.gluten.extension.columnar.heuristic.HeuristicTransform
2323
import org.apache.gluten.extension.columnar.validator.Validators
@@ -35,7 +35,11 @@ class VeloxDeltaComponent extends Component {
3535
}
3636

3737
override def injectRules(injector: Injector): Unit = {
38-
injector.spark.injectPlannerStrategy(DeltaCDFScanStrategy(_))
38+
injector.spark.injectPlannerStrategy(
39+
spark =>
40+
DeltaCDFScanStrategy(
41+
spark,
42+
() => new VeloxDeltaConfig(spark.sessionState.conf).enableChangeDataFeedScan))
3943

4044
val legacy = injector.gluten.legacy
4145
// Deletion-vector scans need no Gluten-side logical preprocessing: Delta's own

backends-velox/src-delta/main/scala/org/apache/gluten/config/VeloxDeltaConfig.scala

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ class VeloxDeltaConfig(conf: SQLConf) extends GlutenCoreConfig(conf) {
2222
import VeloxDeltaConfig._
2323

2424
def enableNativeWrite: Boolean = getConf(ENABLE_NATIVE_WRITE)
25+
26+
def enableChangeDataFeedScan: Boolean = getConf(ENABLE_CHANGE_DATA_FEED_SCAN)
2527
}
2628

2729
object VeloxDeltaConfig extends ConfigRegistry {
@@ -40,4 +42,11 @@ object VeloxDeltaConfig extends ConfigRegistry {
4042
.doc("Enable native Delta Lake write for Velox backend.")
4143
.booleanConf
4244
.createWithDefault(false)
45+
46+
val ENABLE_CHANGE_DATA_FEED_SCAN: ConfigEntry[Boolean] =
47+
buildConf("spark.gluten.sql.columnar.backend.velox.delta.enableChangeDataFeedScan")
48+
.experimental()
49+
.doc("Enable Delta Lake change data feed scan offload for Velox backend.")
50+
.booleanConf
51+
.createWithDefault(true)
4352
}

backends-velox/src-delta/test/scala/org/apache/gluten/execution/VeloxDeltaSuite.scala

Lines changed: 34 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,4 +16,37 @@
1616
*/
1717
package org.apache.gluten.execution
1818

19-
class VeloxDeltaSuite extends DeltaSuite
19+
import org.apache.gluten.config.VeloxDeltaConfig
20+
21+
import org.apache.spark.sql.Row
22+
23+
class VeloxDeltaSuite extends DeltaSuite {
24+
testWithMinSparkVersion("delta: change data feed scan offload can be disabled", "3.2") {
25+
withTable("delta_cdf_disabled") {
26+
spark.sql("""
27+
|create table delta_cdf_disabled (id int, name string) using delta
28+
|tblproperties ("delta.enableChangeDataFeed" = "true")
29+
|""".stripMargin)
30+
spark.sql("""
31+
|insert into delta_cdf_disabled values (1, "v1"), (2, "v2")
32+
|""".stripMargin)
33+
34+
withSQLConf(VeloxDeltaConfig.ENABLE_CHANGE_DATA_FEED_SCAN.key -> "false") {
35+
val df = spark.sql("""
36+
|select id, name, _change_type, _commit_version
37+
|from table_changes('delta_cdf_disabled', 1)
38+
|""".stripMargin)
39+
checkAnswer(
40+
df,
41+
Seq(
42+
Row(1, "v1", "insert", 1L),
43+
Row(2, "v2", "insert", 1L)))
44+
assert(
45+
df.queryExecution.executedPlan.collect {
46+
case scan: DeltaScanTransformer => scan
47+
}.isEmpty,
48+
df.queryExecution.executedPlan)
49+
}
50+
}
51+
}
52+
}

docs/get-started/VeloxDelta.md

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,12 @@ Native Delta write is controlled by:
2828
- Default: `false`
2929
- Type: experimental
3030

31+
Native change data feed scan offload is controlled by:
32+
33+
- `spark.gluten.sql.columnar.backend.velox.delta.enableChangeDataFeedScan`
34+
- Default: `true`
35+
- Type: experimental
36+
3137
| Feature | Delta minWriterVersion | Delta minReaderVersion | Iceberg format-version | Feature type | Supported by Gluten (Velox) |
3238
|---|---:|---:|---:|---|---|
3339
| Basic functionality | 2 | 1 | 1 | Writer | Yes |

gluten-delta/src/main/scala/org/apache/gluten/extension/DeltaCDFScanStrategy.scala

Lines changed: 17 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -25,16 +25,23 @@ import org.apache.spark.sql.delta.commands.cdc.CDCReader
2525
import org.apache.spark.sql.execution.{SparkPlan, SparkStrategy}
2626
import org.apache.spark.sql.execution.datasources.LogicalRelation
2727

28-
case class DeltaCDFScanStrategy(spark: SparkSession) extends SparkStrategy {
29-
override def apply(plan: LogicalPlan): Seq[SparkPlan] = plan match {
30-
case PhysicalOperation(projects, filters, relation: LogicalRelation) =>
31-
relation.relation match {
32-
case cdfRelation: CDCReader.DeltaCDFRelation
33-
if !changesContainDeletionVectors(cdfRelation) =>
34-
planCDFRelation(relation, cdfRelation, projects, filters).map(planLater).toSeq
35-
case _ => Nil
36-
}
37-
case _ => Nil
28+
case class DeltaCDFScanStrategy(spark: SparkSession, offloadEnabled: () => Boolean)
29+
extends SparkStrategy {
30+
override def apply(plan: LogicalPlan): Seq[SparkPlan] = {
31+
if (!offloadEnabled()) {
32+
return Nil
33+
}
34+
35+
plan match {
36+
case PhysicalOperation(projects, filters, relation: LogicalRelation) =>
37+
relation.relation match {
38+
case cdfRelation: CDCReader.DeltaCDFRelation
39+
if !changesContainDeletionVectors(cdfRelation) =>
40+
planCDFRelation(relation, cdfRelation, projects, filters).map(planLater).toSeq
41+
case _ => Nil
42+
}
43+
case _ => Nil
44+
}
3845
}
3946

4047
// Delta CDF over a commit that uses deletion vectors needs Delta's DV-aware, row-level

0 commit comments

Comments
 (0)