This is an automated email from the ASF dual-hosted git repository.

bryanck 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 105d314039 Spark: Fix predicate pushdown on changelog metadata column 
(#18148)
105d314039 is described below

commit 105d314039eac7b67f2067d5a64ae364b8dfe9bb
Author: Yingjian Wu <[email protected]>
AuthorDate: Fri Oct 2 13:29:35 2026 -0700

    Spark: Fix predicate pushdown on changelog metadata column (#18148)
---
 .../spark/extensions/TestChangelogTable.java       | 19 +++++++++
 .../spark/source/SparkChangelogScanBuilder.java    | 46 ++++++++++++++++++++++
 2 files changed, 65 insertions(+)

diff --git 
a/spark/v4.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestChangelogTable.java
 
b/spark/v4.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestChangelogTable.java
index 742235935e..33b48b47df 100644
--- 
a/spark/v4.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestChangelogTable.java
+++ 
b/spark/v4.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestChangelogTable.java
@@ -94,6 +94,25 @@ public class TestChangelogTable extends ExtensionsTestBase {
         sql("SELECT * FROM %s.changes WHERE id = 3 ORDER BY _change_ordinal, 
id", tableName));
   }
 
+  @TestTemplate
+  public void testChangelogMetadataColumnFilter() {
+    createTableWithDefaultRows();
+
+    sql("INSERT INTO %s VALUES (3, 'c')", tableName);
+
+    Table table = validationCatalog.loadTable(tableIdent);
+
+    Snapshot snap3 = table.currentSnapshot();
+
+    assertEquals(
+        "Should have expected row",
+        ImmutableList.of(row(3, "c", "INSERT", 2, snap3.snapshotId())),
+        sql(
+            "SELECT * FROM %s.changes WHERE id = 3 AND _change_type = 'INSERT' 
"
+                + "AND _change_ordinal = 2 AND _commit_snapshot_id = %s",
+            tableName, snap3.snapshotId()));
+  }
+
   @TestTemplate
   public void testOverwrites() {
     createTableWithDefaultRows();
diff --git 
a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkChangelogScanBuilder.java
 
b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkChangelogScanBuilder.java
index 43b8a36507..041df728ae 100644
--- 
a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkChangelogScanBuilder.java
+++ 
b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkChangelogScanBuilder.java
@@ -18,14 +18,22 @@
  */
 package org.apache.iceberg.spark.source;
 
+import java.util.Arrays;
+import java.util.List;
+import java.util.Set;
 import org.apache.iceberg.IncrementalChangelogScan;
+import org.apache.iceberg.MetadataColumns;
 import org.apache.iceberg.Schema;
 import org.apache.iceberg.Snapshot;
 import org.apache.iceberg.Table;
 import org.apache.iceberg.relocated.com.google.common.base.Preconditions;
+import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet;
+import org.apache.iceberg.relocated.com.google.common.collect.Lists;
 import org.apache.iceberg.spark.SparkReadOptions;
 import org.apache.iceberg.util.SnapshotUtil;
 import org.apache.spark.sql.SparkSession;
+import org.apache.spark.sql.connector.expressions.NamedReference;
+import org.apache.spark.sql.connector.expressions.filter.Predicate;
 import org.apache.spark.sql.connector.read.Scan;
 import org.apache.spark.sql.connector.read.SupportsPushDownLimit;
 import org.apache.spark.sql.connector.read.SupportsPushDownRequiredColumns;
@@ -35,11 +43,49 @@ import org.apache.spark.sql.util.CaseInsensitiveStringMap;
 public class SparkChangelogScanBuilder extends BaseSparkScanBuilder
     implements SupportsPushDownV2Filters, SupportsPushDownRequiredColumns, 
SupportsPushDownLimit {
 
+  private static final Set<String> CHANGELOG_METADATA_COLUMNS =
+      ImmutableSet.of(
+          MetadataColumns.CHANGE_TYPE.name(),
+          MetadataColumns.CHANGE_ORDINAL.name(),
+          MetadataColumns.COMMIT_SNAPSHOT_ID.name());
+
   SparkChangelogScanBuilder(
       SparkSession spark, Table table, Schema schema, CaseInsensitiveStringMap 
options) {
     super(spark, table, schema, options);
   }
 
+  @Override
+  public Predicate[] pushPredicates(Predicate[] predicates) {
+    List<Predicate> postScanPredicates = Lists.newArrayList();
+    List<Predicate> pushable = Lists.newArrayList();
+
+    for (Predicate predicate : predicates) {
+      if (isChangelogColumnPredicate(predicate)) {
+        postScanPredicates.add(predicate);
+      } else {
+        pushable.add(predicate);
+      }
+    }
+
+    Predicate[] remainingPredicates = 
super.pushPredicates(pushable.toArray(new Predicate[0]));
+
+    postScanPredicates.addAll(Arrays.asList(remainingPredicates));
+    return postScanPredicates.toArray(new Predicate[0]);
+  }
+
+  // changelog metadata columns are generated by ChangelogRowReader, not part 
of the table's
+  // real schema, so leave them for Spark to evaluate after the scan instead 
of pushing them
+  // down to Iceberg.
+  private static boolean isChangelogColumnPredicate(Predicate predicate) {
+    for (NamedReference ref : predicate.references()) {
+      if (CHANGELOG_METADATA_COLUMNS.contains(ref.fieldNames()[0])) {
+        return true;
+      }
+    }
+
+    return false;
+  }
+
   @Override
   public Scan build() {
     Long startSnapshotId = readConf().startSnapshotId();

Reply via email to