Skip to content
Merged
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
Original file line numberDiff line numberDiff line change
Expand Up@@ -40,10 +40,13 @@
import org.apache.arrow.vector.ValueVector;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericData.Record;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.FileContent;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.ManifestFile;
import org.apache.iceberg.Schema;
import org.apache.iceberg.Table;
import org.apache.iceberg.io.CloseableIterable;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.spark.SparkSchemaUtil;
import org.apache.iceberg.spark.data.vectorized.IcebergArrowColumnVector;
Expand DownExpand Up@@ -796,4 +799,9 @@ public static Dataset<Row> selectNonDerived(Dataset<Row> metadataTable) {
public static Types.StructType nonDerivedSchema(Dataset<Row> metadataTable) {
return SparkSchemaUtil.convert(TestHelpers.selectNonDerived(metadataTable).schema()).asStruct();
}

public static List<DataFile> dataFiles(Table table) {
CloseableIterable<FileScanTask> tasks = table.newScan().includeColumnStats().planFiles();
return Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
}
}
Original file line numberDiff line numberDiff line change
Expand Up@@ -38,7 +38,6 @@
import org.apache.iceberg.DataFile;
import org.apache.iceberg.DeleteFile;
import org.apache.iceberg.FileContent;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.ManifestFile;
import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.Schema;
Expand DownExpand Up@@ -1486,7 +1485,7 @@ public void testPartitionsTableLastUpdatedSnapshot() {
AvroSchemaUtil.convert(
partitionsTable.schema().findType("partition").asStructType(), "partition"));

List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
assertDataFilePartitions(dataFiles, Arrays.asList(1, 2, 2));

List<GenericData.Record> expected = Lists.newArrayList();
Expand DownExpand Up@@ -2058,11 +2057,6 @@ private long totalSizeInBytes(Iterable<DataFile> dataFiles) {
return Lists.newArrayList(dataFiles).stream().mapToLong(DataFile::fileSizeInBytes).sum();
}

private List<DataFile> dataFiles(Table table) {
CloseableIterable<FileScanTask> tasks = table.newScan().planFiles();
return Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
}

