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;

Reply via email to