Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,25 @@ public void testDataFilters() {
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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,14 +18,22 @@
*/
package org.apache.iceberg.spark.source;

import java.util.Arrays;
import java.util.List;
import java.util.Set;
import java.util.stream.Stream;
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.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;
Expand All @@ -35,11 +43,46 @@
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> unpushableChangelogPredicates = Lists.newArrayList();
List<Predicate> pushableCandidates = Lists.newArrayList();

for (Predicate predicate : predicates) {
if (isChangelogColumnPredicate(predicate)) {
unpushableChangelogPredicates.add(predicate);
} else {
pushableCandidates.add(predicate);
}
}

Predicate[] remainingPredicates =
super.pushPredicates(pushableCandidates.toArray(new Predicate[0]));

return Stream.concat(Arrays.stream(remainingPredicates), unpushableChangelogPredicates.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()))
.anyMatch(CHANGELOG_METADATA_COLUMNS::contains);
}

@Override
public Scan build() {
Long startSnapshotId = readConf().startSnapshotId();
Expand Down
Loading