yingjianwu98 commented on code in PR #18148:
URL: https://github.com/apache/iceberg/pull/18148#discussion_r4167892541


##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkChangelogScanBuilder.java:
##########
@@ -35,11 +43,45 @@
 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> changelogPredicates = Lists.newArrayList();
+    List<Predicate> pushable = Lists.newArrayList();
+
+    for (Predicate predicate : predicates) {
+      if (isChangelogColumnPredicate(predicate)) {
+        changelogPredicates.add(predicate);
+      } else {
+        pushable.add(predicate);
+      }
+    }
+
+    Predicate[] remainingPredicates = 
super.pushPredicates(pushable.toArray(new Predicate[0]));
+
+    return Stream.concat(Arrays.stream(remainingPredicates), 
changelogPredicates.stream())
+        .toArray(Predicate[]::new);

Review Comment:
   sounds good to me!



##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkChangelogScanBuilder.java:
##########
@@ -35,11 +43,45 @@
 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> changelogPredicates = Lists.newArrayList();
+    List<Predicate> pushable = Lists.newArrayList();
+
+    for (Predicate predicate : predicates) {
+      if (isChangelogColumnPredicate(predicate)) {
+        changelogPredicates.add(predicate);
+      } else {
+        pushable.add(predicate);
+      }
+    }
+
+    Predicate[] remainingPredicates = 
super.pushPredicates(pushable.toArray(new Predicate[0]));
+
+    return Stream.concat(Arrays.stream(remainingPredicates), 
changelogPredicates.stream())
+        .toArray(Predicate[]::new);
+  }
+
+  // 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) {
+    return Arrays.stream(predicate.references())
+        .flatMap(ref -> Arrays.stream(ref.fieldNames()))

Review Comment:
   Good catch!



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to