chamikaramj commented on code in PR #40229:
URL: https://github.com/apache/beam/pull/40229#discussion_r4085185754


##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -100,14 +123,33 @@ public static Builder builder() {
         "For a streaming pipeline, sets the limit for lifting bundles into the 
direct write path.")
     public abstract @Nullable Integer getDirectWriteByteLimit();
 
+    @SchemaFieldDescription(

Review Comment:
   Looking at this map I feel like it might be better to do move these 
properties to a different SchemaTransform and to the managed API under 
Iceberg.WRITE_CDC.
   
   I think a second nested map that groups different config  properties 
together can end up being confusing the the end user. I think it's much cleaner 
to have these at the top level.



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -18,14 +18,23 @@
 package org.apache.beam.sdk.io.iceberg;
 
 import static 
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.Configuration;

Review Comment:
   Can we update Managed.java in the same PR ? If not let's keep the Github 
issue open to do it in a subsequent PR.



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -180,6 +222,62 @@ public static Builder builder() {
         "Sets the number of parallel buckets/workers used to query the Iceberg 
catalog during refreshes. Defaults to 1.")
     public abstract @Nullable Integer getPollingBuckets();
 
+    @SchemaFieldDescription(
+        "Columns defining row identity (equality-delete fields). Defaults to 
the destination table's "
+            + "identifier (primary-key) fields. Required if the table doesn't 
exist yet. Currently only supported in CDC mode.")

Review Comment:
   Will this and other properties marked as "Currently only supported in CDC 
mode" ever apply to the non-CDC mode ? If not another reason to move CDC stuff 
to a different SchemaTransform.



##########
CHANGES.md:
##########
@@ -102,6 +102,7 @@
 * ClickHouseIO: support writing `Decimal(P, S)` / `Decimal32/64/128/256` 
columns (Java) ([#39840](https://github.com/apache/beam/issues/39840)).
 * SolaceIO now supports reading and writing user properties (message metadata) 
(Java) ([#40099](https://github.com/apache/beam/issues/40099)).
 * [IcebergIO] AddFiles (`IcebergAddFiles` in YAML) can evolve the table schema 
before registering files, with `schema_evolution_options`, `required_columns`, 
`incompatible_schema_handling` and `unverifiable_file_handling` (Java/YAML, 
batch only) ([#40144](https://github.com/apache/beam/issues/40144)).
+* [IcebergIO] Added batch and streaming CDC writes that applies 
INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE changes to Iceberg V2+ tables by 
primary key. Invoke with `IcebergIO.writeCdcRows` (Java) or by setting the 
`cdc` option of Managed `ICEBERG` write (Java, Python, YAML) 
([#X](https://github.com/apache/beam/issues/X)).

Review Comment:
   `s/#X/#39979`



##########
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergSchemaTransformTranslationTest.java:
##########
@@ -98,6 +99,27 @@ public class IcebergSchemaTransformTranslationTest {
           .withFieldValue("keep", Collections.singletonList("str"))

Review Comment:
   Also add tests for `validateModeOptions` ?



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -35,28 +44,37 @@
 import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
 import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
 import org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider;
+import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling;
 import org.apache.beam.sdk.transforms.MapElements;
 import org.apache.beam.sdk.transforms.SimpleFunction;
 import org.apache.beam.sdk.values.KV;
 import org.apache.beam.sdk.values.PCollection;
 import org.apache.beam.sdk.values.PCollectionRowTuple;
 import org.apache.beam.sdk.values.Row;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
 import org.apache.iceberg.DistributionMode;
 import org.apache.iceberg.FileFormat;
 import org.checkerframework.checker.nullness.qual.Nullable;
 import org.joda.time.Duration;
 
 /**
- * SchemaTransform implementation for {@link IcebergIO#writeRows}. Writes Beam 
Rows to Iceberg and
- * outputs a {@code PCollection<Row>} representing snapshots created in the 
process.
+ * SchemaTransform implementation for {@link IcebergIO#writeRows} and, when 
the {@code cdc} block is
+ * set, {@link IcebergIO#writeCdcRows}. Outputs a {@code PCollection<Row>} 
representing the
+ * snapshots created in the process; CDC writes add a {@code dead_letter} 
output of late records.
  */
 @AutoService(SchemaTransformProvider.class)
 public class IcebergWriteSchemaTransformProvider
     extends TypedSchemaTransformProvider<Configuration> {
 
   static final String INPUT_TAG = "input";
   static final String SNAPSHOTS_TAG = "snapshots";
+  static final String DEAD_LETTER_TAG = "dead_letter";
+  static final String ERRORS_TAG = "errors";
+
+  /** The default sequence-number column: what the CDC read source emits. */
+  private static final String DEFAULT_SEQUENCE_NUMBER_COLUMN =

Review Comment:
   We should make following public and reuse instead of redefining here.
   
   
https://github.com/apache/beam/blob/c462e49b115b2c7f2f9159b193d2547c00c6f157/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfig.java#L34



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -230,6 +348,121 @@ public IcebergCatalogConfig getIcebergCatalog() {
           .setConfigProperties(getConfigProperties())
           .build();
     }
+
+    enum Mode {
+      APPEND,
+      CDC
+    }
+
+    /** The write mode this configuration selects. */
+    Mode mode() {
+      return getCdc() == null ? Mode.APPEND : Mode.CDC;
+    }
+
+    /**
+     * Options that only one mode supports today; everything else applies to 
both. Once an option
+     * becomes available to the other mode, just remove it from here.
+     */
+    private static final ImmutableMap<String, Mode> SUPPORTED_MODES =
+        ImmutableMap.<String, Mode>builder()
+            .put("equality_columns", Mode.CDC)
+            .put("num_shards", Mode.CDC)
+            .put("shards_per_partition", Mode.CDC)
+            .put("allowed_lateness_seconds", Mode.CDC)
+            .put("sink_id", Mode.CDC)
+            .put("token_heartbeat_seconds", Mode.CDC)
+            .put("snapshot_properties", Mode.CDC)
+            .put("error_handling", Mode.CDC)
+            .put("sorter_memory_mb", Mode.CDC)
+            .put("direct_write_byte_limit", Mode.APPEND)
+            .put("distribution_mode", Mode.APPEND)
+            .put("autosharding", Mode.APPEND)
+            .put("write_properties", Mode.APPEND)
+            .put("using_side_input_table_cache", Mode.APPEND)
+            .put("table_refresh_interval_seconds", Mode.APPEND)
+            .put("maximum_cache_size", Mode.APPEND)
+            .put("polling_buckets", Mode.APPEND)
+            .build();
+
+    /** Rejects every set option that the selected mode does not support. */
+    void validateModeOptions() {
+      Mode mode = mode();
+      Map<String, @Nullable Object> values = new LinkedHashMap<>();

Review Comment:
   We should derive this from the `Configuration` instead of maintaining a 
separate local map here and a `SUPPORTED_MODES` map above so that the two 
things don't drift in the future.



##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -230,6 +348,121 @@ public IcebergCatalogConfig getIcebergCatalog() {
           .setConfigProperties(getConfigProperties())
           .build();
     }
+
+    enum Mode {
+      APPEND,
+      CDC
+    }
+
+    /** The write mode this configuration selects. */
+    Mode mode() {
+      return getCdc() == null ? Mode.APPEND : Mode.CDC;
+    }
+
+    /**
+     * Options that only one mode supports today; everything else applies to 
both. Once an option
+     * becomes available to the other mode, just remove it from here.
+     */
+    private static final ImmutableMap<String, Mode> SUPPORTED_MODES =
+        ImmutableMap.<String, Mode>builder()
+            .put("equality_columns", Mode.CDC)
+            .put("num_shards", Mode.CDC)
+            .put("shards_per_partition", Mode.CDC)
+            .put("allowed_lateness_seconds", Mode.CDC)
+            .put("sink_id", Mode.CDC)
+            .put("token_heartbeat_seconds", Mode.CDC)
+            .put("snapshot_properties", Mode.CDC)
+            .put("error_handling", Mode.CDC)
+            .put("sorter_memory_mb", Mode.CDC)
+            .put("direct_write_byte_limit", Mode.APPEND)
+            .put("distribution_mode", Mode.APPEND)
+            .put("autosharding", Mode.APPEND)
+            .put("write_properties", Mode.APPEND)
+            .put("using_side_input_table_cache", Mode.APPEND)
+            .put("table_refresh_interval_seconds", Mode.APPEND)
+            .put("maximum_cache_size", Mode.APPEND)
+            .put("polling_buckets", Mode.APPEND)
+            .build();
+
+    /** Rejects every set option that the selected mode does not support. */
+    void validateModeOptions() {
+      Mode mode = mode();
+      Map<String, @Nullable Object> values = new LinkedHashMap<>();
+      values.put("equality_columns", getEqualityColumns());
+      values.put("num_shards", getNumShards());
+      values.put("shards_per_partition", getShardsPerPartition());
+      values.put("allowed_lateness_seconds", getAllowedLatenessSeconds());
+      values.put("sink_id", getSinkId());
+      values.put("token_heartbeat_seconds", getTokenHeartbeatSeconds());
+      values.put("snapshot_properties", getSnapshotProperties());
+      values.put("error_handling", getErrorHandling());
+      values.put("sorter_memory_mb", getSorterMemoryMb());
+      values.put("direct_write_byte_limit", getDirectWriteByteLimit());
+      values.put("distribution_mode", getDistributionMode());
+      values.put("autosharding", getAutosharding());
+      values.put("write_properties", getWriteProperties());
+      values.put("using_side_input_table_cache", 
getUsingSideInputTableCache());
+      values.put("table_refresh_interval_seconds", 
getTableRefreshIntervalSeconds());
+      values.put("maximum_cache_size", getMaximumCacheSize());
+      values.put("polling_buckets", getPollingBuckets());
+      List<String> invalidOptions = new ArrayList<>();
+      for (Map.Entry<String, Mode> option : SUPPORTED_MODES.entrySet()) {
+        String name = option.getKey();
+        if (values.get(name) == null || option.getValue() == mode) {
+          continue;
+        }
+        invalidOptions.add(name);
+      }
+      if (!invalidOptions.isEmpty()) {
+        throw new IllegalArgumentException(
+            String.format(
+                "The following options are not supported by %s writes yet: %s",

Review Comment:
   Probably give a better hint regarding potential action to take in this error 
message based on the property. For example,
   
   * set the cdc block
   * these aren't implemented for CDC yet



-- 
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]

Reply via email to