This is an automated email from the ASF dual-hosted git repository.
huaxingao pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new 411c38b58c Spark: Skip unsupported aggregate pushdown (#18299)
411c38b58c is described below
commit 411c38b58c5208dd8e972a394f0e409db32458ac
Author: Akash Malbari <[email protected]>
AuthorDate: Fri Oct 2 19:34:57 2026 -0400
Spark: Skip unsupported aggregate pushdown (#18299)
* Spark: Skip unsupported aggregate pushdown
* Spark 4.2: Test row lineage aggregate fallback
* Spark: Strengthen row lineage aggregate tests
* Spark: Cover aggregate fallback in all versions
---
.../apache/iceberg/spark/source/SparkScanBuilder.java | 3 ++-
.../apache/iceberg/spark/sql/TestAggregatePushDown.java | 16 ++++++++++++++++
.../apache/iceberg/spark/source/SparkScanBuilder.java | 3 ++-
.../apache/iceberg/spark/sql/TestAggregatePushDown.java | 16 ++++++++++++++++
.../apache/iceberg/spark/source/SparkScanBuilder.java | 3 ++-
.../apache/iceberg/spark/sql/TestAggregatePushDown.java | 16 ++++++++++++++++
.../apache/iceberg/spark/source/SparkScanBuilder.java | 3 ++-
.../apache/iceberg/spark/sql/TestAggregatePushDown.java | 16 ++++++++++++++++
8 files changed, 72 insertions(+), 4 deletions(-)
diff --git
a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
index 8b75906a62..9710b0015d 100644
---
a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
+++
b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
@@ -40,6 +40,7 @@ import org.apache.iceberg.SparkDistributedDataScan;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.Table;
import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.exceptions.ValidationException;
import org.apache.iceberg.expressions.AggregateEvaluator;
import org.apache.iceberg.expressions.Binder;
import org.apache.iceberg.expressions.BoundAggregate;
@@ -225,7 +226,7 @@ public class SparkScanBuilder
aggregateFunc);
return false;
}
- } catch (IllegalArgumentException e) {
+ } catch (IllegalArgumentException | ValidationException e) {
LOG.info("Skipping aggregate pushdown: Bind failed for AggregateFunc
{}", aggregateFunc, e);
return false;
}
diff --git
a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
index 646e96eb54..4de61f9e6b 100644
---
a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
+++
b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
@@ -94,6 +94,22 @@ public class TestAggregatePushDown extends CatalogTestBase {
testDifferentDataTypesAggregatePushDown(false);
}
+ @TestTemplate
+ public void testAggregatePushDownWithRowLineageMetadataColumn() {
+ sql("CREATE TABLE %s (id INT) USING iceberg TBLPROPERTIES
('format-version'='3')", tableName);
+ sql("INSERT INTO %s VALUES (1), (2), (3)", tableName);
+
+ String select = "SELECT max(_row_id) FROM %s";
+ String explainString = sql("EXPLAIN " + select,
tableName).get(0)[0].toString();
+ assertThat(explainString)
+ .as("explain should not contain the pushed down aggregate")
+ .doesNotContain("max(_row_id)");
+
+ List<Object[]> actual = sql(select, tableName);
+ assertThat(actual).hasSize(1);
+ assertThat(actual.get(0)[0]).isEqualTo(2L);
+ }
+
@SuppressWarnings("checkstyle:CyclomaticComplexity")
private void testDifferentDataTypesAggregatePushDown(boolean
hasPartitionCol) {
String createTable;
diff --git
a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
index f594844751..ca00032f16 100644
---
a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
+++
b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
@@ -40,6 +40,7 @@ import org.apache.iceberg.SparkDistributedDataScan;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.Table;
import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.exceptions.ValidationException;
import org.apache.iceberg.expressions.AggregateEvaluator;
import org.apache.iceberg.expressions.Binder;
import org.apache.iceberg.expressions.BoundAggregate;
@@ -225,7 +226,7 @@ public class SparkScanBuilder
aggregateFunc);
return false;
}
- } catch (IllegalArgumentException e) {
+ } catch (IllegalArgumentException | ValidationException e) {
LOG.info("Skipping aggregate pushdown: Bind failed for AggregateFunc
{}", aggregateFunc, e);
return false;
}
diff --git
a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
index 1669301d2d..f46cfdcad6 100644
---
a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
+++
b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
@@ -94,6 +94,22 @@ public class TestAggregatePushDown extends CatalogTestBase {
testDifferentDataTypesAggregatePushDown(false);
}
+ @TestTemplate
+ public void testAggregatePushDownWithRowLineageMetadataColumn() {
+ sql("CREATE TABLE %s (id INT) USING iceberg TBLPROPERTIES
('format-version'='3')", tableName);
+ sql("INSERT INTO %s VALUES (1), (2), (3)", tableName);
+
+ String select = "SELECT max(_row_id) FROM %s";
+ String explainString = sql("EXPLAIN " + select,
tableName).get(0)[0].toString();
+ assertThat(explainString)
+ .as("explain should not contain the pushed down aggregate")
+ .doesNotContain("max(_row_id)");
+
+ List<Object[]> actual = sql(select, tableName);
+ assertThat(actual).hasSize(1);
+ assertThat(actual.get(0)[0]).isEqualTo(2L);
+ }
+
@SuppressWarnings("checkstyle:CyclomaticComplexity")
private void testDifferentDataTypesAggregatePushDown(boolean
hasPartitionCol) {
String createTable;
diff --git
a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
index 6b9e314d35..6e0d895c9b 100644
---
a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
+++
b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
@@ -35,6 +35,7 @@ import org.apache.iceberg.Snapshot;
import org.apache.iceberg.SparkDistributedDataScan;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.Table;
+import org.apache.iceberg.exceptions.ValidationException;
import org.apache.iceberg.expressions.AggregateEvaluator;
import org.apache.iceberg.expressions.Binder;
import org.apache.iceberg.expressions.BoundAggregate;
@@ -147,7 +148,7 @@ public class SparkScanBuilder extends BaseSparkScanBuilder
aggregateFunc);
return false;
}
- } catch (IllegalArgumentException e) {
+ } catch (IllegalArgumentException | ValidationException e) {
LOG.info("Skipping aggregate pushdown: Bind failed for AggregateFunc
{}", aggregateFunc, e);
return false;
}
diff --git
a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
index 1669301d2d..f46cfdcad6 100644
---
a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
+++
b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
@@ -94,6 +94,22 @@ public class TestAggregatePushDown extends CatalogTestBase {
testDifferentDataTypesAggregatePushDown(false);
}
+ @TestTemplate
+ public void testAggregatePushDownWithRowLineageMetadataColumn() {
+ sql("CREATE TABLE %s (id INT) USING iceberg TBLPROPERTIES
('format-version'='3')", tableName);
+ sql("INSERT INTO %s VALUES (1), (2), (3)", tableName);
+
+ String select = "SELECT max(_row_id) FROM %s";
+ String explainString = sql("EXPLAIN " + select,
tableName).get(0)[0].toString();
+ assertThat(explainString)
+ .as("explain should not contain the pushed down aggregate")
+ .doesNotContain("max(_row_id)");
+
+ List<Object[]> actual = sql(select, tableName);
+ assertThat(actual).hasSize(1);
+ assertThat(actual.get(0)[0]).isEqualTo(2L);
+ }
+
@SuppressWarnings("checkstyle:CyclomaticComplexity")
private void testDifferentDataTypesAggregatePushDown(boolean
hasPartitionCol) {
String createTable;
diff --git
a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
index d8e77ad7ff..5d5f181362 100644
---
a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
+++
b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkScanBuilder.java
@@ -35,6 +35,7 @@ import org.apache.iceberg.Snapshot;
import org.apache.iceberg.SparkDistributedDataScan;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.Table;
+import org.apache.iceberg.exceptions.ValidationException;
import org.apache.iceberg.expressions.AggregateEvaluator;
import org.apache.iceberg.expressions.Binder;
import org.apache.iceberg.expressions.BoundAggregate;
@@ -152,7 +153,7 @@ public class SparkScanBuilder extends BaseSparkScanBuilder
aggregateFunc);
return false;
}
- } catch (IllegalArgumentException e) {
+ } catch (IllegalArgumentException | ValidationException e) {
LOG.info("Skipping aggregate pushdown: Bind failed for AggregateFunc
{}", aggregateFunc, e);
return false;
}
diff --git
a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
index 1669301d2d..f46cfdcad6 100644
---
a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
+++
b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
@@ -94,6 +94,22 @@ public class TestAggregatePushDown extends CatalogTestBase {
testDifferentDataTypesAggregatePushDown(false);
}
+ @TestTemplate
+ public void testAggregatePushDownWithRowLineageMetadataColumn() {
+ sql("CREATE TABLE %s (id INT) USING iceberg TBLPROPERTIES
('format-version'='3')", tableName);
+ sql("INSERT INTO %s VALUES (1), (2), (3)", tableName);
+
+ String select = "SELECT max(_row_id) FROM %s";
+ String explainString = sql("EXPLAIN " + select,
tableName).get(0)[0].toString();
+ assertThat(explainString)
+ .as("explain should not contain the pushed down aggregate")
+ .doesNotContain("max(_row_id)");
+
+ List<Object[]> actual = sql(select, tableName);
+ assertThat(actual).hasSize(1);
+ assertThat(actual.get(0)[0]).isEqualTo(2L);
+ }
+
@SuppressWarnings("checkstyle:CyclomaticComplexity")
private void testDifferentDataTypesAggregatePushDown(boolean
hasPartitionCol) {
String createTable;