From ee6859ddb181de433b79e0752e4243284f028078 Mon Sep 17 00:00:00 2001 From: Peifeng Li Date: Sun, 20 Sep 2026 17:56:03 +0800 Subject: [PATCH 1/2] feat: expose Spark scan schema and pushed predicates Signed-off-by: Peifeng Li --- .../dev/vortex/spark/read/VortexScan.java | 24 ++++++++++ .../dev/vortex/spark/read/VortexScanTest.java | 46 +++++++++++++++++++ 2 files changed, 70 insertions(+) create mode 100644 java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java diff --git a/java/vortex-spark/src/main/java/dev/vortex/spark/read/VortexScan.java b/java/vortex-spark/src/main/java/dev/vortex/spark/read/VortexScan.java index 02a7563f925..6da2306ab69 100644 --- a/java/vortex-spark/src/main/java/dev/vortex/spark/read/VortexScan.java +++ b/java/vortex-spark/src/main/java/dev/vortex/spark/read/VortexScan.java @@ -76,6 +76,30 @@ public StructType readSchema() { return CatalogV2Util.v2ColumnsToStructType(readColumns.toArray(new Column[0])); } + /** + * Returns the table schema before projection pushdown. + * + *

An alternative execution engine must bind {@link #pushedPredicates()} against this schema. Columns used only + * by those predicates may be absent from {@link #readSchema()}. + * + * @return the full table schema + */ + public StructType tableSchema() { + return CatalogV2Util.v2ColumnsToStructType(tableColumns.toArray(new Column[0])); + } + + /** + * Returns the predicates this scan has promised to evaluate completely. + * + *

Spark may remove its own filter for these predicates. An alternative execution engine must enforce every + * returned predicate or retain this scan's original reader. + * + * @return a copy of the pushed predicate array + */ + public Predicate[] pushedPredicates() { + return pushedPredicates.clone(); + } + /** Logging-friendly readable description of the scan source. */ @Override public String description() { diff --git a/java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java b/java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java new file mode 100644 index 00000000000..0d755c514a0 --- /dev/null +++ b/java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java @@ -0,0 +1,46 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +package dev.vortex.spark.read; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; + +import java.util.List; +import java.util.Map; +import org.apache.spark.sql.connector.catalog.Column; +import org.apache.spark.sql.connector.expressions.Expression; +import org.apache.spark.sql.connector.expressions.Expressions; +import org.apache.spark.sql.connector.expressions.LiteralValue; +import org.apache.spark.sql.connector.expressions.filter.Predicate; +import org.apache.spark.sql.types.DataTypes; +import org.junit.jupiter.api.Test; + +final class VortexScanTest { + @Test + void exposesPredicateColumnsEvenWhenProjectionIsEmpty() { + Column id = Column.create("id", DataTypes.IntegerType); + Predicate predicate = new Predicate( + "=", new Expression[] {Expressions.column("id"), new LiteralValue<>(7, DataTypes.IntegerType)}); + VortexScan scan = + new VortexScan(List.of("data.vortex"), List.of(id), List.of(), new Predicate[] {predicate}, Map.of()); + + assertEquals(0, scan.readSchema().size()); + assertArrayEquals(new String[] {"id"}, scan.tableSchema().fieldNames()); + assertEquals(DataTypes.IntegerType, scan.tableSchema().apply("id").dataType()); + assertArrayEquals(new Predicate[] {predicate}, scan.pushedPredicates()); + } + + @Test + void predicateArraysCannotReplaceTheScansPredicates() { + Predicate predicate = new Predicate("IS_NULL", new Expression[] {Expressions.column("id")}); + Predicate[] predicates = {predicate}; + VortexScan scan = new VortexScan(List.of(), List.of(), List.of(), predicates, Map.of()); + + predicates[0] = null; + scan.pushedPredicates()[0] = null; + + assertArrayEquals(new Predicate[] {predicate}, scan.pushedPredicates()); + assertEquals(0, new VortexScan(List.of(), List.of(), List.of(), null, Map.of()).pushedPredicates().length); + } +} From 96a750adff827f687205580517f1ce7f20b0ef06 Mon Sep 17 00:00:00 2001 From: Peifeng Li Date: Mon, 21 Sep 2026 13:47:41 +0800 Subject: [PATCH 2/2] test: cover scan metadata after builder column pruning Signed-off-by: Peifeng Li --- .../java/dev/vortex/spark/read/VortexScanTest.java | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java b/java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java index 0d755c514a0..a4eb6899708 100644 --- a/java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java +++ b/java/vortex-spark/src/test/java/dev/vortex/spark/read/VortexScanTest.java @@ -11,19 +11,22 @@ import org.apache.spark.sql.connector.catalog.Column; import org.apache.spark.sql.connector.expressions.Expression; import org.apache.spark.sql.connector.expressions.Expressions; -import org.apache.spark.sql.connector.expressions.LiteralValue; import org.apache.spark.sql.connector.expressions.filter.Predicate; import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.StructType; import org.junit.jupiter.api.Test; final class VortexScanTest { @Test void exposesPredicateColumnsEvenWhenProjectionIsEmpty() { Column id = Column.create("id", DataTypes.IntegerType); - Predicate predicate = new Predicate( - "=", new Expression[] {Expressions.column("id"), new LiteralValue<>(7, DataTypes.IntegerType)}); - VortexScan scan = - new VortexScan(List.of("data.vortex"), List.of(id), List.of(), new Predicate[] {predicate}, Map.of()); + Predicate predicate = new Predicate("IS_NOT_NULL", new Expression[] {Expressions.column("id")}); + VortexScanBuilder builder = + new VortexScanBuilder(Map.of()).addPath("data.vortex").addColumn(id); + + assertEquals(0, builder.pushPredicates(new Predicate[] {predicate}).length); + builder.pruneColumns(new StructType()); + VortexScan scan = (VortexScan) builder.build(); assertEquals(0, scan.readSchema().size()); assertArrayEquals(new String[] {"id"}, scan.tableSchema().fieldNames());