This is an automated email from the ASF dual-hosted git repository.
ahmedabu98 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 8d24582beab (IcebergIO) document writeProperties param more clearly
(#39645)
8d24582beab is described below
commit 8d24582beabad9717ce80ad0a1946a6e3cd57b26
Author: Claire McGinty <[email protected]>
AuthorDate: Mon Aug 10 13:21:27 2026 -0400
(IcebergIO) document writeProperties param more clearly (#39645)
* (IcebergIO) bugfix: wire writeProperties through table create request,
not DataWriteBuilder
* Test all dynamic write properties are propagated in managedio
* Revert changes and document writeProperties
* improve documentation
---
.../org/apache/beam/sdk/io/iceberg/IcebergIO.java | 10 ++++
.../IcebergWriteSchemaTransformProviderTest.java | 63 ++++++++++++++++++++++
2 files changed, 73 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 78a72ccdbb8..e92504d41e6 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
@@ -490,6 +490,16 @@ public class IcebergIO {
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.
+ *
+ * <p>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.
+ *
+ * <p>See:
https://iceberg.apache.org/docs/latest/configuration/#write-properties
+ */
public WriteRows withWriteProperties(Map<String, String> writeProperties) {
return toBuilder().setWriteProperties(writeProperties).build();
}
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 5a7aa11e10a..bfb762f8e89 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.dateFromDays;
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.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Immuta
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 class IcebergWriteSchemaTransformProviderTest {
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<String, Object> 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<Row> rows = new ArrayList<>();
+ for (int i = 0; i < 10; i++) {
+ Row row = Row.withSchema(schema).addValues("str_" + i, i).build();
+ rows.add(row);
+ }
+
+ PCollection<Row> 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<Record> writtenRecords =
ImmutableList.copyOf(IcebergGenerics.read(table).build());
+ assertEquals(10, writtenRecords.size());
+ }
}