wombatu-kun commented on code in PR #17469:
URL: https://github.com/apache/iceberg/pull/17469#discussion_r3697945539
##########
parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java:
##########
@@ -607,6 +612,130 @@ protected int resolveColumnIndex(Void engineSchema,
String columnName) {
}
}
+ @Test
+ public void testFormatModelVariantShreddingWithEncryption() throws
IOException {
+ Schema variantSchema =
+ new Schema(
+ Types.NestedField.required(1, "id", Types.LongType.get()),
+ Types.NestedField.optional(2, "v", Types.VariantType.get()));
+
+ VariantShreddingAnalyzer<Record, Void> analyzer =
Review Comment:
The schema, the anonymous VariantShreddingAnalyzer, the variant fixture and
the ParquetFormatModel.create wiring here are a verbatim copy of
testFormatModelVariantShreddingRoundTrip, differing only in the literal values.
Extract them into a shared helper that both tests call, so this test carries
only the encryption-specific setup.
##########
parquet/src/test/java/org/apache/iceberg/parquet/TestParquetDataWriter.java:
##########
@@ -607,6 +612,130 @@ protected int resolveColumnIndex(Void engineSchema,
String columnName) {
}
}
+ @Test
+ public void testFormatModelVariantShreddingWithEncryption() throws
IOException {
+ Schema variantSchema =
+ new Schema(
+ Types.NestedField.required(1, "id", Types.LongType.get()),
+ Types.NestedField.optional(2, "v", Types.VariantType.get()));
+
+ VariantShreddingAnalyzer<Record, Void> analyzer =
+ new VariantShreddingAnalyzer<Record, Void>() {
+ @Override
+ protected List<VariantValue> extractVariantValues(List<Record> rows,
int idx) {
+ List<VariantValue> values = Lists.newArrayList();
+ for (Record row : rows) {
+ Object obj = row.get(idx);
+ if (obj instanceof Variant) {
+ values.add(((Variant) obj).value());
+ }
+ }
+ return values;
+ }
+
+ @Override
+ protected int resolveColumnIndex(Void engineSchema, String
columnName) {
+ return
variantSchema.columns().indexOf(variantSchema.findField(columnName));
+ }
+ };
+
+ ByteBuffer metadataBuffer =
VariantTestUtil.createMetadata(ImmutableList.of("a", "b"), true);
+ VariantMetadata metadata = Variants.metadata(metadataBuffer);
+ ByteBuffer objectBuffer =
+ VariantTestUtil.createObject(
+ metadataBuffer,
+ ImmutableMap.of("a", Variants.of(123456789), "b",
Variants.of("string")));
+ Variant variant = Variant.of(metadata, Variants.value(metadata,
objectBuffer));
+
+ GenericRecord record = GenericRecord.create(variantSchema);
+ List<Record> variantRecords =
+ ImmutableList.of(
+ record.copy(ImmutableMap.of("id", 1L, "v", variant)),
+ record.copy(ImmutableMap.of("id", 2L, "v", variant)),
+ record.copy(ImmutableMap.of("id", 3L, "v", variant)));
+
+ ParquetFormatModel<Record, Void, ParquetValueReader<?>> model =
+ ParquetFormatModel.create(
+ Record.class,
+ Void.class,
+ (icebergSchema, messageType, engineSchema) ->
+ GenericParquetWriter.create(icebergSchema, messageType),
+ (icebergSchema, fileSchema, engineSchema, idToConstant) ->
+ GenericParquetReaders.buildReader(icebergSchema, fileSchema),
+ analyzer,
+ (Function<Void, UnaryOperator<Record>>) unused -> input -> input);
+
+ OutputFile encryptedFile = Files.localOutput(createTempFile(temp));
+ ByteBuffer fileDek = ByteBuffer.allocate(16);
+ ByteBuffer aadPrefix = ByteBuffer.allocate(16);
+ SecureRandom random = new SecureRandom();
+ random.nextBytes(fileDek.array());
+ random.nextBytes(aadPrefix.array());
+
+ try (FileAppender<Record> appender =
+ model
+ .writeBuilder(EncryptedFiles.plainAsEncryptedOutput(encryptedFile))
+ .schema(variantSchema)
+ .withFileEncryptionKey(fileDek)
+ .withAADPrefix(aadPrefix)
+ .setAll(
+ ImmutableMap.of(
+ TableProperties.PARQUET_SHRED_VARIANTS, "true",
+ TableProperties.PARQUET_VARIANT_BUFFER_SIZE, "2"))
+ .content(FileContent.DATA)
+ .build()) {
+ assertThat(appender).isInstanceOf(BufferedFileAppender.class);
+ appender.addAll(variantRecords);
+ }
+
+ assertThatThrownBy(
+ () ->
+ Parquet.read(encryptedFile.toInputFile())
+ .project(variantSchema)
+ .createReaderFunc(
+ fileSchema ->
GenericParquetReaders.buildReader(variantSchema, fileSchema))
+ .build()
+ .iterator())
+ .isInstanceOf(ParquetCryptoRuntimeException.class)
+ .hasMessage("Trying to read file with encrypted footer. No keys
available");
+
+ List<Record> writtenRecords;
+ try (CloseableIterable<Record> reader =
+ Parquet.read(encryptedFile.toInputFile())
+ .project(variantSchema)
+ .withFileEncryptionKey(fileDek)
+ .withAADPrefix(aadPrefix)
+ .createReaderFunc(
+ fileSchema -> GenericParquetReaders.buildReader(variantSchema,
fileSchema))
+ .build()) {
+ writtenRecords = Lists.newArrayList(reader);
+ }
+
+ assertThat(writtenRecords).hasSameSizeAs(variantRecords);
+ for (int i = 0; i < variantRecords.size(); i++) {
+ InternalTestHelpers.assertEquals(
+ variantSchema.asStruct(), variantRecords.get(i),
writtenRecords.get(i));
+ }
+
+ try (ParquetFileReader fileReader =
+ ParquetFileReader.open(
+ ParquetIO.file(encryptedFile.toInputFile()),
+ ParquetReadOptions.builder()
+ .withDecryption(
+ FileDecryptionProperties.builder()
+ .withFooterKey(fileDek.array())
+ .withAADPrefix(aadPrefix.array())
+ .build())
+ .build())) {
+ GroupType variantType =
+
fileReader.getFooter().getFileMetaData().getSchema().getType("v").asGroupType();
+ assertThat(variantType.containsField("typed_value")).isTrue();
+ GroupType typedValue = variantType.getType("typed_value").asGroupType();
+ assertThat(typedValue.containsField("a")).isTrue();
Review Comment:
testFormatModelVariantShreddingRoundTrip also asserts that value is empty
and that typed_value.a/b hold the values, while this one asserts only the
footer schema. ParquetReader.Builder.withDecryption accepts
FileDecryptionProperties, so those raw-group checks work on the encrypted file
too - was dropping them intentional?
--
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]