From 803b9730000c6097dcf47ca1da6c266b0573909e Mon Sep 17 00:00:00 2001 From: Szehon Ho Date: Tue, 15 Mar 2022 21:29:19 -0700 Subject: [PATCH 1/7] Predicate pushdown for all_data_files table --- .../org/apache/iceberg/AllDataFilesTable.java | 71 +++------------ .../org/apache/iceberg/BaseFilesTable.java | 9 +- .../org/apache/iceberg/DataFilesTable.java | 6 +- .../org/apache/iceberg/DeleteFilesTable.java | 6 +- .../java/org/apache/iceberg/FilesTable.java | 6 +- .../iceberg/TestMetadataTableFilters.java | 16 +++- .../spark/extensions/TestMetadataTables.java | 90 ++++++++++++++++++- 7 files changed, 130 insertions(+), 74 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java index 4eca0b430dd4..c13dc2b289ce 100644 --- a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java @@ -20,18 +20,10 @@ package org.apache.iceberg; import java.io.IOException; -import java.util.List; -import java.util.concurrent.ExecutorService; -import org.apache.iceberg.BaseFilesTable.ManifestReadTask; import org.apache.iceberg.exceptions.RuntimeIOException; -import org.apache.iceberg.expressions.Expression; -import org.apache.iceberg.expressions.Expressions; -import org.apache.iceberg.expressions.ResidualEvaluator; import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.relocated.com.google.common.collect.Sets; -import org.apache.iceberg.types.TypeUtil; -import org.apache.iceberg.types.Types.StructType; import org.apache.iceberg.util.ParallelIterable; /** @@ -41,7 +33,7 @@ *

* This table may return duplicate rows. */ -public class AllDataFilesTable extends BaseMetadataTable { +public class AllDataFilesTable extends BaseFilesTable { AllDataFilesTable(TableOperations ops, Table table) { this(ops, table, table.name() + ".all_data_files"); @@ -56,45 +48,25 @@ public TableScan newScan() { return new AllDataFilesTableScan(operations(), table(), schema()); } - @Override - public Schema schema() { - StructType partitionType = Partitioning.partitionType(table()); - Schema schema = new Schema(DataFile.getType(partitionType).fields()); - if (partitionType.fields().size() < 1) { - // avoid returning an empty struct, which is not always supported. instead, drop the partition field (id 102) - return TypeUtil.selectNot(schema, Sets.newHashSet(102)); - } else { - return schema; - } - } - @Override MetadataTableType metadataTableType() { return MetadataTableType.ALL_DATA_FILES; } - public static class AllDataFilesTableScan extends BaseAllMetadataTableScan { - private final Schema fileSchema; + public static class AllDataFilesTableScan extends BaseFilesTableScan { AllDataFilesTableScan(TableOperations ops, Table table, Schema fileSchema) { - super(ops, table, fileSchema); - this.fileSchema = fileSchema; + super(ops, table, fileSchema, MetadataTableType.ALL_DATA_FILES); } private AllDataFilesTableScan(TableOperations ops, Table table, Schema schema, Schema fileSchema, TableScanContext context) { - super(ops, table, schema, context); - this.fileSchema = fileSchema; - } - - @Override - protected String tableType() { - return MetadataTableType.ALL_DATA_FILES.name(); + super(ops, table, schema, fileSchema, context, MetadataTableType.ALL_DATA_FILES); } @Override protected TableScan newRefinedScan(TableOperations ops, Table table, Schema schema, TableScanContext context) { - return new AllDataFilesTableScan(ops, table, schema, fileSchema, context); + return new AllDataFilesTableScan(ops, table, schema, fileSchema(), context); } @Override @@ -108,30 +80,15 @@ public TableScan asOfTime(long timestampMillis) { } @Override - protected CloseableIterable planFiles( - TableOperations ops, Snapshot snapshot, Expression rowFilter, - boolean ignoreResiduals, boolean caseSensitive, boolean colStats) { - CloseableIterable manifests = allDataManifestFiles( - ops.current().snapshots(), context().planExecutor()); - String schemaString = SchemaParser.toJson(schema()); - String specString = PartitionSpecParser.toJson(PartitionSpec.unpartitioned()); - Expression filter = ignoreResiduals ? Expressions.alwaysTrue() : rowFilter; - ResidualEvaluator residuals = ResidualEvaluator.unpartitioned(filter); - - return CloseableIterable.transform(manifests, manifest -> - new ManifestReadTask(ops.io(), ops.current().specsById(), manifest, schema(), - schemaString, specString, residuals)); - } - } - - private static CloseableIterable allDataManifestFiles( - List snapshots, ExecutorService workerPool) { - try (CloseableIterable iterable = new ParallelIterable<>( - Iterables.transform(snapshots, snapshot -> (Iterable) () -> snapshot.dataManifests().iterator()), - workerPool)) { - return CloseableIterable.withNoopClose(Sets.newHashSet(iterable)); - } catch (IOException e) { - throw new RuntimeIOException(e, "Failed to close parallel iterable"); + protected CloseableIterable manifests() { + try (CloseableIterable iterable = new ParallelIterable<>( + Iterables.transform(table().snapshots(), + snapshot -> (Iterable) () -> snapshot.dataManifests().iterator()), + context().planExecutor())) { + return CloseableIterable.withNoopClose(Sets.newHashSet(iterable)); + } catch (IOException e) { + throw new RuntimeIOException(e, "Failed to close parallel iterable"); + } } } } diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index a0b835a4abcf..48199bbd9991 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java @@ -19,7 +19,6 @@ package org.apache.iceberg; -import java.util.List; import java.util.Map; import org.apache.iceberg.expressions.Expression; import org.apache.iceberg.expressions.Expressions; @@ -110,13 +109,11 @@ protected CloseableIterable planFiles(TableOperations ops, Snapsho /** * @return list of manifest files to explore for this files metadata table scan */ - protected abstract List manifests(); + protected abstract CloseableIterable manifests(); - private CloseableIterable filterManifests(List manifests, + private CloseableIterable filterManifests(CloseableIterable manifests, Expression rowFilter, boolean caseSensitive) { - CloseableIterable manifestIterable = CloseableIterable.withNoopClose(manifests); - // use an inclusive projection to remove the partition name prefix and filter out any non-partition expressions PartitionSpec spec = transformSpec(fileSchema, table().spec(), PARTITION_FIELD_PREFIX); Expression partitionFilter = Projections.inclusive(spec, caseSensitive).project(rowFilter); @@ -124,7 +121,7 @@ private CloseableIterable filterManifests(List manif ManifestEvaluator manifestEval = ManifestEvaluator.forPartitionFilter( partitionFilter, table().spec(), caseSensitive); - return CloseableIterable.filter(manifestIterable, manifestEval::eval); + return CloseableIterable.filter(manifests, manifestEval::eval); } } diff --git a/core/src/main/java/org/apache/iceberg/DataFilesTable.java b/core/src/main/java/org/apache/iceberg/DataFilesTable.java index 6bdf3b4e5ae7..588b66687a66 100644 --- a/core/src/main/java/org/apache/iceberg/DataFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/DataFilesTable.java @@ -19,7 +19,7 @@ package org.apache.iceberg; -import java.util.List; +import org.apache.iceberg.io.CloseableIterable; /** * A {@link Table} implementation that exposes a table's data files as rows. @@ -60,8 +60,8 @@ protected TableScan newRefinedScan(TableOperations ops, Table table, Schema sche } @Override - protected List manifests() { - return snapshot().dataManifests(); + protected CloseableIterable manifests() { + return CloseableIterable.withNoopClose(snapshot().dataManifests()); } } } diff --git a/core/src/main/java/org/apache/iceberg/DeleteFilesTable.java b/core/src/main/java/org/apache/iceberg/DeleteFilesTable.java index 201c76dbd671..28744ce00f48 100644 --- a/core/src/main/java/org/apache/iceberg/DeleteFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/DeleteFilesTable.java @@ -19,7 +19,7 @@ package org.apache.iceberg; -import java.util.List; +import org.apache.iceberg.io.CloseableIterable; /** * A {@link Table} implementation that exposes a table's delete files as rows. @@ -60,8 +60,8 @@ protected TableScan newRefinedScan(TableOperations ops, Table table, Schema sche } @Override - protected List manifests() { - return snapshot().deleteManifests(); + protected CloseableIterable manifests() { + return CloseableIterable.withNoopClose(snapshot().deleteManifests()); } } } diff --git a/core/src/main/java/org/apache/iceberg/FilesTable.java b/core/src/main/java/org/apache/iceberg/FilesTable.java index 077cc679bf96..54a0f17d0007 100644 --- a/core/src/main/java/org/apache/iceberg/FilesTable.java +++ b/core/src/main/java/org/apache/iceberg/FilesTable.java @@ -19,7 +19,7 @@ package org.apache.iceberg; -import java.util.List; +import org.apache.iceberg.io.CloseableIterable; /** * A {@link Table} implementation that exposes a table's files as rows. @@ -60,8 +60,8 @@ protected TableScan newRefinedScan(TableOperations ops, Table table, Schema sche } @Override - protected List manifests() { - return snapshot().allManifests(); + protected CloseableIterable manifests() { + return CloseableIterable.withNoopClose(snapshot().allManifests()); } } } diff --git a/core/src/test/java/org/apache/iceberg/TestMetadataTableFilters.java b/core/src/test/java/org/apache/iceberg/TestMetadataTableFilters.java index a724f9f2db04..81210abf307e 100644 --- a/core/src/test/java/org/apache/iceberg/TestMetadataTableFilters.java +++ b/core/src/test/java/org/apache/iceberg/TestMetadataTableFilters.java @@ -48,7 +48,9 @@ public static Object[][] parameters() { { MetadataTableType.DATA_FILES, 2 }, { MetadataTableType.DELETE_FILES, 2 }, { MetadataTableType.FILES, 1 }, - { MetadataTableType.FILES, 2 } + { MetadataTableType.FILES, 2 }, + { MetadataTableType.ALL_DATA_FILES, 1 }, + { MetadataTableType.ALL_DATA_FILES, 2 } }; } @@ -88,6 +90,14 @@ public void setupTable() throws Exception { .addDeletes(FILE_D2_DELETES) .commit(); } + + if (type.equals(MetadataTableType.ALL_DATA_FILES)) { + // Clear all files from current snapshot to test whether 'all' Files tables scans previous files + table.newDelete().deleteFromRowFilter(Expressions.alwaysTrue()).commit(); // Moves file entries to DELETED state + table.newDelete().deleteFromRowFilter(Expressions.alwaysTrue()).commit(); // Removes all entries + Assert.assertEquals("Current snapshot should be made empty", + 0, table.currentSnapshot().allManifests().size()); + } } private Table createMetadataTable() { @@ -98,6 +108,8 @@ private Table createMetadataTable() { return new DataFilesTable(table.ops(), table); case DELETE_FILES: return new DeleteFilesTable(table.ops(), table); + case ALL_DATA_FILES: + return new AllDataFilesTable(table.ops(), table); default: throw new IllegalArgumentException("Unsupported metadata table type:" + type); } @@ -114,6 +126,8 @@ private int expectedScanTaskCount(int partitions) { case DATA_FILES: case DELETE_FILES: return partitions; + case ALL_DATA_FILES: + return partitions * 2; // ScanTask for Data Manifest in DELETED and ADDED states default: throw new IllegalArgumentException("Unsupported metadata table type:" + type); } diff --git a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java index 4579eb0ae413..8959fcffdafa 100644 --- a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java +++ b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java @@ -130,7 +130,7 @@ public void testPartitionedTable() throws Exception { List recordsB = Lists.newArrayList( new SimpleRecord(1, "b"), new SimpleRecord(2, "b") - ); + ); spark.createDataset(recordsB, Encoders.bean(SimpleRecord.class)) .coalesce(1) .writeTo(tableName) @@ -181,6 +181,94 @@ public void testPartitionedTable() throws Exception { TestHelpers.assertEqualsSafe(filesTableSchema.asStruct(), expectedFiles.get(1), actualFiles.get(1)); } + @Test + public void testAllFiles() throws Exception { + sql("CREATE TABLE %s (id bigint, data string) USING iceberg TBLPROPERTIES" + + "('format-version'='2', 'write.delete.mode'='merge-on-read')", tableName); + + List records = Lists.newArrayList( + new SimpleRecord(1, "a"), + new SimpleRecord(2, "b"), + new SimpleRecord(3, "c"), + new SimpleRecord(4, "d") + ); + spark.createDataset(records, Encoders.bean(SimpleRecord.class)) + .coalesce(1) + .writeTo(tableName) + .append(); + + + Table table = Spark3Util.loadIcebergTable(spark, tableName); + List expectedDataManifests = TestHelpers.dataManifests(table); + Assert.assertEquals("Should have 2 data files", 1, expectedDataManifests.size()); + + // Clear table to test whether 'all_files' can read past files + List results = sql("DELETE FROM %s", tableName); + Assert.assertEquals("Table should be cleared", 0, results.size()); + + Schema entriesTableSchema = Spark3Util.loadIcebergTable(spark, tableName + ".entries").schema(); + Schema filesTableSchema = Spark3Util.loadIcebergTable(spark, tableName + ".all_data_files").schema(); + + List actualDataFiles = spark.sql("SELECT * FROM " + tableName + ".all_data_files").collectAsList(); + + List expectedDataFiles = expectedEntries(table, FileContent.DATA, + entriesTableSchema, expectedDataManifests, null); + + Assert.assertEquals("Should be one data file manifest entry", 1, expectedDataFiles.size()); + Assert.assertEquals("Metadata table should return one data file", 1, actualDataFiles.size()); + + TestHelpers.assertEqualsSafe(filesTableSchema.asStruct(), expectedDataFiles.get(0), actualDataFiles.get(0)); + } + + @Test + public void testAllFilesPartitioned() throws Exception { + sql("CREATE TABLE %s (id bigint, data string) " + + "USING iceberg " + + "PARTITIONED BY (data) " + + "TBLPROPERTIES" + + "('format-version'='2', 'write.delete.mode'='merge-on-read')", tableName); + + List recordsA = Lists.newArrayList( + new SimpleRecord(1, "a"), + new SimpleRecord(2, "a") + ); + spark.createDataset(recordsA, Encoders.bean(SimpleRecord.class)) + .coalesce(1) + .writeTo(tableName) + .append(); + + List recordsB = Lists.newArrayList( + new SimpleRecord(1, "b"), + new SimpleRecord(2, "b") + ); + spark.createDataset(recordsB, Encoders.bean(SimpleRecord.class)) + .coalesce(1) + .writeTo(tableName) + .append(); + + Table table = Spark3Util.loadIcebergTable(spark, tableName); + List expectedDataManifests = TestHelpers.dataManifests(table); + Assert.assertEquals("Should have 2 data files", 2, expectedDataManifests.size()); + + // Clear table to test whether 'all_files' can read past files + List results = sql("DELETE FROM %s", tableName); + Assert.assertEquals("Table should be cleared", 0, results.size()); + + List actualDataFiles = spark.sql("SELECT * FROM " + tableName + ".all_data_files " + + "WHERE partition.data='a'").collectAsList(); + + Schema entriesTableSchema = Spark3Util.loadIcebergTable(spark, tableName + ".entries").schema(); + Schema filesTableSchema = Spark3Util.loadIcebergTable(spark, tableName + ".all_data_files").schema(); + + List expectedDataFiles = expectedEntries(table, FileContent.DATA, + entriesTableSchema, expectedDataManifests, "a"); + + Assert.assertEquals("Should be one data file manifest entry", 1, expectedDataFiles.size()); + Assert.assertEquals("Metadata table should return one data file", 1, actualDataFiles.size()); + + TestHelpers.assertEqualsSafe(filesTableSchema.asStruct(), expectedDataFiles.get(0), actualDataFiles.get(0)); + } + /** * Find matching manifest entries of an Iceberg table * @param table iceberg table From 2608ae3cc5697d83034de41cc7b8f19f94a6132a Mon Sep 17 00:00:00 2001 From: Szehon Ho Date: Wed, 23 Mar 2022 17:53:19 -0700 Subject: [PATCH 2/7] Fix WAP case --- .../org/apache/iceberg/BaseFilesTable.java | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index 48199bbd9991..f6aafafe57d3 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java @@ -20,6 +20,8 @@ package org.apache.iceberg; import java.util.Map; +import org.apache.iceberg.events.Listeners; +import org.apache.iceberg.events.ScanEvent; import org.apache.iceberg.expressions.Expression; import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.expressions.ManifestEvaluator; @@ -32,6 +34,8 @@ import org.apache.iceberg.relocated.com.google.common.collect.Sets; import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.types.Types.StructType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * Base class logic for files metadata tables @@ -55,6 +59,8 @@ public Schema schema() { } abstract static class BaseFilesTableScan extends BaseMetadataTableScan { + private static final Logger LOG = LoggerFactory.getLogger(BaseFilesTableScan.class); + private final Schema fileSchema; private final MetadataTableType type; @@ -87,6 +93,18 @@ public TableScan appendsAfter(long fromSnapshotId) { String.format("Cannot incrementally scan table of type %s", type.name())); } + @Override + public CloseableIterable planFiles() { + if (type.equals(MetadataTableType.ALL_DATA_FILES)) { + LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); + Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); + + return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); + } else { + return super.planFiles(); + } + } + @Override protected CloseableIterable planFiles(TableOperations ops, Snapshot snapshot, Expression rowFilter, boolean ignoreResiduals, boolean caseSensitive, boolean colStats) { From b0deb5516290c6682970fe6e025d9bcb4e7b2935 Mon Sep 17 00:00:00 2001 From: Szehon Ho Date: Thu, 24 Mar 2022 14:04:10 -0700 Subject: [PATCH 3/7] Refactor to make WAP-handling slightly more clear --- .../org/apache/iceberg/AllDataFilesTable.java | 5 +++++ .../iceberg/BaseAllMetadataTableScan.java | 10 +--------- .../java/org/apache/iceberg/BaseFilesTable.java | 17 ----------------- .../apache/iceberg/BaseMetadataTableScan.java | 16 ++++++++++++++++ 4 files changed, 22 insertions(+), 26 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java index c13dc2b289ce..6ee1a2693ff3 100644 --- a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java @@ -79,6 +79,11 @@ public TableScan asOfTime(long timestampMillis) { throw new UnsupportedOperationException("Cannot select snapshot: all_data_files is for all snapshots"); } + @Override + public CloseableIterable planFiles() { + return super.planAllFiles(); + } + @Override protected CloseableIterable manifests() { try (CloseableIterable iterable = new ParallelIterable<>( diff --git a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java index 5b7a5365ba44..70b49cea1a57 100644 --- a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java @@ -19,14 +19,9 @@ package org.apache.iceberg; -import org.apache.iceberg.events.Listeners; -import org.apache.iceberg.events.ScanEvent; import org.apache.iceberg.io.CloseableIterable; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; abstract class BaseAllMetadataTableScan extends BaseMetadataTableScan { - private static final Logger LOG = LoggerFactory.getLogger(BaseAllMetadataTableScan.class); BaseAllMetadataTableScan(TableOperations ops, Table table, Schema fileSchema) { super(ops, table, fileSchema); @@ -58,9 +53,6 @@ public TableScan appendsAfter(long fromSnapshotId) { @Override public CloseableIterable planFiles() { - LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); - Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); - - return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); + return super.planAllFiles(); } } diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index f6aafafe57d3..f999a021cf1c 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java @@ -20,8 +20,6 @@ package org.apache.iceberg; import java.util.Map; -import org.apache.iceberg.events.Listeners; -import org.apache.iceberg.events.ScanEvent; import org.apache.iceberg.expressions.Expression; import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.expressions.ManifestEvaluator; @@ -34,8 +32,6 @@ import org.apache.iceberg.relocated.com.google.common.collect.Sets; import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.types.Types.StructType; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * Base class logic for files metadata tables @@ -59,7 +55,6 @@ public Schema schema() { } abstract static class BaseFilesTableScan extends BaseMetadataTableScan { - private static final Logger LOG = LoggerFactory.getLogger(BaseFilesTableScan.class); private final Schema fileSchema; private final MetadataTableType type; @@ -93,18 +88,6 @@ public TableScan appendsAfter(long fromSnapshotId) { String.format("Cannot incrementally scan table of type %s", type.name())); } - @Override - public CloseableIterable planFiles() { - if (type.equals(MetadataTableType.ALL_DATA_FILES)) { - LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); - Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); - - return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); - } else { - return super.planFiles(); - } - } - @Override protected CloseableIterable planFiles(TableOperations ops, Snapshot snapshot, Expression rowFilter, boolean ignoreResiduals, boolean caseSensitive, boolean colStats) { diff --git a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java index b5db3e5a8035..42aa5338659c 100644 --- a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java @@ -19,9 +19,15 @@ package org.apache.iceberg; +import org.apache.iceberg.events.Listeners; +import org.apache.iceberg.events.ScanEvent; +import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.util.PropertyUtil; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; abstract class BaseMetadataTableScan extends BaseTableScan { + private static final Logger LOG = LoggerFactory.getLogger(BaseMetadataTableScan.class); protected BaseMetadataTableScan(TableOperations ops, Table table, Schema schema) { super(ops, table, schema); @@ -38,4 +44,14 @@ public long targetSplitSize() { TableProperties.METADATA_SPLIT_SIZE_DEFAULT); return PropertyUtil.propertyAsLong(options(), TableProperties.SPLIT_SIZE, tableValue); } + + /** + * Alternative to {@link #planFiles()}, allows exploring old snapshots even for empty table. + */ + protected CloseableIterable planAllFiles() { + LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); + Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); + + return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); + } } From 1db465ee2d0359b3d83c5eff673fedd58c766e41 Mon Sep 17 00:00:00 2001 From: Szehon Ho Date: Fri, 25 Mar 2022 11:05:15 -0700 Subject: [PATCH 4/7] Add allScan() flag in BaseMetadataTableScan and restore the planFiles override --- .../org/apache/iceberg/AllDataFilesTable.java | 10 ++++---- .../iceberg/BaseAllMetadataTableScan.java | 6 ++--- .../org/apache/iceberg/BaseFilesTable.java | 2 +- .../apache/iceberg/BaseMetadataTableScan.java | 23 +++++++++++++------ 4 files changed, 24 insertions(+), 17 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java index 6ee1a2693ff3..b3d298f55fb7 100644 --- a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java @@ -64,6 +64,11 @@ private AllDataFilesTableScan(TableOperations ops, Table table, Schema schema, S super(ops, table, schema, fileSchema, context, MetadataTableType.ALL_DATA_FILES); } + @Override + public boolean allScan() { + return true; + } + @Override protected TableScan newRefinedScan(TableOperations ops, Table table, Schema schema, TableScanContext context) { return new AllDataFilesTableScan(ops, table, schema, fileSchema(), context); @@ -79,11 +84,6 @@ public TableScan asOfTime(long timestampMillis) { throw new UnsupportedOperationException("Cannot select snapshot: all_data_files is for all snapshots"); } - @Override - public CloseableIterable planFiles() { - return super.planAllFiles(); - } - @Override protected CloseableIterable manifests() { try (CloseableIterable iterable = new ParallelIterable<>( diff --git a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java index 70b49cea1a57..201fa499e74c 100644 --- a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java @@ -19,8 +19,6 @@ package org.apache.iceberg; -import org.apache.iceberg.io.CloseableIterable; - abstract class BaseAllMetadataTableScan extends BaseMetadataTableScan { BaseAllMetadataTableScan(TableOperations ops, Table table, Schema fileSchema) { @@ -52,7 +50,7 @@ public TableScan appendsAfter(long fromSnapshotId) { } @Override - public CloseableIterable planFiles() { - return super.planAllFiles(); + protected boolean allScan() { + return true; } } diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index f999a021cf1c..a1034bd0fade 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java @@ -108,7 +108,7 @@ protected CloseableIterable planFiles(TableOperations ops, Snapsho } /** - * @return list of manifest files to explore for this files metadata table scan + * @return iterable of manifest files to explore for this files metadata table scan */ protected abstract CloseableIterable manifests(); diff --git a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java index 42aa5338659c..10b314837089 100644 --- a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java @@ -37,6 +37,13 @@ protected BaseMetadataTableScan(TableOperations ops, Table table, Schema schema, super(ops, table, schema, context); } + /** + * @return if metadata table scan is for all snapshots, ie 'all_x' metadata tables + */ + protected boolean allScan() { + return false; + } + @Override public long targetSplitSize() { long tableValue = tableOps().current().propertyAsLong( @@ -45,13 +52,15 @@ public long targetSplitSize() { return PropertyUtil.propertyAsLong(options(), TableProperties.SPLIT_SIZE, tableValue); } - /** - * Alternative to {@link #planFiles()}, allows exploring old snapshots even for empty table. - */ - protected CloseableIterable planAllFiles() { - LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); - Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); + @Override + public CloseableIterable planFiles() { + if (allScan()) { // Avoid returning early for empty tables + LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); + Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); - return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); + return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); + } else { + return super.planFiles(); + } } } From b6b5596ad16a2af684f641816f8b77a1993cbc0d Mon Sep 17 00:00:00 2001 From: Szehon Ho Date: Mon, 4 Apr 2022 14:40:41 -0700 Subject: [PATCH 5/7] Revert "Add allScan() flag in BaseMetadataTableScan and restore the planFiles override" This reverts commit 1db465ee2d0359b3d83c5eff673fedd58c766e41. --- .../org/apache/iceberg/AllDataFilesTable.java | 10 ++++---- .../iceberg/BaseAllMetadataTableScan.java | 6 +++-- .../org/apache/iceberg/BaseFilesTable.java | 2 +- .../apache/iceberg/BaseMetadataTableScan.java | 23 ++++++------------- 4 files changed, 17 insertions(+), 24 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java index b3d298f55fb7..6ee1a2693ff3 100644 --- a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java @@ -64,11 +64,6 @@ private AllDataFilesTableScan(TableOperations ops, Table table, Schema schema, S super(ops, table, schema, fileSchema, context, MetadataTableType.ALL_DATA_FILES); } - @Override - public boolean allScan() { - return true; - } - @Override protected TableScan newRefinedScan(TableOperations ops, Table table, Schema schema, TableScanContext context) { return new AllDataFilesTableScan(ops, table, schema, fileSchema(), context); @@ -84,6 +79,11 @@ public TableScan asOfTime(long timestampMillis) { throw new UnsupportedOperationException("Cannot select snapshot: all_data_files is for all snapshots"); } + @Override + public CloseableIterable planFiles() { + return super.planAllFiles(); + } + @Override protected CloseableIterable manifests() { try (CloseableIterable iterable = new ParallelIterable<>( diff --git a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java index 201fa499e74c..70b49cea1a57 100644 --- a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java @@ -19,6 +19,8 @@ package org.apache.iceberg; +import org.apache.iceberg.io.CloseableIterable; + abstract class BaseAllMetadataTableScan extends BaseMetadataTableScan { BaseAllMetadataTableScan(TableOperations ops, Table table, Schema fileSchema) { @@ -50,7 +52,7 @@ public TableScan appendsAfter(long fromSnapshotId) { } @Override - protected boolean allScan() { - return true; + public CloseableIterable planFiles() { + return super.planAllFiles(); } } diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index a1034bd0fade..f999a021cf1c 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java @@ -108,7 +108,7 @@ protected CloseableIterable planFiles(TableOperations ops, Snapsho } /** - * @return iterable of manifest files to explore for this files metadata table scan + * @return list of manifest files to explore for this files metadata table scan */ protected abstract CloseableIterable manifests(); diff --git a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java index 10b314837089..42aa5338659c 100644 --- a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java @@ -37,13 +37,6 @@ protected BaseMetadataTableScan(TableOperations ops, Table table, Schema schema, super(ops, table, schema, context); } - /** - * @return if metadata table scan is for all snapshots, ie 'all_x' metadata tables - */ - protected boolean allScan() { - return false; - } - @Override public long targetSplitSize() { long tableValue = tableOps().current().propertyAsLong( @@ -52,15 +45,13 @@ public long targetSplitSize() { return PropertyUtil.propertyAsLong(options(), TableProperties.SPLIT_SIZE, tableValue); } - @Override - public CloseableIterable planFiles() { - if (allScan()) { // Avoid returning early for empty tables - LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); - Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); + /** + * Alternative to {@link #planFiles()}, allows exploring old snapshots even for empty table. + */ + protected CloseableIterable planAllFiles() { + LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); + Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); - return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); - } else { - return super.planFiles(); - } + return planFiles(tableOps(), snapshot(), filter(), shouldIgnoreResiduals(), isCaseSensitive(), colStats()); } } From 5843e6f8308a74975ec81fece331213a99821110 Mon Sep 17 00:00:00 2001 From: Szehon Ho Date: Mon, 4 Apr 2022 14:59:57 -0700 Subject: [PATCH 6/7] Address review comments --- core/src/main/java/org/apache/iceberg/AllDataFilesTable.java | 5 ++--- .../java/org/apache/iceberg/BaseAllMetadataTableScan.java | 2 +- core/src/main/java/org/apache/iceberg/BaseFilesTable.java | 4 ++-- .../main/java/org/apache/iceberg/BaseMetadataTableScan.java | 4 ++-- .../apache/iceberg/spark/extensions/TestMetadataTables.java | 4 ++-- 5 files changed, 9 insertions(+), 10 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java index 6ee1a2693ff3..36f331637181 100644 --- a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java @@ -81,14 +81,13 @@ public TableScan asOfTime(long timestampMillis) { @Override public CloseableIterable planFiles() { - return super.planAllFiles(); + return super.planFilesAllSnapshots(); } @Override protected CloseableIterable manifests() { try (CloseableIterable iterable = new ParallelIterable<>( - Iterables.transform(table().snapshots(), - snapshot -> (Iterable) () -> snapshot.dataManifests().iterator()), + Iterables.transform(table().snapshots(), snapshot -> () -> snapshot.dataManifests().iterator()), context().planExecutor())) { return CloseableIterable.withNoopClose(Sets.newHashSet(iterable)); } catch (IOException e) { diff --git a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java index 70b49cea1a57..e6df849a663c 100644 --- a/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java @@ -53,6 +53,6 @@ public TableScan appendsAfter(long fromSnapshotId) { @Override public CloseableIterable planFiles() { - return super.planAllFiles(); + return super.planFilesAllSnapshots(); } } diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index f999a021cf1c..410023c98198 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java @@ -90,7 +90,7 @@ public TableScan appendsAfter(long fromSnapshotId) { @Override protected CloseableIterable planFiles(TableOperations ops, Snapshot snapshot, Expression rowFilter, - boolean ignoreResiduals, boolean caseSensitive, boolean colStats) { + boolean ignoreResiduals, boolean caseSensitive, boolean colStats) { CloseableIterable filtered = filterManifests(manifests(), rowFilter, caseSensitive); String schemaString = SchemaParser.toJson(schema()); @@ -108,7 +108,7 @@ protected CloseableIterable planFiles(TableOperations ops, Snapsho } /** - * @return list of manifest files to explore for this files metadata table scan + * Returns an iterable of manifest files to explore for this Files metadata table scan */ protected abstract CloseableIterable manifests(); diff --git a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java index 42aa5338659c..59c75f382726 100644 --- a/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java +++ b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java @@ -46,9 +46,9 @@ public long targetSplitSize() { } /** - * Alternative to {@link #planFiles()}, allows exploring old snapshots even for empty table. + * Alternative to {@link #planFiles()}, allows exploring old snapshots even for an empty table. */ - protected CloseableIterable planAllFiles() { + protected CloseableIterable planFilesAllSnapshots() { LOG.info("Scanning metadata table {} with filter {}.", table(), filter()); Listeners.notifyAll(new ScanEvent(table().name(), 0L, filter(), schema())); diff --git a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java index 8959fcffdafa..968ed79dfed6 100644 --- a/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java +++ b/spark/v3.2/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestMetadataTables.java @@ -200,7 +200,7 @@ public void testAllFiles() throws Exception { Table table = Spark3Util.loadIcebergTable(spark, tableName); List expectedDataManifests = TestHelpers.dataManifests(table); - Assert.assertEquals("Should have 2 data files", 1, expectedDataManifests.size()); + Assert.assertEquals("Should have 1 data manifest", 1, expectedDataManifests.size()); // Clear table to test whether 'all_files' can read past files List results = sql("DELETE FROM %s", tableName); @@ -248,7 +248,7 @@ public void testAllFilesPartitioned() throws Exception { Table table = Spark3Util.loadIcebergTable(spark, tableName); List expectedDataManifests = TestHelpers.dataManifests(table); - Assert.assertEquals("Should have 2 data files", 2, expectedDataManifests.size()); + Assert.assertEquals("Should have 2 data manifests", 2, expectedDataManifests.size()); // Clear table to test whether 'all_files' can read past files List results = sql("DELETE FROM %s", tableName); From 51ba2305356101c74be8590354a767600b5f3e35 Mon Sep 17 00:00:00 2001 From: Szehon Ho Date: Mon, 4 Apr 2022 15:09:32 -0700 Subject: [PATCH 7/7] Fix build --- core/src/main/java/org/apache/iceberg/AllDataFilesTable.java | 3 ++- core/src/main/java/org/apache/iceberg/BaseFilesTable.java | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java index 36f331637181..a352c9f369ae 100644 --- a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java @@ -87,7 +87,8 @@ public CloseableIterable planFiles() { @Override protected CloseableIterable manifests() { try (CloseableIterable iterable = new ParallelIterable<>( - Iterables.transform(table().snapshots(), snapshot -> () -> snapshot.dataManifests().iterator()), + Iterables.transform(table().snapshots(), + snapshot -> (Iterable) () -> snapshot.dataManifests().iterator()), context().planExecutor())) { return CloseableIterable.withNoopClose(Sets.newHashSet(iterable)); } catch (IOException e) { diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index 410023c98198..a2c7f11078fb 100644 --- a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java +++ b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java @@ -90,7 +90,8 @@ public TableScan appendsAfter(long fromSnapshotId) { @Override protected CloseableIterable planFiles(TableOperations ops, Snapshot snapshot, Expression rowFilter, - boolean ignoreResiduals, boolean caseSensitive, boolean colStats) { + boolean ignoreResiduals, boolean caseSensitive, + boolean colStats) { CloseableIterable filtered = filterManifests(manifests(), rowFilter, caseSensitive); String schemaString = SchemaParser.toJson(schema());