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();