JandyTenedora commented on code in PR #16960:
URL: https://github.com/apache/iceberg/pull/16960#discussion_r3602495871
##########
core/src/test/java/org/apache/iceberg/TestRewriteTablePathUtil.java:
##########
@@ -376,4 +392,136 @@ private ManifestFile
deleteManifestWithLiveAndDeletedEntry(DeleteFile live, Dele
return writer.toManifestFile();
}
+
+ @TestTemplate
+ void rewriteDVFileRewritesReferencedDataFileInBlobMetadata() throws
IOException {
+ assumeThat(formatVersion).as("DVs require format version
3+").isGreaterThanOrEqualTo(3);
+
+ String sourcePrefix = temp.resolve("source").toString();
+ String targetPrefix = temp.resolve("target").toString();
+ String sourceDataFile = sourcePrefix + "/data/file-a.parquet";
+ String externalDataFile =
temp.resolve("external/data/file-b.parquet").toString();
+
+ PositionDeleteIndex sourceDeletes = positionDeleteIndex(1L, 4L);
+ PositionDeleteIndex externalDeletes = positionDeleteIndex(2L, 5L, 8L);
+ byte[] sourcePayload = serializedDV(sourceDeletes);
+ byte[] externalPayload = serializedDV(externalDeletes);
+
+ OutputFile sourceDVFile =
+ Files.localOutput(
+ temp.resolve("source/metadata/dv-" + System.nanoTime() +
".puffin").toString());
+ try (PuffinWriter writer =
Puffin.write(sourceDVFile).createdBy("test").build()) {
+ writer.write(newDVBlob(sourceDeletes, sourcePayload, sourceDataFile));
+ writer.write(newDVBlob(externalDeletes, externalPayload,
externalDataFile));
+ }
+
+ List<BlobMetadata> sourceBlobMetadata;
+ try (PuffinReader reader =
Puffin.read(sourceDVFile.toInputFile()).build()) {
+ sourceBlobMetadata = reader.fileMetadata().blobs();
+ }
+ assertThat(sourceBlobMetadata).hasSize(2);
+
+ DeleteFile dvDeleteFile =
+ FileMetadata.deleteFileBuilder(table.spec())
+ .ofPositionDeletes()
+ .withFormat(FileFormat.PUFFIN)
+ .withPath(sourceDVFile.location())
+ .withFileSizeInBytes(sourceDVFile.toInputFile().getLength())
+ .withPartition(FILE_A.partition())
+ .withRecordCount(sourceDeletes.cardinality())
+ .withReferencedDataFile(sourceDataFile)
+ .withContentOffset(sourceBlobMetadata.get(0).offset())
+ .withContentSizeInBytes(sourceBlobMetadata.get(0).length())
+ .build();
+
+ OutputFile rewrittenDVFile =
+ Files.localOutput(
+ temp.resolve("target/metadata/dv-rewritten-" + System.nanoTime() +
".puffin")
+ .toString());
+ RewriteTablePathUtil.rewritePositionDeleteFile(
+ dvDeleteFile, rewrittenDVFile, table.io(), table.spec(), sourcePrefix,
targetPrefix, null);
+
+ try (PuffinReader reader =
Puffin.read(rewrittenDVFile.toInputFile()).build()) {
+ List<BlobMetadata> rewrittenBlobMetadata = reader.fileMetadata().blobs();
+ assertThat(rewrittenBlobMetadata).hasSize(2);
+
+ BlobMetadata rewrittenSourceBlob = rewrittenBlobMetadata.get(0);
+ BlobMetadata rewrittenExternalBlob = rewrittenBlobMetadata.get(1);
+ assertDVBlobMetadata(
+ rewrittenSourceBlob, targetPrefix + "/data/file-a.parquet",
sourceDeletes.cardinality());
+ assertDVBlobMetadata(rewrittenExternalBlob, externalDataFile,
externalDeletes.cardinality());
+
+
assertThat(rewrittenSourceBlob.offset()).isEqualTo(dvDeleteFile.contentOffset());
+
assertThat(rewrittenSourceBlob.length()).isEqualTo(dvDeleteFile.contentSizeInBytes());
+ assertThat(rewrittenExternalBlob.offset())
+ .isEqualTo(rewrittenSourceBlob.offset() +
rewrittenSourceBlob.length());
+
+ List<Pair<BlobMetadata, ByteBuffer>> blobs =
+ ImmutableList.copyOf(reader.readAll(rewrittenBlobMetadata));
+
assertThat(ByteBuffers.toByteArray(blobs.get(0).second())).isEqualTo(sourcePayload);
+
assertThat(ByteBuffers.toByteArray(blobs.get(1).second())).isEqualTo(externalPayload);
+ }
+
+ DeleteFile rewrittenDVDeleteFile =
+ FileMetadata.deleteFileBuilder(table.spec())
+ .copy(dvDeleteFile)
+ .withPath(rewrittenDVFile.location())
+ .withFileSizeInBytes(rewrittenDVFile.toInputFile().getLength())
+ .withReferencedDataFile(targetPrefix + "/data/file-a.parquet")
+ .build();
+ assertDeletedPositions(DVUtil.readDV(rewrittenDVDeleteFile, table.io()),
1L, 4L);
+ }
+
+ private static Blob newDVBlob(
+ PositionDeleteIndex deletes, byte[] payload, String referencedDataFile) {
+ return new Blob(
+ StandardBlobTypes.DV_V1,
+ ImmutableList.of(MetadataColumns.ROW_POSITION.fieldId()),
+ -1L,
+ -1L,
+ ByteBuffer.wrap(payload),
+ null,
+ ImmutableMap.of(
+ REFERENCED_DATA_FILE,
+ referencedDataFile,
+ CARDINALITY,
+ String.valueOf(deletes.cardinality())));
+ }
+
+ private static PositionDeleteIndex positionDeleteIndex(long... positions) {
+ ImmutableList.Builder<Long> builder = ImmutableList.builder();
+ for (long position : positions) {
+ builder.add(position);
+ }
+
+ return
Deletes.toPositionIndex(CloseableIterable.withNoopClose(builder.build()));
+ }
+
+ private static byte[] serializedDV(PositionDeleteIndex deletes) {
+ return ByteBuffers.toByteArray(deletes.serialize());
+ }
+
+ private static void assertDVBlobMetadata(
+ BlobMetadata blobMetadata, String referencedDataFile, long cardinality) {
+ assertThat(blobMetadata.type()).isEqualTo(StandardBlobTypes.DV_V1);
+
assertThat(blobMetadata.inputFields()).containsExactly(MetadataColumns.ROW_POSITION.fieldId());
+ assertThat(blobMetadata.snapshotId()).isEqualTo(-1L);
+ assertThat(blobMetadata.sequenceNumber()).isEqualTo(-1L);
+ assertThat(blobMetadata.compressionCodec()).isNull();
+ assertThat(blobMetadata.offset()).isPositive();
+ assertThat(blobMetadata.length()).isPositive();
+ assertThat(blobMetadata.properties())
+ .containsExactlyInAnyOrderEntriesOf(
+ ImmutableMap.of(
+ REFERENCED_DATA_FILE,
+ referencedDataFile,
+ CARDINALITY,
+ String.valueOf(cardinality)));
+ }
+
+ private static void assertDeletedPositions(PositionDeleteIndex deletes,
long... positions) {
+ for (long position : positions) {
+ assertThat(deletes.isDeleted(position)).isTrue();
+ }
+ }
Review Comment:
Resolved in
[commit](https://github.com/apache/iceberg/pull/16960/commits/3b1f5a8fab90c6f2824782c3a2e28af01fd6ee69)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]