private void assertDataFilePartitions(
List<DataFile> dataFiles, List<Integer> expectedPartitionIds) {
Assert.assertEquals(
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,7 +30,6 @@
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
Expand DownExpand Up@@ -145,11 +144,11 @@ public void testDeleteFileThenMetadataDelete() throws Exception {

// Metadata Delete
Table table = Spark3Util.loadIcebergTable(spark, tableName);
Set<DataFile> dataFilesBefore = TestHelpers.dataFiles(table);
List<DataFile> dataFilesBefore = TestHelpers.dataFiles(table);

sql("DELETE FROM %s AS t WHERE t.id = 1", tableName);

Set<DataFile> dataFilesAfter = TestHelpers.dataFiles(table);
List<DataFile> dataFilesAfter = TestHelpers.dataFiles(table);
Assert.assertTrue(
"Data file should have been removed", dataFilesBefore.size() > dataFilesAfter.size());

Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -94,6 +94,7 @@
import org.apache.iceberg.spark.SparkTableUtil;
import org.apache.iceberg.spark.SparkTestBase;
import org.apache.iceberg.spark.actions.RewriteDataFilesSparkAction.RewriteExecutionContext;
import org.apache.iceberg.spark.data.TestHelpers;
import org.apache.iceberg.spark.source.ThreeColumnRecord;
import org.apache.iceberg.types.Comparators;
import org.apache.iceberg.types.Conversions;
Expand DownExpand Up@@ -253,9 +254,7 @@ public void testBinPackWithDeletes() throws Exception {
shouldHaveFiles(table, 8);
table.refresh();

CloseableIterable<FileScanTask> tasks = table.newScan().planFiles();
List<DataFile> dataFiles =
Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
int total = (int) dataFiles.stream().mapToLong(ContentFile::recordCount).sum();

RowDelta rowDelta = table.newRowDelta();
Expand DownExpand Up@@ -299,9 +298,7 @@ public void testBinPackWithDeleteAllData() {
shouldHaveFiles(table, 1);
table.refresh();

CloseableIterable<FileScanTask> tasks = table.newScan().planFiles();
List<DataFile> dataFiles =
Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
int total = (int) dataFiles.stream().mapToLong(ContentFile::recordCount).sum();

RowDelta rowDelta = table.newRowDelta();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,6 +49,8 @@
import org.apache.iceberg.ManifestFile;
import org.apache.iceberg.Schema;
import org.apache.iceberg.Table;
import org.apache.iceberg.TableScan;
import org.apache.iceberg.io.CloseableIterable;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.relocated.com.google.common.collect.Sets;
import org.apache.iceberg.relocated.com.google.common.collect.Streams;
Expand DownExpand Up@@ -781,14 +783,18 @@ public static List<ManifestFile> deleteManifests(Table table) {
return table.currentSnapshot().deleteManifests(table.io());
}

public static Set<DataFile> dataFiles(Table table) {
Set<DataFile> dataFiles = Sets.newHashSet();
public static List<DataFile> dataFiles(Table table) {
return dataFiles(table, null);
}

for (FileScanTask task : table.newScan().planFiles()) {
dataFiles.add(task.file());
public static List<DataFile> dataFiles(Table table, String branch) {
TableScan scan = table.newScan();
if (branch != null) {
scan.useRef(branch);
}

return dataFiles;
CloseableIterable<FileScanTask> tasks = scan.includeColumnStats().planFiles();
return Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
}

public static Set<DeleteFile> deleteFiles(Table table) {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -37,7 +37,6 @@
import org.apache.iceberg.AssertHelpers;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.DeleteFile;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.Files;
import org.apache.iceberg.ManifestFile;
import org.apache.iceberg.PartitionSpec;
Expand DownExpand Up@@ -1491,7 +1490,7 @@ public void testPartitionsTableLastUpdatedSnapshot() {
AvroSchemaUtil.convert(
partitionsTable.schema().findType("partition").asStructType(), "partition"));

List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
assertDataFilePartitions(dataFiles, Arrays.asList(1, 2, 2));

List<GenericData.Record> expected = Lists.newArrayList();
Expand DownExpand Up@@ -2195,11 +2194,6 @@ private long totalSizeInBytes(Iterable<DataFile> dataFiles) {
return Lists.newArrayList(dataFiles).stream().mapToLong(DataFile::fileSizeInBytes).sum();
}

private List<DataFile> dataFiles(Table table) {
CloseableIterable<FileScanTask> tasks = table.newScan().planFiles();
return Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
}

private void assertDataFilePartitions(
List<DataFile> dataFiles, List<Integer> expectedPartitionIds) {
Assert.assertEquals(
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,7 +30,6 @@
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
Expand DownExpand Up@@ -152,11 +151,11 @@ public void testDeleteFileThenMetadataDelete() throws Exception {

// Metadata Delete
Table table = Spark3Util.loadIcebergTable(spark, tableName);
Set<DataFile> dataFilesBefore = TestHelpers.dataFiles(table, branch);
List<DataFile> dataFilesBefore = TestHelpers.dataFiles(table, branch);

sql("DELETE FROM %s AS t WHERE t.id = 1", commitTarget());

Set<DataFile> dataFilesAfter = TestHelpers.dataFiles(table, branch);
List<DataFile> dataFilesAfter = TestHelpers.dataFiles(table, branch);
Assert.assertTrue(
"Data file should have been removed", dataFilesBefore.size() > dataFilesAfter.size());

Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -94,6 +94,7 @@
import org.apache.iceberg.spark.SparkTableUtil;
import org.apache.iceberg.spark.SparkTestBase;
import org.apache.iceberg.spark.actions.RewriteDataFilesSparkAction.RewriteExecutionContext;
import org.apache.iceberg.spark.data.TestHelpers;
import org.apache.iceberg.spark.source.ThreeColumnRecord;
import org.apache.iceberg.types.Comparators;
import org.apache.iceberg.types.Conversions;
Expand DownExpand Up@@ -254,9 +255,7 @@ public void testBinPackWithDeletes() throws Exception {
shouldHaveFiles(table, 8);
table.refresh();

CloseableIterable<FileScanTask> tasks = table.newScan().planFiles();
List<DataFile> dataFiles =
Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
int total = (int) dataFiles.stream().mapToLong(ContentFile::recordCount).sum();

RowDelta rowDelta = table.newRowDelta();
Expand DownExpand Up@@ -300,9 +299,7 @@ public void testBinPackWithDeleteAllData() {
shouldHaveFiles(table, 1);
table.refresh();

CloseableIterable<FileScanTask> tasks = table.newScan().planFiles();
List<DataFile> dataFiles =
Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
int total = (int) dataFiles.stream().mapToLong(ContentFile::recordCount).sum();

RowDelta rowDelta = table.newRowDelta();
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -32,7 +32,6 @@
import org.apache.iceberg.DataFile;
import org.apache.iceberg.DeleteFile;
import org.apache.iceberg.FileFormat;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.Files;
import org.apache.iceberg.MetadataTableType;
import org.apache.iceberg.MetadataTableUtils;
Expand DownExpand Up@@ -60,6 +59,7 @@
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.spark.SparkCatalogConfig;
import org.apache.iceberg.spark.SparkCatalogTestBase;
import org.apache.iceberg.spark.data.TestHelpers;
import org.apache.iceberg.spark.source.FourColumnRecord;
import org.apache.iceberg.spark.source.ThreeColumnRecord;
import org.apache.iceberg.types.Types;
Expand DownExpand Up@@ -136,7 +136,7 @@ public void testEmptyTable() {
@Test
public void testUnpartitioned() throws Exception {
Table table = createTableUnpartitioned(2, SCALE);
List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
writePosDeletesForFiles(table, 2, DELETES_SCALE, dataFiles);
Assert.assertEquals(2, dataFiles.size());

Expand DownExpand Up@@ -170,7 +170,7 @@ public void testUnpartitioned() throws Exception {
public void testRewriteAll() throws Exception {
Table table = createTablePartitioned(4, 2, SCALE);

List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
writePosDeletesForFiles(table, 2, DELETES_SCALE, dataFiles);
Assert.assertEquals(4, dataFiles.size());

Expand DownExpand Up@@ -206,7 +206,7 @@ public void testRewriteAll() throws Exception {
public void testRewriteToSmallerTarget() throws Exception {
Table table = createTablePartitioned(4, 2, SCALE);

List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
writePosDeletesForFiles(table, 2, DELETES_SCALE, dataFiles);
Assert.assertEquals(4, dataFiles.size());

Expand DownExpand Up@@ -243,7 +243,7 @@ public void testRewriteToSmallerTarget() throws Exception {
public void testRemoveDanglingDeletes() throws Exception {
Table table = createTablePartitioned(4, 2, SCALE);

List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
writePosDeletesForFiles(
table,
2,
Expand DownExpand Up@@ -288,7 +288,7 @@ public void testRemoveDanglingDeletes() throws Exception {
public void testSomePartitionsDanglingDeletes() throws Exception {
Table table = createTablePartitioned(4, 2, SCALE);

List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
writePosDeletesForFiles(table, 2, DELETES_SCALE, dataFiles);
Assert.assertEquals(4, dataFiles.size());

Expand DownExpand Up@@ -340,7 +340,7 @@ public void testSomePartitionsDanglingDeletes() throws Exception {
@Test
public void testPartitionEvolutionAdd() throws Exception {
Table table = createTableUnpartitioned(2, SCALE);
List<DataFile> unpartitionedDataFiles = dataFiles(table);
List<DataFile> unpartitionedDataFiles = TestHelpers.dataFiles(table);
writePosDeletesForFiles(table, 2, DELETES_SCALE, unpartitionedDataFiles);
Assert.assertEquals(2, unpartitionedDataFiles.size());

Expand All@@ -354,7 +354,8 @@ public void testPartitionEvolutionAdd() throws Exception {

table.updateSpec().addField("c1").commit();
writeRecords(table, 2, SCALE, 2);
List<DataFile> partitionedDataFiles = except(dataFiles(table), unpartitionedDataFiles);
List<DataFile> partitionedDataFiles =
except(TestHelpers.dataFiles(table), unpartitionedDataFiles);
writePosDeletesForFiles(table, 2, DELETES_SCALE, partitionedDataFiles);
Assert.assertEquals(2, partitionedDataFiles.size());

Expand DownExpand Up@@ -391,7 +392,7 @@ public void testPartitionEvolutionAdd() throws Exception {
@Test
public void testPartitionEvolutionRemove() throws Exception {
Table table = createTablePartitioned(2, 2, SCALE);
List<DataFile> dataFilesUnpartitioned = dataFiles(table);
List<DataFile> dataFilesUnpartitioned = TestHelpers.dataFiles(table);
writePosDeletesForFiles(table, 2, DELETES_SCALE, dataFilesUnpartitioned);
Assert.assertEquals(2, dataFilesUnpartitioned.size());

Expand All@@ -401,7 +402,8 @@ public void testPartitionEvolutionRemove() throws Exception {
table.updateSpec().removeField("c1").commit();

writeRecords(table, 2, SCALE);
List<DataFile> dataFilesPartitioned = except(dataFiles(table), dataFilesUnpartitioned);
List<DataFile> dataFilesPartitioned =
except(TestHelpers.dataFiles(table), dataFilesUnpartitioned);
writePosDeletesForFiles(table, 2, DELETES_SCALE, dataFilesPartitioned);
Assert.assertEquals(2, dataFilesPartitioned.size());

Expand DownExpand Up@@ -438,7 +440,7 @@ public void testPartitionEvolutionRemove() throws Exception {
@Test
public void testSchemaEvolution() throws Exception {
Table table = createTablePartitioned(2, 2, SCALE);
List<DataFile> dataFiles = dataFiles(table);
List<DataFile> dataFiles = TestHelpers.dataFiles(table);
writePosDeletesForFiles(table, 2, DELETES_SCALE, dataFiles);
Assert.assertEquals(2, dataFiles.size());

Expand All@@ -450,7 +452,7 @@ public void testSchemaEvolution() throws Exception {

int newColId = table.schema().findField("c4").fieldId();
List<DataFile> newSchemaDataFiles =
dataFiles(table).stream()
TestHelpers.dataFiles(table).stream()
.filter(f -> f.upperBounds().containsKey(newColId))
.collect(Collectors.toList());
writePosDeletesForFiles(table, 2, DELETES_SCALE, newSchemaDataFiles);
Expand DownExpand Up@@ -679,11 +681,6 @@ private void writePosDeletesForFiles(
}
}

private List<DataFile> dataFiles(Table table) {
CloseableIterable<FileScanTask> tasks = table.newScan().includeColumnStats().planFiles();
return Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
}

private List<DeleteFile> deleteFiles(Table table) {
Table deletesTable =
MetadataTableUtils.createMetadataTableInstance(table, MetadataTableType.POSITION_DELETES);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@
import org.apache.iceberg.Schema;
import org.apache.iceberg.Table;
import org.apache.iceberg.TableScan;
import org.apache.iceberg.io.CloseableIterable;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.relocated.com.google.common.collect.Sets;
import org.apache.iceberg.relocated.com.google.common.collect.Streams;
Expand DownExpand Up@@ -782,22 +783,18 @@ public static List<ManifestFile> deleteManifests(Table table) {
return table.currentSnapshot().deleteManifests(table.io());
}

public static Set<DataFile> dataFiles(Table table) {
public static List<DataFile> dataFiles(Table table) {
return dataFiles(table, null);
}

public static Set<DataFile> dataFiles(Table table, String branch) {
Set<DataFile> dataFiles = Sets.newHashSet();
public static List<DataFile> dataFiles(Table table, String branch) {
TableScan scan = table.newScan();
if (branch != null) {
scan.useRef(branch);
}

for (FileScanTask task : scan.planFiles()) {
dataFiles.add(task.file());
}

return dataFiles;
CloseableIterable<FileScanTask> tasks = scan.includeColumnStats().planFiles();
return Lists.newArrayList(CloseableIterable.transform(tasks, FileScanTask::file));
}

public static Set<DeleteFile> deleteFiles(Table table) {
Expand Down
Loading