Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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 @@ -162,9 +162,9 @@ private String jobDesc() {
}

private DeleteOrphanFiles.Result doExecute() {
Dataset<Row> validDataFileDF = buildValidDataFileDF(table);
Dataset<Row> validContentFileDF = buildValidContentFileDF(table);
Dataset<Row> validMetadataFileDF = buildValidMetadataFileDF(table);
Dataset<Row> validFileDF = validDataFileDF.union(validMetadataFileDF);
Dataset<Row> validFileDF = validContentFileDF.union(validMetadataFileDF);
Dataset<Row> actualFileDF = buildActualFileDF();

Column actualFileName = filenameUDF.apply(actualFileDF.col("file_path"));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ public class BaseDeleteReachableFilesSparkAction
extends BaseSparkAction<DeleteReachableFiles, DeleteReachableFiles.Result> implements DeleteReachableFiles {
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";
Expand Down Expand Up @@ -138,7 +138,7 @@ private Dataset<Row> projectFilePathWithType(Dataset<Row> ds, String type) {

private Dataset<Row> 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));
Expand Down Expand Up @@ -177,9 +177,9 @@ private BaseDeleteReachableFilesActionResult deleteFiles(Iterator<Row> 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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ public class BaseExpireSnapshotsSparkAction
extends BaseSparkAction<ExpireSnapshots, ExpireSnapshots.Result> implements ExpireSnapshots {
private static final Logger LOG = LoggerFactory.getLogger(BaseExpireSnapshotsSparkAction.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";

Expand Down Expand Up @@ -226,7 +226,7 @@ private Dataset<Row> appendTypeString(Dataset<Row> ds, String type) {

private Dataset<Row> 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));
}
Expand Down Expand Up @@ -255,9 +255,9 @@ private BaseExpireSnapshotsActionResult deleteFiles(Iterator<Row> 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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,8 @@ protected Table newStaticTable(TableMetadata metadata, FileIO io) {
return new BaseTable(ops, metadataFileLocation);
}

protected Dataset<Row> buildValidDataFileDF(Table table) {
// builds a DF of delete and data file locations by reading all manifests
protected Dataset<Row> buildValidContentFileDF(Table table) {
JavaSparkContext context = JavaSparkContext.fromSparkContext(spark.sparkContext());
Broadcast<FileIO> ioBroadcast = context.broadcast(SparkUtil.serializableFileIO(table));

Expand Down