diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteOrphanFilesSparkAction.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteOrphanFilesSparkAction.java index a79f075ef442..72b6268026ab 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteOrphanFilesSparkAction.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteOrphanFilesSparkAction.java @@ -60,9 +60,9 @@ import org.slf4j.LoggerFactory; /** - * An action that removes orphan metadata and data files by listing a given location and comparing - * the actual files in that location with data and metadata files referenced by all valid snapshots. - * The location must be accessible for listing via the Hadoop {@link FileSystem}. + * An action that removes orphan metadata, data and delete files by listing a given location and + * comparing the actual files in that location with content and metadata files referenced by all + * valid snapshots. The location must be accessible for listing via the Hadoop {@link FileSystem}. * *

By default, this action cleans up the table location returned by {@link Table#location()} and * removes unreachable files that are older than 3 days using {@link Table#io()}. The behavior can @@ -169,7 +169,7 @@ private String jobDesc() { } private DeleteOrphanFiles.Result doExecute() { - Dataset validDataFileDF = buildValidDataFileDF(table); + Dataset validDataFileDF = buildValidContentFileDF(table); Dataset validMetadataFileDF = buildValidMetadataFileDF(table); Dataset validFileDF = validDataFileDF.union(validMetadataFileDF); Dataset actualFileDF = buildActualFileDF(); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteReachableFilesSparkAction.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteReachableFilesSparkAction.java index 1431ae5d78ec..a1bc19d7dcc0 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteReachableFilesSparkAction.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseDeleteReachableFilesSparkAction.java @@ -60,7 +60,7 @@ public class BaseDeleteReachableFilesSparkAction private static final Logger LOG = LoggerFactory.getLogger(BaseDeleteReachableFilesSparkAction.class); - private static final String DATA_FILE = "Data File"; + private static final String CONTENT_FILE = "Content File"; private static final String MANIFEST = "Manifest"; private static final String MANIFEST_LIST = "Manifest List"; private static final String OTHERS = "Others"; @@ -140,7 +140,7 @@ private Dataset projectFilePathWithType(Dataset ds, String type) { private Dataset buildValidFileDF(TableMetadata metadata) { Table staticTable = newStaticTable(metadata, io); - return projectFilePathWithType(buildValidDataFileDF(staticTable), DATA_FILE) + return projectFilePathWithType(buildValidContentFileDF(staticTable), CONTENT_FILE) .union(projectFilePathWithType(buildManifestFileDF(staticTable), MANIFEST)) .union(projectFilePathWithType(buildManifestListDF(staticTable), MANIFEST_LIST)) .union(projectFilePathWithType(buildOtherMetadataFileDF(staticTable), OTHERS)); @@ -183,9 +183,9 @@ private BaseDeleteReachableFilesActionResult deleteFiles(Iterator deleted) String type = fileInfo.getString(1); removeFunc.accept(file); switch (type) { - case DATA_FILE: + case CONTENT_FILE: dataFileCount.incrementAndGet(); - LOG.trace("Deleted Data File: {}", file); + LOG.trace("Deleted Content File: {}", file); break; case MANIFEST: manifestCount.incrementAndGet(); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseExpireSnapshotsSparkAction.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseExpireSnapshotsSparkAction.java index 2e1f0c079eca..da9907fe325a 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseExpireSnapshotsSparkAction.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseExpireSnapshotsSparkAction.java @@ -57,7 +57,7 @@ * *

This action first leverages {@link org.apache.iceberg.ExpireSnapshots} to expire snapshots and * then uses metadata tables to find files that can be safely deleted. This is done by anti-joining - * two Datasets that contain all manifest and data files before and after the expiration. The + * two Datasets that contain all manifest and content files before and after the expiration. The * snapshot expiration will be fully committed before any deletes are issued. * *

This operation performs a shuffle so the parallelism can be controlled through @@ -72,7 +72,7 @@ public class BaseExpireSnapshotsSparkAction public static final String STREAM_RESULTS = "stream-results"; - private static final String DATA_FILE = "Data File"; + private static final String CONTENT_FILE = "Content File"; private static final String MANIFEST = "Manifest"; private static final String MANIFEST_LIST = "Manifest List"; @@ -233,7 +233,7 @@ private Dataset appendTypeString(Dataset ds, String type) { private Dataset buildValidFileDF(TableMetadata metadata) { Table staticTable = newStaticTable(metadata, this.table.io()); - return appendTypeString(buildValidDataFileDF(staticTable), DATA_FILE) + return appendTypeString(buildValidContentFileDF(staticTable), CONTENT_FILE) .union(appendTypeString(buildManifestFileDF(staticTable), MANIFEST)) .union(appendTypeString(buildManifestListDF(staticTable), MANIFEST_LIST)); } @@ -266,9 +266,9 @@ private BaseExpireSnapshotsActionResult deleteFiles(Iterator expired) { String type = fileInfo.getString(1); deleteFunc.accept(file); switch (type) { - case DATA_FILE: + case CONTENT_FILE: dataFileCount.incrementAndGet(); - LOG.trace("Deleted Data File: {}", file); + LOG.trace("Deleted Content File: {}", file); break; case MANIFEST: manifestCount.incrementAndGet(); diff --git a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseSparkAction.java b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseSparkAction.java index c9d93ce9de5f..5abfcc4482a4 100644 --- a/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseSparkAction.java +++ b/spark/v3.1/spark/src/main/java/org/apache/iceberg/spark/actions/BaseSparkAction.java @@ -110,7 +110,8 @@ protected Table newStaticTable(TableMetadata metadata, FileIO io) { return new BaseTable(ops, metadataFileLocation); } - protected Dataset buildValidDataFileDF(Table table) { + // builds a DF of delete and data file locations by reading all manifests + protected Dataset buildValidContentFileDF(Table table) { JavaSparkContext context = JavaSparkContext.fromSparkContext(spark.sparkContext()); Broadcast ioBroadcast = context.broadcast(SparkUtil.serializableFileIO(table));