diff --git a/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java b/core/src/main/java/org/apache/iceberg/AllDataFilesTable.java index 4eca0b430dd4..a352c9f369ae 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,20 @@ 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)); + public CloseableIterable planFiles() { + return super.planFilesAllSnapshots(); } - } - 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"); + @Override + 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/BaseAllMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseAllMetadataTableScan.java index 5b7a5365ba44..e6df849a663c 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.planFilesAllSnapshots(); } } diff --git a/core/src/main/java/org/apache/iceberg/BaseFilesTable.java b/core/src/main/java/org/apache/iceberg/BaseFilesTable.java index a0b835a4abcf..a2c7f11078fb 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; @@ -56,6 +55,7 @@ public Schema schema() { } abstract static class BaseFilesTableScan extends BaseMetadataTableScan { + private final Schema fileSchema; private final MetadataTableType type; @@ -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()); @@ -108,15 +109,13 @@ 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 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 +123,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/BaseMetadataTableScan.java b/core/src/main/java/org/apache/iceberg/BaseMetadataTableScan.java index b5db3e5a8035..59c75f382726 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 an empty table. + */ + protected CloseableIterable planFilesAllSnapshots() { + 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()); + } } 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..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 @@ -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 1 data manifest", 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 manifests", 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