diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java index 5c5f934ea205..04feea5037d1 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java @@ -22,6 +22,7 @@ import com.google.auto.value.AutoValue; import java.util.Arrays; import java.util.List; +import java.util.Map; import org.apache.beam.sdk.annotations.Internal; import org.apache.beam.sdk.io.Read; import org.apache.beam.sdk.schemas.Schema; @@ -406,6 +407,8 @@ public abstract static class WriteRows extends PTransform, Iceb abstract boolean getAutoSharding(); + abstract @Nullable Map getWriteProperties(); + abstract Builder toBuilder(); @AutoValue.Builder @@ -424,6 +427,8 @@ abstract static class Builder { abstract Builder setAutoSharding(boolean autoSharding); + abstract Builder setWriteProperties(Map writeProperties); + abstract WriteRows build(); } @@ -474,6 +479,10 @@ public WriteRows withAutosharding() { return toBuilder().setAutoSharding(true).build(); } + public WriteRows withWriteProperties(Map writeProperties) { + return toBuilder().setWriteProperties(writeProperties).build(); + } + @Override public IcebergWriteResult expand(PCollection input) { List allToArgs = Arrays.asList(getTableIdentifier(), getDynamicDestinations()); @@ -509,7 +518,8 @@ public IcebergWriteResult expand(PCollection input) { getCatalogConfig(), destinations, getTriggeringFrequency(), - getDirectWriteByteLimit())); + getDirectWriteByteLimit(), + getWriteProperties())); case HASH: return input .apply( @@ -521,7 +531,8 @@ public IcebergWriteResult expand(PCollection input) { getCatalogConfig(), destinations, getTriggeringFrequency(), - getAutoSharding())); + getAutoSharding(), + getWriteProperties())); default: throw new UnsupportedOperationException( "Unsupported distribution mode: " + getDistributionMode()); diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java index 8db4fb77a8e8..0ae9d5fb0ecb 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java @@ -158,6 +158,11 @@ public static Builder builder() { + "during high-throughput writes. Only available with 'hash' distribution mode.") public abstract @Nullable Boolean getAutosharding(); + @SchemaFieldDescription( + "Properties applied to the underlying file writer (e.g. Parquet write properties like " + + "'write.parquet.bloom-filter-enabled.column.').") + public abstract @Nullable Map getWriteProperties(); + @AutoValue.Builder public abstract static class Builder { public abstract Builder setTable(String table); @@ -188,6 +193,8 @@ public abstract static class Builder { public abstract Builder setAutosharding(Boolean autosharding); + public abstract Builder setWriteProperties(Map writeProperties); + public abstract Configuration build(); } @@ -279,6 +286,11 @@ public PCollectionRowTuple expand(PCollectionRowTuple input) { writeTransform = writeTransform.withAutosharding(); } + @Nullable Map writeProperties = configuration.getWriteProperties(); + if (writeProperties != null && !writeProperties.isEmpty()) { + writeTransform = writeTransform.withWriteProperties(writeProperties); + } + // TODO: support dynamic destinations IcebergWriteResult result = rows.apply(writeTransform); diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriter.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriter.java index fd3d5d63327c..c3b63b2a336f 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriter.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriter.java @@ -18,6 +18,7 @@ package org.apache.beam.sdk.io.iceberg; import java.io.IOException; +import java.util.Map; import org.apache.beam.sdk.metrics.Counter; import org.apache.beam.sdk.metrics.Metrics; import org.apache.iceberg.DataFile; @@ -34,6 +35,7 @@ import org.apache.iceberg.io.DataWriter; import org.apache.iceberg.io.OutputFile; import org.apache.iceberg.parquet.Parquet; +import org.checkerframework.checker.nullness.qual.Nullable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -54,11 +56,22 @@ class RecordWriter { catalog.loadTable(destination.getTableIdentifier()), destination.getFileFormat(), filename, - partitionKey); + partitionKey, + null); } RecordWriter(Table table, FileFormat fileFormat, String filename, StructLike partitionKey) throws IOException { + this(table, fileFormat, filename, partitionKey, null); + } + + RecordWriter( + Table table, + FileFormat fileFormat, + String filename, + StructLike partitionKey, + @Nullable Map writeProperties) + throws IOException { this.table = table; this.fileFormat = fileFormat; @@ -91,14 +104,17 @@ class RecordWriter { .build(); break; case PARQUET: - icebergDataWriter = + Parquet.DataWriteBuilder parquetBuilder = Parquet.writeData(outputFile) .forTable(table) .createWriterFunc(GenericParquetWriter::create) .withPartition(partitionKey) .withKeyMetadata(keyMetadata) - .overwrite() - .build(); + .overwrite(); + if (writeProperties != null && !writeProperties.isEmpty()) { + parquetBuilder.setAll(writeProperties); + } + icebergDataWriter = parquetBuilder.build(); break; case ORC: throw new UnsupportedOperationException("ORC file format not currently supported."); diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java index d5439ad83831..6893c743f431 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java @@ -202,7 +202,8 @@ private RecordWriter createWriter(PartitionKey partitionKey) { table, icebergDestination.getFileFormat(), filePrefix + "_" + stateToken + "_" + recordIndex, - partitionKey); + partitionKey, + writeProperties); openWriters++; return writer; } catch (IOException e) { @@ -253,6 +254,7 @@ static String getPartitionDataPath( private final String filePrefix; private final long maxFileSize; private final int maxNumWriters; + private final @Nullable Map writeProperties; @VisibleForTesting int openWriters = 0; @VisibleForTesting @@ -265,10 +267,20 @@ static String getPartitionDataPath( RecordWriterManager( IcebergCatalogConfig catalogConfig, String filePrefix, long maxFileSize, int maxNumWriters) { + this(catalogConfig, filePrefix, maxFileSize, maxNumWriters, null); + } + + RecordWriterManager( + IcebergCatalogConfig catalogConfig, + String filePrefix, + long maxFileSize, + int maxNumWriters, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.filePrefix = filePrefix; this.maxFileSize = maxFileSize; this.maxNumWriters = maxNumWriters; + this.writeProperties = writeProperties; } /** diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteDirectRowsToFiles.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteDirectRowsToFiles.java index 5cf095dc3c30..e03085e6be78 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteDirectRowsToFiles.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteDirectRowsToFiles.java @@ -40,16 +40,19 @@ class WriteDirectRowsToFiles private final IcebergCatalogConfig catalogConfig; private final String filePrefix; private final long maxBytesPerFile; + private final @Nullable Map writeProperties; WriteDirectRowsToFiles( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, String filePrefix, - long maxBytesPerFile) { + long maxBytesPerFile, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filePrefix = filePrefix; this.maxBytesPerFile = maxBytesPerFile; + this.writeProperties = writeProperties; } @Override @@ -57,7 +60,7 @@ public PCollection expand(PCollection> input) { return input.apply( ParDo.of( new WriteDirectRowsToFilesDoFn( - catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix))); + catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix, writeProperties))); } private static class WriteDirectRowsToFilesDoFn extends DoFn, FileWriteResult> { @@ -66,24 +69,28 @@ private static class WriteDirectRowsToFilesDoFn extends DoFn, Fi private final IcebergCatalogConfig catalogConfig; private final String filePrefix; private final long maxFileSize; + private final @Nullable Map writeProperties; private transient @Nullable RecordWriterManager recordWriterManager; WriteDirectRowsToFilesDoFn( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, long maxFileSize, - String filePrefix) { + String filePrefix, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filePrefix = filePrefix; this.maxFileSize = maxFileSize; + this.writeProperties = writeProperties; this.recordWriterManager = null; } @StartBundle public void startBundle() { recordWriterManager = - new RecordWriterManager(catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE); + new RecordWriterManager( + catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE, writeProperties); } @ProcessElement diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteGroupedRowsToFiles.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteGroupedRowsToFiles.java index b16496240d18..e74715a7eebc 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteGroupedRowsToFiles.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteGroupedRowsToFiles.java @@ -18,6 +18,7 @@ package org.apache.beam.sdk.io.iceberg; import java.util.List; +import java.util.Map; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.PTransform; import org.apache.beam.sdk.transforms.ParDo; @@ -30,6 +31,7 @@ import org.apache.beam.sdk.values.WindowedValue; import org.apache.beam.sdk.values.WindowedValues; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions; +import org.checkerframework.checker.nullness.qual.Nullable; class WriteGroupedRowsToFiles extends PTransform< @@ -39,16 +41,19 @@ class WriteGroupedRowsToFiles private final DynamicDestinations dynamicDestinations; private final IcebergCatalogConfig catalogConfig; private final String filePrefix; + private final @Nullable Map writeProperties; WriteGroupedRowsToFiles( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, String filePrefix, - long maxBytesPerFile) { + long maxBytesPerFile, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filePrefix = filePrefix; this.maxBytesPerFile = maxBytesPerFile; + this.writeProperties = writeProperties; } @Override @@ -57,7 +62,7 @@ public PCollection expand( return input.apply( ParDo.of( new WriteGroupedRowsToFilesDoFn( - catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix))); + catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix, writeProperties))); } private static class WriteGroupedRowsToFilesDoFn @@ -67,16 +72,19 @@ private static class WriteGroupedRowsToFilesDoFn private final IcebergCatalogConfig catalogConfig; private final String filePrefix; private final long maxFileSize; + private final @Nullable Map writeProperties; WriteGroupedRowsToFilesDoFn( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, long maxFileSize, - String filePrefix) { + String filePrefix, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filePrefix = filePrefix; this.maxFileSize = maxFileSize; + this.writeProperties = writeProperties; } @ProcessElement @@ -93,7 +101,8 @@ public void processElement( WindowedValues.of(destination, window.maxTimestamp(), window, paneInfo); RecordWriterManager writer; try (RecordWriterManager openWriter = - new RecordWriterManager(catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE)) { + new RecordWriterManager( + catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE, writeProperties)) { writer = openWriter; for (Row e : element.getValue()) { writer.write(windowedDestination, e); diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java index 8b4ae0863f72..84e7d05b124a 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java @@ -59,14 +59,17 @@ class WritePartitionedRowsToFiles private final DynamicDestinations dynamicDestinations; private final IcebergCatalogConfig catalogConfig; private final String filePrefix; + private final @Nullable Map writeProperties; WritePartitionedRowsToFiles( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, - String filePrefix) { + String filePrefix, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filePrefix = filePrefix; + this.writeProperties = writeProperties; } @Override @@ -78,7 +81,9 @@ public PCollection expand(PCollection>> i .getElemCoder()) .getSchema(); return input.apply( - ParDo.of(new WriteDoFn(catalogConfig, dynamicDestinations, filePrefix, dataSchema))); + ParDo.of( + new WriteDoFn( + catalogConfig, dynamicDestinations, filePrefix, dataSchema, writeProperties))); } private static class WriteDoFn extends DoFn>, FileWriteResult> { @@ -87,6 +92,7 @@ private static class WriteDoFn extends DoFn>, FileWriteRes private final IcebergCatalogConfig catalogConfig; private final String filePrefix; private final Schema dataSchema; + private final @Nullable Map writeProperties; private transient @MonotonicNonNull Map specIds; private transient @MonotonicNonNull Map> partitionFieldMaps; @@ -95,11 +101,13 @@ private static class WriteDoFn extends DoFn>, FileWriteRes IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, String filePrefix, - Schema dataSchema) { + Schema dataSchema, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filePrefix = filePrefix; this.dataSchema = dataSchema; + this.writeProperties = writeProperties; } @Setup @@ -132,7 +140,8 @@ public void processElement( .addExtension(String.format("%s-%s", filePrefix, UUID.randomUUID())); RecordWriter writer = - new RecordWriter(table, destination.getFileFormat(), fileName, partitionData); + new RecordWriter( + table, destination.getFileFormat(), fileName, partitionData, writeProperties); try { for (Row row : element.getValue()) { Record record = IcebergUtils.beamRowToIcebergRecord(table.schema(), row); diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteToDestinations.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteToDestinations.java index bea84fc826b7..684ef350a20f 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteToDestinations.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteToDestinations.java @@ -59,16 +59,19 @@ class WriteToDestinations extends PTransform>, Icebe private final @Nullable Duration triggeringFrequency; private final String filePrefix; private final @Nullable Integer directWriteByteLimit; + private final @Nullable Map writeProperties; WriteToDestinations( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, @Nullable Duration triggeringFrequency, - @Nullable Integer directWriteByteLimit) { + @Nullable Integer directWriteByteLimit, + @Nullable Map writeProperties) { this.dynamicDestinations = dynamicDestinations; this.catalogConfig = catalogConfig; this.triggeringFrequency = triggeringFrequency; this.directWriteByteLimit = directWriteByteLimit; + this.writeProperties = writeProperties; // single unique prefix per write transform this.filePrefix = UUID.randomUUID().toString(); } @@ -112,7 +115,11 @@ private PCollection groupAndWriteRecords(PCollection applyUserTriggering(PCollection input) { @@ -157,7 +164,11 @@ private PCollection writeTriggeredWithBundleLifting( largeBatches.apply( "WriteDirectRowsToFiles", new WriteDirectRowsToFiles( - catalogConfig, dynamicDestinations, filePrefix, DEFAULT_MAX_BYTES_PER_FILE)); + catalogConfig, + dynamicDestinations, + filePrefix, + DEFAULT_MAX_BYTES_PER_FILE, + writeProperties)); PCollection groupedFileWrites = groupAndWriteRecords(smallBatches); @@ -189,7 +200,11 @@ private PCollection writeUntriggered(PCollection writeGroupedResult = @@ -199,7 +214,11 @@ private PCollection writeUntriggered(PCollection>, IcebergWri private final @Nullable Duration triggeringFrequency; private final String filePrefix; private final boolean autoSharding; + private final @Nullable Map writeProperties; WriteToPartitions( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, @Nullable Duration triggeringFrequency, - boolean autoSharding) { + boolean autoSharding, + @Nullable Map writeProperties) { this.dynamicDestinations = dynamicDestinations; this.catalogConfig = catalogConfig; this.triggeringFrequency = triggeringFrequency; // single unique prefix per write transform this.filePrefix = UUID.randomUUID().toString(); this.autoSharding = autoSharding; + this.writeProperties = writeProperties; } private PCollection>> groupByPartition(PCollection> input) { @@ -95,7 +99,8 @@ public IcebergWriteResult expand(PCollection> input) { PCollection writtenFiles = groupedRows.apply( - new WritePartitionedRowsToFiles(catalogConfig, dynamicDestinations, filePrefix)); + new WritePartitionedRowsToFiles( + catalogConfig, dynamicDestinations, filePrefix, writeProperties)); if (IcebergUtils.isUnbounded(input) && triggeringFrequency != null) { writtenFiles = diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteUngroupedRowsToFiles.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteUngroupedRowsToFiles.java index ff9a98c62005..7c780e6395df 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteUngroupedRowsToFiles.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteUngroupedRowsToFiles.java @@ -72,16 +72,19 @@ class WriteUngroupedRowsToFiles private final DynamicDestinations dynamicDestinations; private final IcebergCatalogConfig catalogConfig; private final long maxBytesPerFile; + private final @Nullable Map writeProperties; WriteUngroupedRowsToFiles( IcebergCatalogConfig catalogConfig, DynamicDestinations dynamicDestinations, String filePrefix, - long maxBytesPerFile) { + long maxBytesPerFile, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filePrefix = filePrefix; this.maxBytesPerFile = maxBytesPerFile; + this.writeProperties = writeProperties; } @Override @@ -95,7 +98,8 @@ public Result expand(PCollection> input) { dynamicDestinations, filePrefix, DEFAULT_MAX_WRITERS_PER_BUNDLE, - maxBytesPerFile)) + maxBytesPerFile, + writeProperties)) .withOutputTags( WRITTEN_FILES_TAG, TupleTagList.of(ImmutableList.of(WRITTEN_ROWS_TAG, SPILLED_ROWS_TAG)))); @@ -191,6 +195,7 @@ private static class WriteUngroupedRowsToFilesDoFn private final long maxFileSize; private final DynamicDestinations dynamicDestinations; private final IcebergCatalogConfig catalogConfig; + private final @Nullable Map writeProperties; private transient @Nullable RecordWriterManager recordWriterManager; private int spilledShardNumber; @@ -199,18 +204,21 @@ public WriteUngroupedRowsToFilesDoFn( DynamicDestinations dynamicDestinations, String filename, int maximumWritersPerBundle, - long maxFileSize) { + long maxFileSize, + @Nullable Map writeProperties) { this.catalogConfig = catalogConfig; this.dynamicDestinations = dynamicDestinations; this.filename = filename; this.maxWritersPerBundle = maximumWritersPerBundle; this.maxFileSize = maxFileSize; + this.writeProperties = writeProperties; } @StartBundle public void startBundle() { recordWriterManager = - new RecordWriterManager(catalogConfig, filename, maxFileSize, maxWritersPerBundle); + new RecordWriterManager( + catalogConfig, filename, maxFileSize, maxWritersPerBundle, writeProperties); this.spilledShardNumber = ThreadLocalRandom.current().nextInt(SPILLED_RECORD_SHARDING_FACTOR); } diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/RecordWriterManagerTest.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/RecordWriterManagerTest.java index 390e8d87af28..821fb2ac7b24 100644 --- a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/RecordWriterManagerTest.java +++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/RecordWriterManagerTest.java @@ -53,6 +53,8 @@ import org.apache.beam.sdk.values.WindowedValues; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; import org.apache.commons.lang3.RandomStringUtils; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; import org.apache.iceberg.AppendFiles; import org.apache.iceberg.DataFile; import org.apache.iceberg.FileFormat; @@ -77,6 +79,10 @@ import org.apache.iceberg.types.Type; import org.apache.iceberg.types.Types; import org.apache.iceberg.util.DateTimeUtil; +import org.apache.parquet.hadoop.ParquetFileReader; +import org.apache.parquet.hadoop.metadata.BlockMetaData; +import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData; +import org.apache.parquet.hadoop.util.HadoopInputFile; import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.DateTime; import org.joda.time.DateTimeZone; @@ -1300,4 +1306,67 @@ public void testFileIOSurvivesAcrossBundles() throws IOException { assertTrue( "Bundle 2 should produce data files", bundle2.getSerializableDataFiles().containsKey(dest)); } + + @Test + public void testWritePropertiesAppliedToParquetFiles() throws IOException { + Schema bloomSchema = + Schema.builder().addInt32Field("colWithBf").addInt32Field("colWithoutBf").build(); + org.apache.iceberg.Schema icebergBloomSchema = + IcebergUtils.beamSchemaToIcebergSchema(bloomSchema); + + TableIdentifier tableId = TableIdentifier.of("default", "test_write_properties"); + warehouse.createTable(tableId, icebergBloomSchema); + + Map writeProperties = + ImmutableMap.of( + "write.parquet.bloom-filter-enabled.column.colWithBf", "true", + "write.parquet.bloom-filter-enabled.column.colWithoutBf", "false"); + + IcebergDestination destination = + IcebergDestination.builder() + .setTableIdentifier(tableId) + .setFileFormat(FileFormat.PARQUET) + .build(); + WindowedValue dest = WindowedValues.valueInGlobalWindow(destination); + + RecordWriterManager writerManager = + new RecordWriterManager(catalogConfig, "test_bloom", Long.MAX_VALUE, 3, writeProperties); + for (int i = 0; i < 10; i++) { + Row row = Row.withSchema(bloomSchema).addValues(i, 100 + i).build(); + assertTrue(writerManager.write(dest, row)); + } + writerManager.close(); + + List dataFiles = writerManager.getSerializableDataFiles().get(dest); + assertEquals(1, dataFiles.size()); + + String dataFilePath = dataFiles.get(0).getPath(); + assertNotNull(dataFilePath); + + try (ParquetFileReader reader = + ParquetFileReader.open( + HadoopInputFile.fromPath(new Path(dataFilePath), new Configuration()))) { + List blocks = reader.getFooter().getBlocks(); + assertFalse("Parquet file should have at least one row group", blocks.isEmpty()); + + for (int i = 0; i < blocks.size(); i++) { + BlockMetaData block = blocks.get(i); + assertEquals("Each row group should have 2 columns", 2, block.getColumns().size()); + + for (ColumnChunkMetaData col : block.getColumns()) { + boolean hasBloomFilter = col.getBloomFilterOffset() > 0; + String colName = col.getPath().toDotString(); + if (colName.equals("colWithBf")) { + assertTrue( + "Column 'colWithBf' in row group " + i + " should have a bloom filter", + hasBloomFilter); + } else if (colName.equals("colWithoutBf")) { + assertFalse( + "Column 'colWithoutBf' in row group " + i + " should not have a bloom filter", + hasBloomFilter); + } + } + } + } + } }