From f347aa29d7768a5571c4ef6026c8a2f3ad2bd8eb Mon Sep 17 00:00:00 2001 From: Claire McGinty Date: Wed, 5 Aug 2026 16:35:55 -0400 Subject: [PATCH 1/4] (IcebergIO) bugfix: wire writeProperties through table create request, not DataWriteBuilder --- .../beam/sdk/io/iceberg/RecordWriter.java | 24 +--- .../sdk/io/iceberg/RecordWriterManager.java | 8 +- .../iceberg/WritePartitionedRowsToFiles.java | 8 +- .../sdk/io/iceberg/IcebergIOWriteTest.java | 106 ++++++++++++++++++ .../io/iceberg/RecordWriterManagerTest.java | 69 ------------ 5 files changed, 120 insertions(+), 95 deletions(-) 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 c3b63b2a336f..fd3d5d63327c 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,7 +18,6 @@ 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; @@ -35,7 +34,6 @@ 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; @@ -56,22 +54,11 @@ class RecordWriter { catalog.loadTable(destination.getTableIdentifier()), destination.getFileFormat(), filename, - partitionKey, - null); + partitionKey); } 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; @@ -104,17 +91,14 @@ class RecordWriter { .build(); break; case PARQUET: - Parquet.DataWriteBuilder parquetBuilder = + icebergDataWriter = Parquet.writeData(outputFile) .forTable(table) .createWriterFunc(GenericParquetWriter::create) .withPartition(partitionKey) .withKeyMetadata(keyMetadata) - .overwrite(); - if (writeProperties != null && !writeProperties.isEmpty()) { - parquetBuilder.setAll(writeProperties); - } - icebergDataWriter = parquetBuilder.build(); + .overwrite() + .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 6893c743f431..014475714050 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,8 +202,7 @@ private RecordWriter createWriter(PartitionKey partitionKey) { table, icebergDestination.getFileFormat(), filePrefix + "_" + stateToken + "_" + recordIndex, - partitionKey, - writeProperties); + partitionKey); openWriters++; return writer; } catch (IOException e) { @@ -311,8 +310,11 @@ private Table loadOrCreateTable(IcebergDestination destination, Schema dataSchem SortOrder sortOrder = createConfig != null ? createConfig.getSortOrder() : SortOrder.unsorted(); Map tableProperties = createConfig != null && createConfig.getTableProperties() != null - ? createConfig.getTableProperties() + ? Maps.newHashMap(createConfig.getTableProperties()) : Maps.newHashMap(); + if (writeProperties != null) { + tableProperties.putAll(writeProperties); + } // Create namespace if it does not exist yet if (!namespace.isEmpty() && catalog instanceof SupportsNamespaces) { 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 d1a08980fa9d..dc912ff1e4e1 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 @@ -141,8 +141,7 @@ public void processElement( .addExtension(String.format("%s-%s", filePrefix, UUID.randomUUID())); RecordWriter writer = - new RecordWriter( - table, destination.getFileFormat(), fileName, partitionData, writeProperties); + new RecordWriter(table, destination.getFileFormat(), fileName, partitionData); try { for (Row row : element.getValue()) { Record record = IcebergUtils.beamRowToIcebergRecord(table.schema(), row); @@ -194,8 +193,11 @@ private Table loadOrCreateTable( createConfig != null ? createConfig.getSortOrder() : SortOrder.unsorted(); Map tableProperties = createConfig != null && createConfig.getTableProperties() != null - ? createConfig.getTableProperties() + ? Maps.newHashMap(createConfig.getTableProperties()) : Maps.newHashMap(); + if (writeProperties != null) { + tableProperties.putAll(writeProperties); + } // Create namespace if it does not exist yet if (!namespace.isEmpty() && catalog instanceof SupportsNamespaces) { diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java index 384e11b761bb..2d3243d8b9e6 100644 --- a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java +++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java @@ -88,9 +88,13 @@ import org.apache.iceberg.data.IcebergGenerics; import org.apache.iceberg.data.Record; import org.apache.iceberg.data.parquet.GenericParquetWriter; +import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.io.DataWriter; import org.apache.iceberg.io.OutputFile; import org.apache.iceberg.parquet.Parquet; +import org.apache.parquet.hadoop.ParquetFileReader; +import org.apache.parquet.hadoop.metadata.BlockMetaData; +import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData; import org.hamcrest.Matchers; import org.joda.time.Duration; import org.joda.time.Instant; @@ -829,4 +833,106 @@ public void testCreateTableWithPartitionSpecAndSortOrder() { List writtenRecords = ImmutableList.copyOf(IcebergGenerics.read(table).build()); assertThat(writtenRecords, Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray())); } + + @Test + public void testWriteWithParquetProperties() throws Exception { + TableIdentifier tableId = + TableIdentifier.of( + "default", "parquet_props_" + Long.toString(UUID.randomUUID().hashCode(), 16)); + + Schema beamSchema = IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA); + + Map catalogProps = + ImmutableMap.builder() + .put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP) + .put("warehouse", warehouse.location) + .build(); + + IcebergCatalogConfig catalog = + IcebergCatalogConfig.builder() + .setCatalogName("name") + .setCatalogProperties(catalogProps) + .build(); + + testPipeline + .apply("Records To Add", Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1))) + .setRowSchema(beamSchema) + .apply( + "Append To Table", + writeTransform(catalog, tableId) + .withWriteProperties( + ImmutableMap.of("write.parquet.bloom-filter-enabled.column.data", "true"))); + + testPipeline.run().waitUntilFinish(); + + Table table = warehouse.loadTable(tableId); + + List writtenRecords = ImmutableList.copyOf(IcebergGenerics.read(table).build()); + assertThat(writtenRecords, Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray())); + + // verify bloom filter is present on 'data' column in written parquet files + try (CloseableIterable tasks = table.newScan().planFiles()) { + for (org.apache.iceberg.FileScanTask task : tasks) { + String path = task.file().path().toString(); + try (ParquetFileReader reader = + ParquetFileReader.open( + org.apache.parquet.hadoop.util.HadoopInputFile.fromPath( + new org.apache.hadoop.fs.Path(path), + new org.apache.hadoop.conf.Configuration()))) { + for (BlockMetaData block : reader.getFooter().getBlocks()) { + for (ColumnChunkMetaData col : block.getColumns()) { + boolean hasBloom = col.getBloomFilterOffset() > 0; + if (col.getPath().toDotString().equals("data")) { + assertTrue("Expected bloom filter on column 'data', but none was found", hasBloom); + } else { + assertFalse( + "Expected no bloom filter on column '" + col.getPath().toDotString() + "'", + hasBloom); + } + } + } + } + } + } + } + + @Test + public void testWriteWithTableProperties() throws Exception { + TableIdentifier tableId = + TableIdentifier.of( + "default", "table_props_" + Long.toString(UUID.randomUUID().hashCode(), 16)); + + Schema beamSchema = IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA); + + Map catalogProps = + ImmutableMap.builder() + .put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP) + .put("warehouse", warehouse.location) + .build(); + + IcebergCatalogConfig catalog = + IcebergCatalogConfig.builder() + .setCatalogName("name") + .setCatalogProperties(catalogProps) + .build(); + + testPipeline + .apply("Records To Add", Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1))) + .setRowSchema(beamSchema) + .apply( + "Append To Table", + writeTransform(catalog, tableId) + .withWriteProperties( + ImmutableMap.of("write.data.path", warehouse.location + "/custom_data_path"))); + + testPipeline.run().waitUntilFinish(); + + Table table = warehouse.loadTable(tableId); + + List writtenRecords = ImmutableList.copyOf(IcebergGenerics.read(table).build()); + assertThat(writtenRecords, Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray())); + + assertEquals( + warehouse.location + "/custom_data_path", table.properties().get("write.data.path")); + } } 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 821fb2ac7b24..390e8d87af28 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,8 +53,6 @@ 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; @@ -79,10 +77,6 @@ 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; @@ -1306,67 +1300,4 @@ 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); - } - } - } - } - } } From 1f5573e9acf47ba64cd57375371df543fcdd8c67 Mon Sep 17 00:00:00 2001 From: Claire McGinty Date: Thu, 6 Aug 2026 14:37:27 -0400 Subject: [PATCH 2/4] Test all dynamic write properties are propagated in managedio --- ...ebergWriteSchemaTransformProviderTest.java | 63 +++++++++++++++++++ 1 file changed, 63 insertions(+) diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java index 5a7aa11e10a9..bfb762f8e89d 100644 --- a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java +++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderTest.java @@ -25,6 +25,7 @@ import static org.apache.iceberg.util.DateTimeUtil.timestampFromMicros; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; import static org.junit.Assume.assumeTrue; import java.time.LocalDate; @@ -62,6 +63,8 @@ import org.apache.iceberg.CatalogUtil; import org.apache.iceberg.DistributionMode; import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.SortDirection; +import org.apache.iceberg.SortOrder; import org.apache.iceberg.Table; import org.apache.iceberg.catalog.TableIdentifier; import org.apache.iceberg.data.IcebergGenerics; @@ -693,4 +696,64 @@ public void testWriteCreateTableWithTableProperties() { assertEquals("5", table.properties().get("commit.retry.num-retries")); assertEquals("134217728", table.properties().get("read.split.target-size")); } + + @Test + public void testDynamicWriteCreateTableWithTableProperties() { + String identifier = "default.table_" + Long.toString(UUID.randomUUID().hashCode(), 16); + Schema schema = Schema.builder().addStringField("str").addInt32Field("int").build(); + + String customDataPath = warehouse.location + "/custom_data_path"; + + Map config = + ImmutableMap.of( + "table", + identifier, + "catalog_properties", + ImmutableMap.of("type", "hadoop", "warehouse", warehouse.location), + "table_properties", + ImmutableMap.of( + "write.data.path", + customDataPath, + "write.parquet.bloom-filter-enabled.column.int", + "true"), + "sort_fields", + Collections.singletonList("str desc"), + "partition_fields", + Collections.singletonList("int")); + + List rows = new ArrayList<>(); + for (int i = 0; i < 10; i++) { + Row row = Row.withSchema(schema).addValues("str_" + i, i).build(); + rows.add(row); + } + + PCollection result = + testPipeline + .apply("Records To Add", Create.of(rows)) + .setRowSchema(schema) + .apply(Managed.write(Managed.ICEBERG).withConfig(config)) + .get(SNAPSHOTS_TAG); + + PAssert.that(result) + .satisfies(new VerifyOutputs(Collections.singletonList(identifier), "append")); + testPipeline.run().waitUntilFinish(); + + Table table = warehouse.loadTable(TableIdentifier.parse(identifier)); + + PartitionSpec spec = table.spec(); + assertTrue(spec.isPartitioned()); + assertEquals(1, spec.fields().size()); + assertEquals("int", spec.fields().get(0).name()); + + SortOrder sortOrder = table.sortOrder(); + assertTrue(sortOrder.isSorted()); + assertEquals(1, sortOrder.fields().size()); + assertEquals(SortDirection.DESC, sortOrder.fields().get(0).direction()); + + assertEquals(customDataPath, table.properties().get("write.data.path")); + assertEquals("true", table.properties().get("write.parquet.bloom-filter-enabled.column.int")); + + List writtenRecords = ImmutableList.copyOf(IcebergGenerics.read(table).build()); + assertEquals(10, writtenRecords.size()); + } } From 3b6665e739dc37e7c694b831ff14bf02308db123 Mon Sep 17 00:00:00 2001 From: Claire McGinty Date: Fri, 7 Aug 2026 10:52:36 -0400 Subject: [PATCH 3/4] Revert changes and document writeProperties --- .../apache/beam/sdk/io/iceberg/IcebergIO.java | 7 ++ .../beam/sdk/io/iceberg/RecordWriter.java | 24 +++- .../sdk/io/iceberg/RecordWriterManager.java | 8 +- .../iceberg/WritePartitionedRowsToFiles.java | 8 +- .../sdk/io/iceberg/IcebergIOWriteTest.java | 106 ------------------ .../io/iceberg/RecordWriterManagerTest.java | 69 ++++++++++++ 6 files changed, 102 insertions(+), 120 deletions(-) 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 ee5755898b7f..9ee2e395421e 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 @@ -488,6 +488,13 @@ public WriteRows withAutosharding() { return toBuilder().setAutoSharding(true).build(); } + /** + * Defines properties to be passed to the Iceberg writer itself. Note that these properties are + * execution-scoped, meaning that they are applied to a preexisting table and will not mutate + * any table-level properties. + * + *

See: https://iceberg.apache.org/docs/latest/configuration/#write-properties + */ public WriteRows withWriteProperties(Map writeProperties) { return toBuilder().setWriteProperties(writeProperties).build(); } 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 014475714050..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) { @@ -310,11 +311,8 @@ private Table loadOrCreateTable(IcebergDestination destination, Schema dataSchem SortOrder sortOrder = createConfig != null ? createConfig.getSortOrder() : SortOrder.unsorted(); Map tableProperties = createConfig != null && createConfig.getTableProperties() != null - ? Maps.newHashMap(createConfig.getTableProperties()) + ? createConfig.getTableProperties() : Maps.newHashMap(); - if (writeProperties != null) { - tableProperties.putAll(writeProperties); - } // Create namespace if it does not exist yet if (!namespace.isEmpty() && catalog instanceof SupportsNamespaces) { 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 dc912ff1e4e1..d1a08980fa9d 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 @@ -141,7 +141,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); @@ -193,11 +194,8 @@ private Table loadOrCreateTable( createConfig != null ? createConfig.getSortOrder() : SortOrder.unsorted(); Map tableProperties = createConfig != null && createConfig.getTableProperties() != null - ? Maps.newHashMap(createConfig.getTableProperties()) + ? createConfig.getTableProperties() : Maps.newHashMap(); - if (writeProperties != null) { - tableProperties.putAll(writeProperties); - } // Create namespace if it does not exist yet if (!namespace.isEmpty() && catalog instanceof SupportsNamespaces) { diff --git a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java index 2d3243d8b9e6..384e11b761bb 100644 --- a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java +++ b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOWriteTest.java @@ -88,13 +88,9 @@ import org.apache.iceberg.data.IcebergGenerics; import org.apache.iceberg.data.Record; import org.apache.iceberg.data.parquet.GenericParquetWriter; -import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.io.DataWriter; import org.apache.iceberg.io.OutputFile; import org.apache.iceberg.parquet.Parquet; -import org.apache.parquet.hadoop.ParquetFileReader; -import org.apache.parquet.hadoop.metadata.BlockMetaData; -import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData; import org.hamcrest.Matchers; import org.joda.time.Duration; import org.joda.time.Instant; @@ -833,106 +829,4 @@ public void testCreateTableWithPartitionSpecAndSortOrder() { List writtenRecords = ImmutableList.copyOf(IcebergGenerics.read(table).build()); assertThat(writtenRecords, Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray())); } - - @Test - public void testWriteWithParquetProperties() throws Exception { - TableIdentifier tableId = - TableIdentifier.of( - "default", "parquet_props_" + Long.toString(UUID.randomUUID().hashCode(), 16)); - - Schema beamSchema = IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA); - - Map catalogProps = - ImmutableMap.builder() - .put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP) - .put("warehouse", warehouse.location) - .build(); - - IcebergCatalogConfig catalog = - IcebergCatalogConfig.builder() - .setCatalogName("name") - .setCatalogProperties(catalogProps) - .build(); - - testPipeline - .apply("Records To Add", Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1))) - .setRowSchema(beamSchema) - .apply( - "Append To Table", - writeTransform(catalog, tableId) - .withWriteProperties( - ImmutableMap.of("write.parquet.bloom-filter-enabled.column.data", "true"))); - - testPipeline.run().waitUntilFinish(); - - Table table = warehouse.loadTable(tableId); - - List writtenRecords = ImmutableList.copyOf(IcebergGenerics.read(table).build()); - assertThat(writtenRecords, Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray())); - - // verify bloom filter is present on 'data' column in written parquet files - try (CloseableIterable tasks = table.newScan().planFiles()) { - for (org.apache.iceberg.FileScanTask task : tasks) { - String path = task.file().path().toString(); - try (ParquetFileReader reader = - ParquetFileReader.open( - org.apache.parquet.hadoop.util.HadoopInputFile.fromPath( - new org.apache.hadoop.fs.Path(path), - new org.apache.hadoop.conf.Configuration()))) { - for (BlockMetaData block : reader.getFooter().getBlocks()) { - for (ColumnChunkMetaData col : block.getColumns()) { - boolean hasBloom = col.getBloomFilterOffset() > 0; - if (col.getPath().toDotString().equals("data")) { - assertTrue("Expected bloom filter on column 'data', but none was found", hasBloom); - } else { - assertFalse( - "Expected no bloom filter on column '" + col.getPath().toDotString() + "'", - hasBloom); - } - } - } - } - } - } - } - - @Test - public void testWriteWithTableProperties() throws Exception { - TableIdentifier tableId = - TableIdentifier.of( - "default", "table_props_" + Long.toString(UUID.randomUUID().hashCode(), 16)); - - Schema beamSchema = IcebergUtils.icebergSchemaToBeamSchema(TestFixtures.SCHEMA); - - Map catalogProps = - ImmutableMap.builder() - .put("type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP) - .put("warehouse", warehouse.location) - .build(); - - IcebergCatalogConfig catalog = - IcebergCatalogConfig.builder() - .setCatalogName("name") - .setCatalogProperties(catalogProps) - .build(); - - testPipeline - .apply("Records To Add", Create.of(TestFixtures.asRows(TestFixtures.FILE1SNAPSHOT1))) - .setRowSchema(beamSchema) - .apply( - "Append To Table", - writeTransform(catalog, tableId) - .withWriteProperties( - ImmutableMap.of("write.data.path", warehouse.location + "/custom_data_path"))); - - testPipeline.run().waitUntilFinish(); - - Table table = warehouse.loadTable(tableId); - - List writtenRecords = ImmutableList.copyOf(IcebergGenerics.read(table).build()); - assertThat(writtenRecords, Matchers.containsInAnyOrder(TestFixtures.FILE1SNAPSHOT1.toArray())); - - assertEquals( - warehouse.location + "/custom_data_path", table.properties().get("write.data.path")); - } } 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); + } + } + } + } + } } From de097bb632816c299e31b00fc952961d045d88ae Mon Sep 17 00:00:00 2001 From: Claire McGinty Date: Fri, 7 Aug 2026 10:57:31 -0400 Subject: [PATCH 4/4] improve documentation --- .../main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java | 3 +++ 1 file changed, 3 insertions(+) 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 9ee2e395421e..8b718fc42d92 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 @@ -493,6 +493,9 @@ public WriteRows withAutosharding() { * execution-scoped, meaning that they are applied to a preexisting table and will not mutate * any table-level properties. * + *

To set table-level properties that will be applied to dynamically created tables, use the + * managed Iceberg transform instead, setting the `table_properties` config property. + * *

See: https://iceberg.apache.org/docs/latest/configuration/#write-properties */ public WriteRows withWriteProperties(Map writeProperties) {