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
25 changes: 25 additions & 0 deletions core/src/main/java/org/apache/iceberg/BaseDistributedDataScan.java
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import org.apache.iceberg.expressions.Projections;
import org.apache.iceberg.expressions.ResidualEvaluator;
import org.apache.iceberg.io.CloseableIterable;
import org.apache.iceberg.io.InputFile;
import org.apache.iceberg.metrics.ScanMetricsUtil;
import org.apache.iceberg.relocated.com.google.common.collect.Iterables;
import org.apache.iceberg.relocated.com.google.common.collect.Maps;
Expand Down Expand Up @@ -144,6 +145,10 @@ protected PlanningMode deletePlanningMode() {

@Override
protected CloseableIterable<ScanTask> doPlanFiles() {
if (TableUtil.formatVersion(table()) >= TableMetadata.MIN_FORMAT_VERSION_PARQUET_MANIFESTS) {
return doPlanFilesV4();
}

Snapshot snapshot = snapshot();

List<ManifestFile> deleteManifests = findMatchingDeleteManifests(snapshot);
Expand Down Expand Up @@ -187,6 +192,26 @@ protected CloseableIterable<ScanTask> doPlanFiles() {
}
}

private CloseableIterable<ScanTask> doPlanFilesV4() {
Snapshot snapshot = snapshot();
InputFile rootManifest = table().io().newInputFile(snapshot.manifestListLocation());
ScanTaskPlanner.Builder planner =
ScanTaskPlanner.builder(table().io(), rootManifest, specs(), table().location())
.filterData(filter())
.caseSensitive(isCaseSensitive())
.scanMetrics(scanMetrics());

if (shouldIgnoreResiduals()) {
planner = planner.ignoreResiduals();
}

if (shouldPlanWithExecutor()) {
planner = planner.planWith(planExecutor());
}

return CloseableIterable.transform(planner.build().planFiles(), task -> (ScanTask) task);
}

@Override
public CloseableIterable<ScanTaskGroup<ScanTask>> planTasks() {
return TableScanUtil.planTaskGroups(
Expand Down
8 changes: 8 additions & 0 deletions core/src/main/java/org/apache/iceberg/BaseSnapshot.java
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,14 @@ private void cacheManifests(FileIO fileIO) {
throw new IllegalArgumentException("Cannot cache changes: FileIO is null");
}

if (allManifests == null
&& manifestListLocation != null
&& FileFormat.fromFileName(manifestListLocation) == FileFormat.PARQUET) {
// A v4 flat-tree root manifest stores data entries inline rather than referencing leaf
// manifests, so it exposes no ManifestFiles. Scans plan directly from the root manifest.
this.allManifests = ImmutableList.of();
}

if (allManifests == null && v1ManifestLocations != null) {
// if we have a collection of manifest locations, then we need to load them here
allManifests =
Expand Down
29 changes: 29 additions & 0 deletions core/src/main/java/org/apache/iceberg/DataTableScan.java
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import java.util.List;
import org.apache.iceberg.io.CloseableIterable;
import org.apache.iceberg.io.FileIO;
import org.apache.iceberg.io.InputFile;
import org.apache.iceberg.relocated.com.google.common.base.Preconditions;

public class DataTableScan extends BaseTableScan {
Expand Down Expand Up @@ -62,6 +63,34 @@ protected TableScan newRefinedScan(Table table, Schema schema, TableScanContext

@Override
public CloseableIterable<FileScanTask> doPlanFiles() {
if (TableUtil.formatVersion(table()) >= TableMetadata.MIN_FORMAT_VERSION_PARQUET_MANIFESTS) {
return doPlanFilesV4();
}

return doPlanFilesV3();
}

private CloseableIterable<FileScanTask> doPlanFilesV4() {
Snapshot snapshot = snapshot();
InputFile rootManifest = table().io().newInputFile(snapshot.manifestListLocation());
ScanTaskPlanner.Builder planner =
ScanTaskPlanner.builder(table().io(), rootManifest, specs(), table().location())
.filterData(filter())
.caseSensitive(isCaseSensitive())
.scanMetrics(scanMetrics());

if (shouldIgnoreResiduals()) {
planner = planner.ignoreResiduals();
}

if (shouldPlanWithExecutor()) {
planner = planner.planWith(planExecutor());
}

return planner.build().planFiles();
}

private CloseableIterable<FileScanTask> doPlanFilesV3() {
Snapshot snapshot = snapshot();

FileIO io = table().io();
Expand Down
Loading
Loading