nssalian commented on code in PR #17469:
URL: https://github.com/apache/iceberg/pull/17469#discussion_r3698130168
##########
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:
I dropped the footer-schema check since reading those groups already
requires the columns to exist. But I can use
`ParquetReader.Builder.withDecryption`, thanks for the pointer.
--
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]