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 e3b3e7ca1f7 [Iceberg CDC sink] SchemaTransform (ManagedIO) layer 
(#40229)
e3b3e7ca1f7 is described below

commit e3b3e7ca1f7a372d44ad563b1d4e4f4acb30b298
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Wed Sep 23 20:41:51 2026 -0700

    [Iceberg CDC sink] SchemaTransform (ManagedIO) layer (#40229)
---
 CHANGES.md                                         |   1 +
 .../IcebergWriteSchemaTransformProvider.java       | 360 ++++++++++-
 .../sdk/io/iceberg/cdc/sink/CdcWriteConfig.java    |   7 +-
 .../beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java |   5 +-
 .../IcebergSchemaTransformTranslationTest.java     |  27 +
 ...IcebergWriteSchemaTransformProviderCdcTest.java | 696 +++++++++++++++++++++
 .../sdk/io/iceberg/cdc/sink/AssignCdcKeysTest.java |   3 +-
 .../io/iceberg/cdc/sink/CdcWriteConfigTest.java    |   4 +-
 sdks/python/apache_beam/yaml/yaml_io.py            |  75 ++-
 9 files changed, 1160 insertions(+), 18 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 87b8ad672eb..63f84e249ec 100644
--- a/CHANGES.md
+++ b/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 `mode: 
merge-on-read` on the Managed `ICEBERG` write (Java, Python, YAML) 
([#39979](https://github.com/apache/beam/issues/39979)).
 
 ## New Features / Improvements
 
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 e5a81e41a8b..4f568fb6562 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
@@ -18,14 +18,21 @@
 package org.apache.beam.sdk.io.iceberg;
 
 import static 
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.Configuration;
+import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
 import static org.apache.beam.sdk.util.construction.BeamUrns.getUrn;
 
 import com.google.auto.service.AutoService;
 import com.google.auto.value.AutoValue;
+import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
+import java.util.LinkedHashSet;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import org.apache.beam.model.pipeline.v1.ExternalTransforms;
+import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
+import org.apache.beam.sdk.io.iceberg.cdc.sink.WriteCdcRows;
 import org.apache.beam.sdk.schemas.AutoValueSchema;
 import org.apache.beam.sdk.schemas.NoSuchSchemaException;
 import org.apache.beam.sdk.schemas.Schema;
@@ -35,6 +42,7 @@ import 
org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription;
 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;
@@ -48,8 +56,10 @@ 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, in 
{@code merge-on-read}
+ * mode, {@link IcebergIO#writeCdcRows}. Outputs a {@code PCollection<Row>} 
representing the
+ * snapshots created in the process; merge-on-read writes add a {@code 
dead_letter} output of late
+ * records.
  */
 @AutoService(SchemaTransformProvider.class)
 public class IcebergWriteSchemaTransformProvider
@@ -57,6 +67,8 @@ public class IcebergWriteSchemaTransformProvider
 
   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";
 
   static final Schema OUTPUT_SCHEMA =
       Schema.builder()
@@ -66,9 +78,14 @@ public class IcebergWriteSchemaTransformProvider
 
   @Override
   public String description() {
-    return "Writes Beam Rows to Iceberg.\n"
-        + "Returns a PCollection representing the snapshots produced in the 
process, with the following schema:\n"
-        + "{\"table\" (str), \"operation\" (str), \"summary\" (map[str, str]), 
\"manifestListLocation\" (str)}";
+    return "Writes Beam Rows to Iceberg, appending them by default. Set mode 
to 'merge-on-read' to "
+        + "apply them as a stream of row-level changes 
(INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE) "
+        + "by primary key instead.\n"
+        + "Returns a 'snapshots' PCollection representing the snapshots 
produced in the process, "
+        + "with the following schema:\n"
+        + "{\"table\" (str), \"operation\" (str), \"summary\" (map[str, str]), 
\"manifestListLocation\" (str)}\n"
+        + "Merge-on-read mode also returns a 'dead_letter' PCollection 
representing late data, and an "
+        + "'errors' PCollection representing invalid records.";
   }
 
   @DefaultSchema(AutoValueSchema.class)
@@ -100,14 +117,44 @@ public class IcebergWriteSchemaTransformProvider
         "For a streaming pipeline, sets the limit for lifting bundles into the 
direct write path.")
     public abstract @Nullable Integer getDirectWriteByteLimit();
 
+    @SchemaFieldDescription(
+        "Controls how rows are written. 'append' (default) appends every row 
as new data. "
+            + "'merge-on-read' treats each row as a change (INSERT, 
UPDATE_BEFORE, UPDATE_AFTER, "
+            + "or DELETE) applied to the table by primary key.")
+    public abstract @Nullable String getMode();
+
+    @SchemaFieldDescription(
+        "Merge-on-read only. The required column name representing the 
monotonic sequence number used to "
+            + "order a single key's changes. Defaults to 
'_commit_snapshot_sequence_number'. This column will be "
+            + "stripped from the data row before writing to Iceberg.")
+    public abstract @Nullable String getSequenceNumberColumn();
+
+    @SchemaFieldDescription(
+        "Merge-on-read only. The optional column name representing the row's 
change type (INSERT, "
+            + "UPDATE_BEFORE, UPDATE_AFTER,  or DELETE). This column will be 
stripped from the data row "
+            + "before writing to Iceberg. If unset, the sink will use the 
element's native ValueKind")
+    public abstract @Nullable String getChangeTypeColumn();
+
+    @SchemaFieldDescription(
+        "Merge-on-read only. Optional map from a change_type_column value to 
the canonical change "
+            + "type name (see above).")
+    public abstract @Nullable Map<String, String> getChangeTypeMap();
+
+    @SchemaFieldDescription(
+        "Merge-on-read only. If true, only the after-image of each change 
(INSERT/UPDATE_AFTER) "
+            + "is applied, as an upsert; UPDATE_BEFORE records are dropped. 
Default: false.")
+    public abstract @Nullable Boolean getUpsert();
+
     @SchemaFieldDescription(
         "A list of field names to keep in the input record. All other fields 
are dropped before writing. "
-            + "Is mutually exclusive with 'drop' and 'only'.")
+            + "Is mutually exclusive with 'drop' and 'only'. In merge-on-read 
mode the control columns are "
+            + "dropped unless listed here.")
     public abstract @Nullable List<String> getKeep();
 
     @SchemaFieldDescription(
         "A list of field names to drop from the input record before writing. "
-            + "Is mutually exclusive with 'keep' and 'only'.")
+            + "Is mutually exclusive with 'keep' and 'only'. In merge-on-read 
mode the control columns are "
+            + "always dropped.")
     public abstract @Nullable List<String> getDrop();
 
     @SchemaFieldDescription(
@@ -180,6 +227,63 @@ public class IcebergWriteSchemaTransformProvider
         "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 'merge-on-read' mode.")
+    public abstract @Nullable List<String> getEqualityColumns();
+
+    @SchemaFieldDescription(
+        "The number of deterministic primary-key-hash shards per destination, 
i.e. the max "
+            + "write parallelism per destination. Too low may bottleneck 
writes, and too high may "
+            + "produce more files. Defaults to 16. Currently only supported in 
'merge-on-read' mode.")
+    public abstract @Nullable Integer getNumShards();
+
+    @SchemaFieldDescription(
+        "Maximum number of shards a single partition's rows may occupy. Lower 
values "
+            + "write fewer files per commit, but also reduces per-partition 
write parallelism. A value "
+            + "of 1 pins each partition to one writer. Ignored for 
unpartitioned tables. Must be between 1 "
+            + "and `num_shards`; defaults to `num_shards`. Currently only 
supported in 'merge-on-read' mode.")
+    public abstract @Nullable Integer getShardsPerPartition();
+
+    @SchemaFieldDescription(
+        "How long a late record may lag behind the watermark before it is "
+            + "dropped entirely, rather than routed to the dead_letter output. 
Defaults to 21600 "
+            + "(6 hours). Currently only supported in 'merge-on-read' mode.")
+    public abstract @Nullable Integer getAllowedLatenessSeconds();
+
+    @SchemaFieldDescription(
+        "A stable identifier for this sink, used to namespace the idempotency 
tokens "
+            + "written to each commit's Iceberg snapshot summary. Defaults to 
a unique per-write UUID. "
+            + "Set it explicitly (and keep it stable across relaunches) for 
exactly-once commits across "
+            + "relaunches of a particular streaming write. A batch load with a 
stable sink_id "
+            + "commits only once (later batch loads with the same sink_id are 
skipped). "
+            + "Currently only supported in 'merge-on-read' mode.")
+    public abstract @Nullable String getSinkId();
+
+    @SchemaFieldDescription(
+        "Streaming only. If set, the sink will emit a periodic empty 
token-refresh commit while idle, "
+            + "so its thread of `sink_id` stamped snapshot stays recent and is 
less likely to be "
+            + "lost to `expire_snapshots`. Disabled by default. Currently only 
supported in 'merge-on-read' mode.")
+    public abstract @Nullable Integer getTokenHeartbeatSeconds();
+
+    @SchemaFieldDescription(
+        "Extra key/value properties to add to every commit's Iceberg snapshot 
summary. "
+            + "Keys prefixed with 'beam.cdc.' are reserved and rejected. 
Currently only supported in 'merge-on-read' mode.")
+    public abstract @Nullable Map<String, String> getSnapshotProperties();
+
+    @SchemaFieldDescription(
+        "Whether and where to output per-record invalid rows (null or missing 
sequence "
+            + "value, unknown change type, null equality value, unresolvable 
destination). Fails the pipeline "
+            + "if unset (default). Distinct from the `dead_letter` output, 
which is for late-but-valid "
+            + "rows. Currently only supported in 'merge-on-read' mode.")
+    public abstract @Nullable ErrorHandling getErrorHandling();
+
+    @SchemaFieldDescription(
+        "The in-memory buffer size (MB) for the pre-write sort; groups larger 
than this "
+            + "spill to disk. Must be >= 1. Defaults to 100. Currently only 
supported in 'merge-on-read' mode.")
+    public abstract @Nullable Integer getSorterMemoryMb();
+
     @AutoValue.Builder
     public abstract static class Builder {
       public abstract Builder setTable(String table);
@@ -220,6 +324,34 @@ public class IcebergWriteSchemaTransformProvider
 
       public abstract Builder setPollingBuckets(Integer pollingBuckets);
 
+      public abstract Builder setMode(String mode);
+
+      public abstract Builder setSequenceNumberColumn(String 
sequenceNumberColumn);
+
+      public abstract Builder setChangeTypeColumn(String changeTypeColumn);
+
+      public abstract Builder setChangeTypeMap(Map<String, String> 
changeTypeMap);
+
+      public abstract Builder setUpsert(Boolean upsert);
+
+      public abstract Builder setEqualityColumns(List<String> equalityColumns);
+
+      public abstract Builder setNumShards(Integer numShards);
+
+      public abstract Builder setShardsPerPartition(Integer 
shardsPerPartition);
+
+      public abstract Builder setAllowedLatenessSeconds(Integer 
allowedLatenessSeconds);
+
+      public abstract Builder setSinkId(String sinkId);
+
+      public abstract Builder setTokenHeartbeatSeconds(Integer 
tokenHeartbeatSeconds);
+
+      public abstract Builder setSnapshotProperties(Map<String, String> 
snapshotProperties);
+
+      public abstract Builder setErrorHandling(ErrorHandling errorHandling);
+
+      public abstract Builder setSorterMemoryMb(Integer sorterMemoryMb);
+
       public abstract Configuration build();
     }
 
@@ -230,6 +362,86 @@ public class IcebergWriteSchemaTransformProvider
           .setConfigProperties(getConfigProperties())
           .build();
     }
+
+    enum Mode {
+      APPEND("append"),
+      MERGE_ON_READ("merge-on-read");
+
+      /** The value users set {@code mode} to. */
+      final String optionValue;
+
+      Mode(String optionValue) {
+        this.optionValue = optionValue;
+      }
+    }
+
+    /** The write mode this configuration selects; unset means append. */
+    Mode mode() {
+      @Nullable String mode = getMode();
+      if (mode == null || mode.equalsIgnoreCase(Mode.APPEND.optionValue)) {
+        return Mode.APPEND;
+      }
+      if (mode.equalsIgnoreCase(Mode.MERGE_ON_READ.optionValue)) {
+        return Mode.MERGE_ON_READ;
+      }
+      throw new IllegalArgumentException(
+          String.format(
+              "Unknown mode '%s'; expected '%s' or '%s'.",
+              mode, Mode.APPEND.optionValue, Mode.MERGE_ON_READ.optionValue));
+    }
+
+    /** Rejects every set option that the selected mode does not support. */
+    void validateModeOptions() {
+      // Resolve the mode first so an unknown value fails here, not only once 
an option trips it.
+      Mode mode = mode();
+      List<String> unsupported = new ArrayList<>();
+      // Merge-on-read only: the change-stream contract.
+      requireMode(
+          Mode.MERGE_ON_READ, "sequence_number_column", 
getSequenceNumberColumn(), unsupported);
+      requireMode(Mode.MERGE_ON_READ, "change_type_column", 
getChangeTypeColumn(), unsupported);
+      requireMode(Mode.MERGE_ON_READ, "change_type_map", getChangeTypeMap(), 
unsupported);
+      requireMode(Mode.MERGE_ON_READ, "upsert", getUpsert(), unsupported);
+      // Merge-on-read only, until the append write grows these features.
+      requireMode(Mode.MERGE_ON_READ, "equality_columns", 
getEqualityColumns(), unsupported);
+      requireMode(Mode.MERGE_ON_READ, "num_shards", getNumShards(), 
unsupported);
+      requireMode(Mode.MERGE_ON_READ, "shards_per_partition", 
getShardsPerPartition(), unsupported);
+      requireMode(
+          Mode.MERGE_ON_READ, "allowed_lateness_seconds", 
getAllowedLatenessSeconds(), unsupported);
+      requireMode(Mode.MERGE_ON_READ, "sink_id", getSinkId(), unsupported);
+      requireMode(
+          Mode.MERGE_ON_READ, "token_heartbeat_seconds", 
getTokenHeartbeatSeconds(), unsupported);
+      requireMode(Mode.MERGE_ON_READ, "snapshot_properties", 
getSnapshotProperties(), unsupported);
+      requireMode(Mode.MERGE_ON_READ, "error_handling", getErrorHandling(), 
unsupported);
+      requireMode(Mode.MERGE_ON_READ, "sorter_memory_mb", getSorterMemoryMb(), 
unsupported);
+      // Append only: merge-on-read has neither a direct-write path nor a 
side-input table cache.
+      requireMode(Mode.APPEND, "direct_write_byte_limit", 
getDirectWriteByteLimit(), unsupported);
+      requireMode(Mode.APPEND, "distribution_mode", getDistributionMode(), 
unsupported);
+      requireMode(Mode.APPEND, "autosharding", getAutosharding(), unsupported);
+      requireMode(Mode.APPEND, "write_properties", getWriteProperties(), 
unsupported);
+      requireMode(
+          Mode.APPEND, "using_side_input_table_cache", 
getUsingSideInputTableCache(), unsupported);
+      requireMode(
+          Mode.APPEND,
+          "table_refresh_interval_seconds",
+          getTableRefreshIntervalSeconds(),
+          unsupported);
+      requireMode(Mode.APPEND, "maximum_cache_size", getMaximumCacheSize(), 
unsupported);
+      requireMode(Mode.APPEND, "polling_buckets", getPollingBuckets(), 
unsupported);
+      if (!unsupported.isEmpty()) {
+        throw new IllegalArgumentException(
+            String.format(
+                "The following options are not supported in '%s' mode yet: %s",
+                mode.optionValue, unsupported));
+      }
+    }
+
+    /** Records {@code option} as unsupported when it is set under a mode 
other than its own. */
+    private void requireMode(
+        Mode supported, String option, @Nullable Object value, List<String> 
unsupported) {
+      if (value != null && mode() != supported) {
+        unsupported.add(option);
+      }
+    }
   }
 
   @Override
@@ -244,7 +456,7 @@ public class IcebergWriteSchemaTransformProvider
 
   @Override
   public List<String> outputCollectionNames() {
-    return Collections.singletonList(SNAPSHOTS_TAG);
+    return Arrays.asList(SNAPSHOTS_TAG, DEAD_LETTER_TAG, ERRORS_TAG);
   }
 
   @Override
@@ -276,7 +488,13 @@ public class IcebergWriteSchemaTransformProvider
     @Override
     public PCollectionRowTuple expand(PCollectionRowTuple input) {
       PCollection<Row> rows = input.get(INPUT_TAG);
+      configuration.validateModeOptions();
+      return configuration.mode() == Configuration.Mode.MERGE_ON_READ
+          ? expandCdc(rows)
+          : expandAppend(rows);
+    }
 
+    private PCollectionRowTuple expandAppend(PCollection<Row> rows) {
       IcebergIO.WriteRows writeTransform =
           IcebergIO.writeRows(configuration.getIcebergCatalog())
               .to(
@@ -360,6 +578,132 @@ public class IcebergWriteSchemaTransformProvider
       return PCollectionRowTuple.of(SNAPSHOTS_TAG, snapshots);
     }
 
+    private PCollectionRowTuple expandCdc(PCollection<Row> rows) {
+      Schema inputSchema = rows.getSchema();
+      @Nullable List<String> drop = configuration.getDrop();
+      @Nullable List<String> keep = configuration.getKeep();
+      @Nullable String only = configuration.getOnly();
+      @Nullable String changeTypeColumn = configuration.getChangeTypeColumn();
+      @Nullable String configuredSeq = configuration.getSequenceNumberColumn();
+      String seqColumn =
+          configuredSeq != null
+              ? configuredSeq
+              : IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER;
+
+      // The sink reads the control columns from the raw element, so by 
default they are dropped
+      // from the written row. Listing them in keep writes them too.
+      if (keep == null && only == null) {
+        Set<String> effectiveDrop =
+            new LinkedHashSet<>(controlColumnsPresent(inputSchema, 
changeTypeColumn, seqColumn));
+        if (drop != null) {
+          effectiveDrop.addAll(drop);
+        }
+        drop = effectiveDrop.isEmpty() ? null : new ArrayList<>(effectiveDrop);
+      }
+
+      WriteCdcRows write =
+          IcebergIO.writeCdcRows(configuration.getIcebergCatalog())
+              .to(
+                  new PortableIcebergDestinations(
+                      configuration.getTable(),
+                      FileFormat.PARQUET.toString(),
+                      inputSchema,
+                      configuration.getPartitionFields(),
+                      configuration.getSortFields(),
+                      configuration.getTableProperties(),
+                      drop,
+                      keep,
+                      only));
+      IcebergWriteResult result = rows.apply(applyOptions(write));
+
+      PCollection<Row> snapshots =
+          result
+              .getSnapshots()
+              .apply(MapElements.via(new SnapshotToRow()))
+              .setRowSchema(OUTPUT_SCHEMA);
+      PCollectionRowTuple output =
+          PCollectionRowTuple.of(SNAPSHOTS_TAG, snapshots)
+              .and(DEAD_LETTER_TAG, result.getDeadLetterRows());
+      @Nullable ErrorHandling errorHandling = configuration.getErrorHandling();
+      if (ErrorHandling.hasOutput(errorHandling)) {
+        output = output.and(checkStateNotNull(errorHandling).getOutput(), 
result.getFailedRows());
+      }
+      return output;
+    }
+
+    /** Threads every set option onto {@code write}. */
+    private WriteCdcRows applyOptions(WriteCdcRows write) {
+      @Nullable List<String> equalityColumns = 
configuration.getEqualityColumns();
+      if (equalityColumns != null) {
+        write = write.withEqualityColumns(equalityColumns);
+      }
+      @Nullable String sequenceNumberColumn = 
configuration.getSequenceNumberColumn();
+      if (sequenceNumberColumn != null) {
+        write = write.withSequenceNumberColumn(sequenceNumberColumn);
+      }
+      @Nullable String changeTypeColumn = configuration.getChangeTypeColumn();
+      if (changeTypeColumn != null) {
+        write = write.withChangeTypeColumn(changeTypeColumn);
+      }
+      @Nullable Map<String, String> changeTypeMap = 
configuration.getChangeTypeMap();
+      if (changeTypeMap != null) {
+        write = write.withChangeTypeMap(changeTypeMap);
+      }
+      @Nullable Boolean upsert = configuration.getUpsert();
+      if (upsert != null) {
+        write = write.withUpsert(upsert);
+      }
+      @Nullable Integer sorterMemoryMb = configuration.getSorterMemoryMb();
+      if (sorterMemoryMb != null) {
+        write = write.withSorterMemoryMB(sorterMemoryMb);
+      }
+      @Nullable Integer numShards = configuration.getNumShards();
+      if (numShards != null) {
+        write = write.withNumShards(numShards);
+      }
+      @Nullable Integer shardsPerPartition = 
configuration.getShardsPerPartition();
+      if (shardsPerPartition != null) {
+        write = write.withShardsPerPartition(shardsPerPartition);
+      }
+      @Nullable String sinkId = configuration.getSinkId();
+      if (sinkId != null) {
+        write = write.withSinkId(sinkId);
+      }
+      @Nullable Integer triggeringFrequencySeconds = 
configuration.getTriggeringFrequencySeconds();
+      if (triggeringFrequencySeconds != null) {
+        write = 
write.withTriggeringFrequency(Duration.standardSeconds(triggeringFrequencySeconds));
+      }
+      @Nullable Integer allowedLatenessSeconds = 
configuration.getAllowedLatenessSeconds();
+      if (allowedLatenessSeconds != null) {
+        write = 
write.withAllowedLateness(Duration.standardSeconds(allowedLatenessSeconds));
+      }
+      if (ErrorHandling.hasOutput(configuration.getErrorHandling())) {
+        write = write.withErrorHandling();
+      }
+      @Nullable Map<String, String> snapshotProperties = 
configuration.getSnapshotProperties();
+      if (snapshotProperties != null) {
+        write = write.withSnapshotProperties(snapshotProperties);
+      }
+      @Nullable Integer tokenHeartbeatSeconds = 
configuration.getTokenHeartbeatSeconds();
+      if (tokenHeartbeatSeconds != null) {
+        write = 
write.withTokenHeartbeat(Duration.standardSeconds(tokenHeartbeatSeconds));
+      }
+      return write;
+    }
+
+    /** The control columns present in the input; a missing one is left to the 
sink to report. */
+    private static List<String> controlColumnsPresent(
+        Schema inputSchema, @Nullable String changeTypeColumn, String 
seqColumn) {
+      List<String> controls = new ArrayList<>();
+      if (changeTypeColumn != null && inputSchema.hasField(changeTypeColumn)) {
+        controls.add(changeTypeColumn);
+      }
+      if (inputSchema.hasField(seqColumn)) {
+        controls.add(seqColumn);
+      }
+      return controls;
+    }
+
     @VisibleForTesting
     static class SnapshotToRow extends SimpleFunction<KV<String, 
SnapshotInfo>, Row> {
       @Override
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfig.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfig.java
index 537ff6aaebb..32335dd14fb 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfig.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfig.java
@@ -31,9 +31,6 @@ import org.checkerframework.checker.nullness.qual.Nullable;
 /** Configuration for the CDC sink. */
 @AutoValue
 abstract class CdcWriteConfig implements Serializable {
-  static final String DEFAULT_SEQUENCE_NUMBER_COLUMN =
-      IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER;
-
   /**
    * Every shard that touches a partition writes a file per commit window, so 
{@code num_shards x
    * touched partitions x windows per day} files. On a <b>partitioned</b> 
table {@link
@@ -52,7 +49,7 @@ abstract class CdcWriteConfig implements Serializable {
 
   /**
    * The column holding the per-primary-key monotonic sequence number used to 
order a single key's
-   * changes. Defaults to {@value #DEFAULT_SEQUENCE_NUMBER_COLUMN}.
+   * changes. Defaults to {@value 
IcebergCdcMetadataColumns#COMMIT_SNAPSHOT_SEQUENCE_NUMBER}.
    */
   abstract String getSequenceNumberColumn();
 
@@ -131,7 +128,7 @@ abstract class CdcWriteConfig implements Serializable {
 
   static Builder builder() {
     return new AutoValue_CdcWriteConfig.Builder()
-        .setSequenceNumberColumn(DEFAULT_SEQUENCE_NUMBER_COLUMN)
+        
.setSequenceNumberColumn(IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER)
         .setNumShards(DEFAULT_NUM_SHARDS)
         .setShardsPerPartition(DEFAULT_NUM_SHARDS)
         .setSorterMemoryMB(DEFAULT_SORTER_MEMORY_MB)
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java
index 9d240c0e7d7..e32ad2f8d5a 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java
@@ -29,6 +29,7 @@ import org.apache.beam.sdk.io.iceberg.DynamicDestinations;
 import org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig;
 import org.apache.beam.sdk.io.iceberg.IcebergWriteResult;
 import org.apache.beam.sdk.io.iceberg.SnapshotInfo;
+import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.transforms.PTransform;
 import org.apache.beam.sdk.transforms.display.DisplayData;
@@ -184,7 +185,7 @@ public abstract class WriteCdcRows extends 
PTransform<PCollection<Row>, IcebergW
   public static WriteCdcRows of(IcebergCatalogConfig catalogConfig) {
     return new AutoValue_WriteCdcRows.Builder()
         .setCatalogConfig(catalogConfig)
-        .setSequenceNumberColumn(CdcWriteConfig.DEFAULT_SEQUENCE_NUMBER_COLUMN)
+        
.setSequenceNumberColumn(IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER)
         .setNumShards(CdcWriteConfig.DEFAULT_NUM_SHARDS)
         .setSorterMemoryMB(CdcWriteConfig.DEFAULT_SORTER_MEMORY_MB)
         .setUpsert(false)
@@ -223,7 +224,7 @@ public abstract class WriteCdcRows extends 
PTransform<PCollection<Row>, IcebergW
    * The column holding the per-primary-key monotonic sequence number used to 
order a single key's
    * changes. Must be declared as a non-nullable {@code INT64} in the input 
schema. The column is
    * stripped from the written rows. Defaults to {@value
-   * CdcWriteConfig#DEFAULT_SEQUENCE_NUMBER_COLUMN}.
+   * IcebergCdcMetadataColumns#COMMIT_SNAPSHOT_SEQUENCE_NUMBER}.
    */
   public WriteCdcRows withSequenceNumberColumn(String column) {
     return toBuilder().setSequenceNumberColumn(column).build();
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergSchemaTransformTranslationTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergSchemaTransformTranslationTest.java
index 675b4aafe76..3c07c97b4a5 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergSchemaTransformTranslationTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergSchemaTransformTranslationTest.java
@@ -98,6 +98,18 @@ public class IcebergSchemaTransformTranslationTest {
           .withFieldValue("keep", Collections.singletonList("str"))
           .build();
 
+  /** A merge-on-read write config: the mode plus the options gated to it. */
+  private static final Row WRITE_CDC_CONFIG_ROW =
+      Row.withSchema(WRITE_PROVIDER.configurationSchema())
+          .withFieldValue("table", "test_table_identifier")
+          .withFieldValue("catalog_properties", CATALOG_PROPERTIES)
+          .withFieldValue("mode", "merge-on-read")
+          .withFieldValue("change_type_column", "op")
+          .withFieldValue("upsert", true)
+          .withFieldValue("equality_columns", Collections.singletonList("id"))
+          .withFieldValue("sink_id", "stable-sink")
+          .build();
+
   private static final Row READ_CONFIG_ROW =
       Row.withSchema(READ_PROVIDER.configurationSchema())
           .withFieldValue("table", "test_table_identifier")
@@ -136,6 +148,21 @@ public class IcebergSchemaTransformTranslationTest {
     assertEquals(WRITE_CONFIG_ROW, 
writeTransformFromRow.getConfigurationRow());
   }
 
+  @Test
+  public void testReCreateCdcWriteTransformFromRow() {
+    IcebergWriteSchemaTransform writeTransform =
+        (IcebergWriteSchemaTransform) 
WRITE_PROVIDER.from(WRITE_CDC_CONFIG_ROW);
+
+    IcebergSchemaTransformTranslation.IcebergWriteSchemaTransformTranslator 
translator =
+        new 
IcebergSchemaTransformTranslation.IcebergWriteSchemaTransformTranslator();
+    Row row = translator.toConfigRow(writeTransform);
+
+    IcebergWriteSchemaTransform writeTransformFromRow =
+        translator.fromConfigRow(row, PipelineOptionsFactory.create());
+
+    assertEquals(WRITE_CDC_CONFIG_ROW, 
writeTransformFromRow.getConfigurationRow());
+  }
+
   @Test
   public void testWriteTransformProtoTranslation()
       throws InvalidProtocolBufferException, IOException {
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderCdcTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderCdcTest.java
new file mode 100644
index 00000000000..8f2199e6e52
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProviderCdcTest.java
@@ -0,0 +1,696 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+import static 
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.Configuration;
+import static 
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.DEAD_LETTER_TAG;
+import static 
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.INPUT_TAG;
+import static 
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.SNAPSHOTS_TAG;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsInAnyOrder;
+import static org.hamcrest.Matchers.containsString;
+import static org.hamcrest.Matchers.equalTo;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+import java.util.List;
+import java.util.Map;
+import org.apache.beam.sdk.managed.Managed;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
+import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling;
+import org.apache.beam.sdk.testing.PAssert;
+import org.apache.beam.sdk.testing.TestPipeline;
+import org.apache.beam.sdk.transforms.Create;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.ParDo;
+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.sdk.values.ValueKind;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet;
+import org.apache.iceberg.CatalogUtil;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.catalog.Catalog;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.IcebergGenerics;
+import org.apache.iceberg.hadoop.HadoopCatalog;
+import org.apache.iceberg.types.Types;
+import org.junit.Before;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * Tests for the {@code cdc} mode of {@link 
IcebergWriteSchemaTransformProvider}. End-to-end tests
+ * drive a real {@link HadoopCatalog} V2 table; each test uses a unique table 
identifier so the
+ * process-wide TableCache never sees a repeat.
+ */
+@RunWith(JUnit4.class)
+public class IcebergWriteSchemaTransformProviderCdcTest {
+
+  @Rule public transient TestPipeline p = TestPipeline.create();
+  @Rule public transient TemporaryFolder tmp = new TemporaryFolder();
+
+  /** Canonical test table schema (id INT primary key, name/data STRING). */
+  private static final org.apache.iceberg.Schema ICEBERG_SCHEMA =
+      new org.apache.iceberg.Schema(
+          Types.NestedField.required(1, "id", Types.IntegerType.get()),
+          Types.NestedField.optional(2, "name", Types.StringType.get()),
+          Types.NestedField.optional(3, "data", Types.StringType.get()));
+
+  private static final Schema DATA_SCHEMA =
+      Schema.builder()
+          .addInt32Field("id")
+          .addNullableField("name", Schema.FieldType.STRING)
+          .addNullableField("data", Schema.FieldType.STRING)
+          .build();
+
+  /** The default sequence-number column name (see {@code 
cdc.sink.CdcWriteConfig}). */
+  private static final String DEFAULT_SEQ_COL = 
"_commit_snapshot_sequence_number";
+
+  /** Input schema = data schema + a {@code change_type} column (string op 
code). */
+  private static final Schema INPUT_SCHEMA_WITH_CHANGE_TYPE =
+      
Schema.builder().addFields(DATA_SCHEMA.getFields()).addStringField("change_type").build();
+
+  /** {@link #INPUT_SCHEMA_WITH_CHANGE_TYPE} + the default sequence-number 
column. */
+  private static final Schema INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ =
+      Schema.builder()
+          .addFields(INPUT_SCHEMA_WITH_CHANGE_TYPE.getFields())
+          .addInt64Field(DEFAULT_SEQ_COL)
+          .build();
+
+  private Catalog catalog;
+  private String warehousePath;
+
+  @Before
+  public void setUp() {
+    warehousePath = tmp.getRoot().getAbsolutePath();
+    catalog = new HadoopCatalog(new org.apache.hadoop.conf.Configuration(), 
warehousePath);
+  }
+
+  private Map<String, String> catalogProperties() {
+    return ImmutableMap.of(
+        "type", CatalogUtil.ICEBERG_CATALOG_TYPE_HADOOP, "warehouse", "file:" 
+ warehousePath);
+  }
+
+  /** Creates a fresh unpartitioned V2 table (PK = {@code id}) with a unique 
identifier. */
+  private TableIdentifier v2Table() {
+    TableIdentifier id = TableIdentifier.of("db", "t" + System.nanoTime());
+    createV2Table(id);
+    return id;
+  }
+
+  private void createV2Table(TableIdentifier id) {
+    org.apache.iceberg.Schema schemaWithIds =
+        new org.apache.iceberg.Schema(ICEBERG_SCHEMA.columns(), 
ImmutableSet.of(1));
+    catalog.createTable(
+        id, schemaWithIds, PartitionSpec.unpartitioned(), 
ImmutableMap.of("format-version", "2"));
+  }
+
+  /** A config builder pre-wired to this test's catalog and a fresh single 
table. */
+  private Configuration.Builder configFor(TableIdentifier id) {
+    return Configuration.builder()
+        .setTable(id.toString())
+        .setCatalogProperties(catalogProperties());
+  }
+
+  /** {@link #configFor} in merge-on-read mode, reading kinds from a {@code 
change_type} column. */
+  private Configuration.Builder cdcConfigFor(TableIdentifier id) {
+    return 
configFor(id).setMode("merge-on-read").setChangeTypeColumn("change_type");
+  }
+
+  /** A merge-on-read config builder for the {@code db.{name}} destination 
template. */
+  private Configuration.Builder templateConfig() {
+    return Configuration.builder()
+        .setTable("db.{name}")
+        .setCatalogProperties(catalogProperties())
+        .setMode("merge-on-read")
+        .setChangeTypeColumn("change_type");
+  }
+
+  /**
+   * Two fresh V2 tables, {@code db.a_<suffix>} and {@code db.b_<suffix>}, for 
the template path.
+   */
+  private List<Table> templateTables(String suffix) {
+    TableIdentifier idA = TableIdentifier.of("db", "a_" + suffix);
+    TableIdentifier idB = TableIdentifier.of("db", "b_" + suffix);
+    createV2Table(idA);
+    createV2Table(idB);
+    return ImmutableList.of(catalog.loadTable(idA), catalog.loadTable(idB));
+  }
+
+  /** The provider's {@link SchemaTransform} for {@code config}. */
+  private static SchemaTransform transformFor(Configuration config) {
+    return new IcebergWriteSchemaTransformProvider().from(config);
+  }
+
+  /** Applies {@code config}'s CDC write to {@code rows}, returning the output 
tuple. */
+  private PCollectionRowTuple applyCdcWrite(Configuration config, Schema 
schema, Row... rows) {
+    PCollection<Row> input = 
p.apply(Create.of(ImmutableList.copyOf(rows)).withRowSchema(schema));
+    return PCollectionRowTuple.of(INPUT_TAG, input).apply("CdcWrite", 
transformFor(config));
+  }
+
+  /** {@link #applyCdcWrite} followed by a pipeline run to completion. */
+  private void runCdcWrite(Configuration config, Schema schema, Row... rows) {
+    applyCdcWrite(config, schema, rows);
+    p.run().waitUntilFinish();
+  }
+
+  /** Reads all live rows of {@code table} as sorted {@code id:name:data} 
strings. */
+  private static List<String> readRows(Table table) {
+    table.refresh();
+    return ImmutableList.copyOf(IcebergGenerics.read(table).build()).stream()
+        .map(r -> r.getField("id") + ":" + r.getField("name") + ":" + 
r.getField("data"))
+        .sorted()
+        .collect(ImmutableList.toImmutableList());
+  }
+
+  /** An input row with a {@code change_type} op column and an explicit 
sequence number. */
+  private static Row rowWithSeq(int id, String name, String data, String 
changeType, long seq) {
+    return Row.withSchema(INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ)
+        .addValues(id, name, data, changeType, seq)
+        .build();
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Mode gating
+  // 
---------------------------------------------------------------------------------------------
+
+  /** Options of the other mode are rejected together; shared options pass in 
both modes. */
+  @Test
+  public void validateModeOptionsRejectsOptionsOfTheOtherMode() {
+    TableIdentifier id = TableIdentifier.of("db", "t");
+
+    Configuration append = 
configFor(id).setUpsert(true).setNumShards(4).setSinkId("s").build();
+    IllegalArgumentException appendError =
+        assertThrows(IllegalArgumentException.class, 
append::validateModeOptions);
+    assertThat(
+        appendError.getMessage(),
+        equalTo(
+            "The following options are not supported in 'append' mode yet: "
+                + "[upsert, num_shards, sink_id]"));
+
+    Configuration mergeOnRead =
+        
cdcConfigFor(id).setDistributionMode("hash").setAutosharding(true).build();
+    IllegalArgumentException mergeOnReadError =
+        assertThrows(IllegalArgumentException.class, 
mergeOnRead::validateModeOptions);
+    assertThat(
+        mergeOnReadError.getMessage(),
+        equalTo(
+            "The following options are not supported in 'merge-on-read' mode 
yet: "
+                + "[distribution_mode, autosharding]"));
+
+    
configFor(id).setTriggeringFrequencySeconds(30).build().validateModeOptions();
+    
cdcConfigFor(id).setTriggeringFrequencySeconds(30).build().validateModeOptions();
+
+    IllegalArgumentException unknown =
+        assertThrows(
+            IllegalArgumentException.class,
+            () -> 
configFor(id).setMode("copy-on-write").build().validateModeOptions());
+    assertThat(unknown.getMessage(), containsString("Unknown mode 
'copy-on-write'"));
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Config values with an observable effect
+  // 
---------------------------------------------------------------------------------------------
+
+  /**
+   * {@code equality_columns} is the most dangerous field to lose (a drop 
falls back to the table's
+   * identifier fields and deletes key on the WRONG column): the identifier 
field is {@code id},
+   * {@code equality_columns} names {@code code}, and the DELETE carries a 
non-matching {@code id} ,
+   * only the configured key can apply it.
+   */
+  @Test
+  public void equalityColumnsConfigDecidesTheDeleteKey() {
+    TableIdentifier id = TableIdentifier.of("db", "eqcols" + 
System.nanoTime());
+    org.apache.iceberg.Schema schema =
+        new org.apache.iceberg.Schema(
+            ImmutableList.of(
+                Types.NestedField.required(1, "id", Types.IntegerType.get()),
+                Types.NestedField.required(2, "code", Types.StringType.get()),
+                Types.NestedField.optional(3, "val", Types.StringType.get())),
+            ImmutableSet.of(1)); // identifier field = id, deliberately NOT 
the equality column
+    catalog.createTable(
+        id, schema, PartitionSpec.unpartitioned(), 
ImmutableMap.of("format-version", "2"));
+    Table table = catalog.loadTable(id);
+
+    Schema inputSchema =
+        Schema.builder()
+            .addInt32Field("id")
+            .addStringField("code")
+            .addNullableField("val", Schema.FieldType.STRING)
+            .addStringField("change_type")
+            .addInt64Field(DEFAULT_SEQ_COL)
+            .build();
+
+    runCdcWrite(
+        cdcConfigFor(id).setEqualityColumns(ImmutableList.of("code")).build(),
+        inputSchema,
+        Row.withSchema(inputSchema).addValues(1, "k1", "x", "INSERT", 
1L).build(),
+        Row.withSchema(inputSchema).addValues(2, "k2", "y", "INSERT", 
1L).build(),
+        // A DELETE whose id (99) matches nothing: only the `code` key can 
apply it.
+        Row.withSchema(inputSchema).addValues(99, "k1", "x", "DELETE", 
2L).build());
+
+    table.refresh();
+    assertThat(
+        ImmutableList.copyOf(IcebergGenerics.read(table).build()).stream()
+            .map(r -> r.getField("id") + ":" + r.getField("code"))
+            .collect(ImmutableList.toImmutableList()),
+        containsInAnyOrder("2:k2"));
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Projection
+  // 
---------------------------------------------------------------------------------------------
+
+  /**
+   * The control columns are dropped from the written row by default; listing 
them in {@code keep}
+   * writes their raw values too, for tables that carry a last-change or 
last-sequence column.
+   */
+  @Test
+  public void keepCanWriteTheControlColumns() {
+    TableIdentifier id = TableIdentifier.of("db", "ctl" + System.nanoTime());
+    org.apache.iceberg.Schema schema =
+        new org.apache.iceberg.Schema(
+            ImmutableList.of(
+                Types.NestedField.required(1, "id", Types.IntegerType.get()),
+                Types.NestedField.optional(2, "name", Types.StringType.get()),
+                Types.NestedField.optional(3, "data", Types.StringType.get()),
+                Types.NestedField.required(4, "change_type", 
Types.StringType.get()),
+                Types.NestedField.required(5, DEFAULT_SEQ_COL, 
Types.LongType.get())),
+            ImmutableSet.of(1));
+    catalog.createTable(
+        id, schema, PartitionSpec.unpartitioned(), 
ImmutableMap.of("format-version", "2"));
+    Table table = catalog.loadTable(id);
+
+    runCdcWrite(
+        cdcConfigFor(id)
+            .setKeep(ImmutableList.of("id", "name", "data", "change_type", 
DEFAULT_SEQ_COL))
+            .build(),
+        INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ,
+        rowWithSeq(1, "a", "x", "INSERT", 1L),
+        rowWithSeq(1, "a", "z", "UPDATE_AFTER", 2L),
+        rowWithSeq(2, "b", "y", "INSERT", 1L),
+        rowWithSeq(2, "b", "y", "DELETE", 2L));
+
+    // The sink still read both columns: id=1 shows its seq=2 image and id=2 
was deleted.
+    table.refresh();
+    assertThat(
+        ImmutableList.copyOf(IcebergGenerics.read(table).build()).stream()
+            .map(
+                r ->
+                    r.getField("id")
+                        + ":"
+                        + r.getField("data")
+                        + ":"
+                        + r.getField("change_type")
+                        + ":"
+                        + r.getField(DEFAULT_SEQ_COL))
+            .collect(ImmutableList.toImmutableList()),
+        equalTo(ImmutableList.of("1:z:UPDATE_AFTER:2")));
+  }
+
+  /**
+   * The unnested dead-letter shape re-feeds into a single-table sink: {@code 
destination} dropped
+   * by the projection, the control columns consumed (then stripped) by the 
sink.
+   */
+  @Test
+  public void singleTableReplaysDeadLetterShape() {
+    TableIdentifier id = v2Table();
+    Table table = catalog.loadTable(id);
+
+    Schema replaySchema =
+        Schema.builder()
+            .addFields(DATA_SCHEMA.getFields())
+            .addStringField("change_type")
+            .addInt64Field("sequence_number")
+            .addStringField("destination")
+            .build();
+    runCdcWrite(
+        cdcConfigFor(id)
+            .setSequenceNumberColumn("sequence_number")
+            .setDrop(ImmutableList.of("destination"))
+            .build(),
+        replaySchema,
+        Row.withSchema(replaySchema).addValues(1, "a", "x", "INSERT", 1L, 
id.toString()).build());
+
+    assertThat(readRows(table), equalTo(ImmutableList.of("1:a:x")));
+  }
+
+  /**
+   * A {@code keep} whitelist that omits the control columns must still let 
both through to the
+   * sink: the UPDATE lands and the DELETE deletes only if kinds and order 
arrived.
+   */
+  @Test
+  public void controlColumnsSurviveKeepProjection() {
+    TableIdentifier id = v2Table();
+    Table table = catalog.loadTable(id);
+
+    Schema inputSchema =
+        Schema.builder()
+            .addFields(INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ.getFields())
+            .addStringField("source")
+            .build();
+
+    // keep lists only the data columns: no control columns, no "source" 
envelope column.
+    runCdcWrite(
+        cdcConfigFor(id).setKeep(ImmutableList.of("id", "name", 
"data")).build(),
+        inputSchema,
+        Row.withSchema(inputSchema).addValues(1, "a", "x", "INSERT", 1L, 
"m").build(),
+        Row.withSchema(inputSchema).addValues(2, "b", "y", "INSERT", 1L, 
"m").build(),
+        Row.withSchema(inputSchema).addValues(1, "a", "x", "UPDATE_BEFORE", 
2L, "m").build(),
+        Row.withSchema(inputSchema).addValues(1, "a", "z", "UPDATE_AFTER", 2L, 
"m").build(),
+        Row.withSchema(inputSchema).addValues(2, "b", "y", "DELETE", 2L, 
"m").build());
+
+    assertThat(readRows(table), equalTo(ImmutableList.of("1:a:z")));
+  }
+
+  /**
+   * Merge-on-read mode with nothing else set uses every default: kinds come 
from each element's
+   * native {@link ValueKind}, which the projection must preserve (a plain 
{@code output()} re-emit
+   * would stamp everything INSERT).
+   */
+  @Test
+  public void mergeOnReadDefaultsUseNativeValueKinds() {
+    TableIdentifier id = v2Table();
+    Table table = catalog.loadTable(id);
+
+    Schema inputSchema =
+        Schema.builder()
+            .addFields(DATA_SCHEMA.getFields())
+            .addInt64Field(DEFAULT_SEQ_COL)
+            .addStringField("source")
+            .build();
+
+    // No change_type_column: change kinds come only from each element's 
native ValueKind.
+    SchemaTransform transform =
+        transformFor(
+            
configFor(id).setMode("merge-on-read").setDrop(ImmutableList.of("source")).build());
+
+    PCollection<Row> input =
+        withKinds(
+                p.apply(
+                    Create.of(
+                        ImmutableList.of(
+                            KV.of(ValueKind.INSERT, kindRow(inputSchema, 1, 
"a", "x", 1L)),
+                            KV.of(ValueKind.INSERT, kindRow(inputSchema, 2, 
"b", "y", 1L)),
+                            KV.of(ValueKind.UPDATE_BEFORE, 
kindRow(inputSchema, 1, "a", "x", 2L)),
+                            KV.of(ValueKind.UPDATE_AFTER, kindRow(inputSchema, 
1, "a", "z", 2L)),
+                            KV.of(ValueKind.DELETE, kindRow(inputSchema, 2, 
"b", "y", 2L))))))
+            .setRowSchema(inputSchema);
+
+    PCollectionRowTuple.of(INPUT_TAG, input).apply("CdcWrite", transform);
+    p.run().waitUntilFinish();
+
+    // id=1 updated, id=2 deleted , possible only if the native kinds survived 
the ParDo.
+    assertThat(readRows(table), equalTo(ImmutableList.of("1:a:z")));
+  }
+
+  /** Mirrors {@code cdc.sink.CdcSinkTestUtils#withKinds} (package-private 
there). */
+  private static PCollection<Row> withKinds(PCollection<KV<ValueKind, Row>> 
tagged) {
+    return tagged.apply(
+        "AttachKinds",
+        ParDo.of(
+            new DoFn<KV<ValueKind, Row>, Row>() {
+              @ProcessElement
+              public void process(@Element KV<ValueKind, Row> e, 
OutputReceiver<Row> out) {
+                out.builder(e.getValue()).setValueKind(e.getKey()).output();
+              }
+            }));
+  }
+
+  private static Row kindRow(Schema schema, int id, String name, String data, 
long seq) {
+    return Row.withSchema(schema).addValues(id, name, data, seq, 
"src").build();
+  }
+
+  /**
+   * The {@code only} projection extracts a nested payload row (the Debezium 
{@code after} pattern)
+   * while carrying the top-level control columns through to the sink.
+   */
+  @Test
+  public void onlyProjectionExtractsPayloadAndPreservesControlColumns() {
+    TableIdentifier id = v2Table();
+    Table table = catalog.loadTable(id);
+
+    Schema envelopeSchema = envelopeSchema(/* nullablePayload= */ false);
+
+    runCdcWrite(
+        configFor(id)
+            .setOnly("after")
+            .setMode("merge-on-read")
+            .setChangeTypeColumn("op")
+            .setChangeTypeMap(ImmutableMap.of("c", "INSERT", "u", 
"UPDATE_AFTER", "d", "DELETE"))
+            .setSequenceNumberColumn("seq")
+            .setUpsert(true)
+            .build(),
+        envelopeSchema,
+        envelopeRow(envelopeSchema, 1, "a", "x", "c", 1L),
+        envelopeRow(envelopeSchema, 2, "b", "y", "c", 1L),
+        envelopeRow(envelopeSchema, 1, "a", "z", "u", 2L),
+        envelopeRow(envelopeSchema, 2, "b", "y", "d", 2L));
+
+    // "u" upserted id=1, "d" deleted id=2; ts_ms and op never reached the 
table.
+    assertThat(readRows(table), equalTo(ImmutableList.of("1:a:z")));
+  }
+
+  /** A Debezium-ish envelope: {@code after} payload row + {@code op}/{@code 
seq}/{@code ts_ms}. */
+  private static Schema envelopeSchema(boolean nullablePayload) {
+    Schema.Builder builder = Schema.builder();
+    if (nullablePayload) {
+      builder.addNullableField("after", Schema.FieldType.row(DATA_SCHEMA));
+    } else {
+      builder.addRowField("after", DATA_SCHEMA);
+    }
+    return 
builder.addStringField("op").addInt64Field("seq").addInt64Field("ts_ms").build();
+  }
+
+  private static Row envelopeRow(
+      Schema envelopeSchema, int id, String name, String data, String op, long 
seq) {
+    return Row.withSchema(envelopeSchema)
+        .addValues(Row.withSchema(DATA_SCHEMA).addValues(id, name, 
data).build(), op, seq, 1234L)
+        .build();
+  }
+
+  /**
+   * {@code drop} removes extra Debezium-ish envelope columns on the templated 
path too, so the
+   * write succeeds and only the data columns land.
+   */
+  @Test
+  public void dynamicDestinationDropsExtraEnvelopeColumns() {
+    String suffix = "t" + System.nanoTime();
+    List<Table> tables = templateTables(suffix);
+
+    // Input carries two extra Debezium envelope columns not present in the 
table.
+    Schema envelopeSchema =
+        Schema.builder()
+            .addFields(INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ.getFields())
+            .addInt64Field("ts_ms")
+            .addStringField("source")
+            .build();
+
+    runCdcWrite(
+        templateConfig().setDrop(ImmutableList.of("ts_ms", "source")).build(),
+        envelopeSchema,
+        Row.withSchema(envelopeSchema)
+            .addValues(1, "a_" + suffix, "x", "INSERT", 1L, 1234L, "mysql")
+            .build(),
+        Row.withSchema(envelopeSchema)
+            .addValues(2, "b_" + suffix, "y", "INSERT", 1L, 5678L, "postgres")
+            .build());
+
+    // The extra envelope columns were dropped; only id:name:data landed.
+    assertThat(readRows(tables.get(0)), equalTo(ImmutableList.of("1:a_" + 
suffix + ":x")));
+    assertThat(readRows(tables.get(1)), equalTo(ImmutableList.of("2:b_" + 
suffix + ":y")));
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Destinations: single table and dynamic templates
+  // 
---------------------------------------------------------------------------------------------
+
+  /**
+   * A destination template routes by a data column, and each table sees its 
own changes applied:
+   * the sink consumes both control columns, so neither leaks into the written 
data.
+   */
+  @Test
+  public void dynamicDestinationTemplateRoutesAndAppliesChanges() {
+    String suffix = "t" + System.nanoTime();
+    List<Table> tables = templateTables(suffix);
+
+    PCollectionRowTuple output =
+        applyCdcWrite(
+            templateConfig().build(),
+            INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ,
+            rowWithSeq(1, "a_" + suffix, "x", "INSERT", 1L),
+            rowWithSeq(2, "b_" + suffix, "y", "INSERT", 1L),
+            rowWithSeq(1, "a_" + suffix, "z", "UPDATE_AFTER", 2L));
+    assertTrue(output.has(SNAPSHOTS_TAG));
+    assertTrue(output.has(DEAD_LETTER_TAG));
+
+    p.run().waitUntilFinish();
+
+    // id=1 reflects the seq=2 UPDATE_AFTER; no control column leaked into 
either table.
+    assertThat(readRows(tables.get(0)), equalTo(ImmutableList.of("1:a_" + 
suffix + ":z")));
+    assertThat(readRows(tables.get(1)), equalTo(ImmutableList.of("2:b_" + 
suffix + ":y")));
+  }
+
+  /**
+   * A destination template may reference a field the projection drops: 
routing sees the raw record
+   * while the written row is the filtered one, so a routing-only column never 
lands in the table.
+   */
+  @Test
+  public void templateRoutesByDroppedField() {
+    String suffix = "t" + System.nanoTime();
+    List<Table> tables = templateTables(suffix);
+    Schema inputSchema =
+        Schema.builder()
+            .addFields(INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ.getFields())
+            .addStringField("source")
+            .build();
+
+    runCdcWrite(
+        
templateConfig().setTable("db.{source}").setDrop(ImmutableList.of("source")).build(),
+        inputSchema,
+        Row.withSchema(inputSchema).addValues(1, "a", "x", "INSERT", 1L, "a_" 
+ suffix).build(),
+        Row.withSchema(inputSchema).addValues(2, "b", "y", "INSERT", 1L, "b_" 
+ suffix).build());
+
+    assertThat(readRows(tables.get(0)), equalTo(ImmutableList.of("1:a:x")));
+    assertThat(readRows(tables.get(1)), equalTo(ImmutableList.of("2:b:y")));
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Change kinds and defaults
+  // 
---------------------------------------------------------------------------------------------
+
+  @Test
+  public void schemaTransformAppliesCdcEndToEnd() {
+    TableIdentifier id = v2Table();
+    Table table = catalog.loadTable(id);
+
+    PCollectionRowTuple output =
+        applyCdcWrite(
+            cdcConfigFor(id).build(),
+            INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ,
+            rowWithSeq(1, "a", "x", "INSERT", 1L),
+            rowWithSeq(2, "b", "y", "INSERT", 1L),
+            rowWithSeq(1, "a", "x", "UPDATE_BEFORE", 2L),
+            rowWithSeq(1, "a", "z", "UPDATE_AFTER", 2L),
+            rowWithSeq(2, "b", "y", "DELETE", 2L));
+
+    // Output tags are pinned: snapshots + dead_letter, and no error output 
when error_handling
+    // is not configured.
+    assertTrue(output.has(SNAPSHOTS_TAG));
+    assertTrue(output.has(DEAD_LETTER_TAG));
+    assertFalse(output.has("errors"));
+
+    PAssert.that(output.get(DEAD_LETTER_TAG)).empty();
+
+    p.run().waitUntilFinish();
+
+    // id=1 updated to (a,z); id=2 deleted.
+    assertThat(readRows(table), equalTo(ImmutableList.of("1:a:z")));
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Error handling
+  // 
---------------------------------------------------------------------------------------------
+
+  /**
+   * With {@code error_handling}, an invalid record (unknown change-type 
value) is routed to the
+   * configured named output; the good records commit.
+   */
+  @Test
+  public void errorHandlingRoutesInvalidRecords() {
+    TableIdentifier id = v2Table();
+    Table table = catalog.loadTable(id);
+
+    PCollectionRowTuple output =
+        applyCdcWrite(
+            cdcConfigFor(id)
+                
.setErrorHandling(ErrorHandling.builder().setOutput("errors").build())
+                .build(),
+            INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ,
+            rowWithSeq(1, "a", "x", "INSERT", 1L),
+            rowWithSeq(2, "b", "y", "BOGUS", 1L),
+            rowWithSeq(3, "c", "z", "INSERT", 1L));
+    assertTrue(output.has("errors"));
+    assertTrue(output.has(SNAPSHOTS_TAG));
+    assertTrue(output.has(DEAD_LETTER_TAG));
+
+    PAssert.that(output.get("errors"))
+        .satisfies(
+            rows -> {
+              Row err = ImmutableList.copyOf(rows).get(0);
+              Row failedRow = err.getRow("failed_row");
+              assertNotNull(failedRow);
+              assertEquals(Integer.valueOf(2), failedRow.getInt32("id"));
+              return null;
+            });
+
+    p.run().waitUntilFinish();
+
+    assertThat(readRows(table), containsInAnyOrder("1:a:x", "3:c:z"));
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Managed
+  // 
---------------------------------------------------------------------------------------------
+
+  /**
+   * {@code Managed.write(Managed.ICEBERG)} with a {@code cdc} map writes CDC 
end-to-end through the
+   * snake_case config surface.
+   */
+  @Test
+  public void managedIcebergMergeOnReadWritesEndToEnd() {
+    TableIdentifier id = v2Table();
+    Table table = catalog.loadTable(id);
+
+    Map<String, Object> configMap =
+        ImmutableMap.<String, Object>builder()
+            .put("table", id.toString())
+            .put("catalog_properties", catalogProperties())
+            .put("mode", "merge-on-read")
+            .put("change_type_column", "change_type")
+            .build();
+
+    PCollection<Row> input =
+        p.apply(
+                Create.of(
+                    rowWithSeq(1, "a", "x", "INSERT", 1L),
+                    rowWithSeq(1, "a", "x", "UPDATE_BEFORE", 2L),
+                    rowWithSeq(1, "a", "z", "UPDATE_AFTER", 2L)))
+            .setRowSchema(INPUT_SCHEMA_WITH_CHANGE_TYPE_AND_SEQ);
+
+    PCollectionRowTuple output = 
input.apply(Managed.write(Managed.ICEBERG).withConfig(configMap));
+    assertTrue(output.has(SNAPSHOTS_TAG));
+    assertTrue(output.has(DEAD_LETTER_TAG));
+
+    p.run().waitUntilFinish();
+
+    assertThat(readRows(table), equalTo(ImmutableList.of("1:a:z")));
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/AssignCdcKeysTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/AssignCdcKeysTest.java
index 7d7ae19230e..232ca3d6bc8 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/AssignCdcKeysTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/AssignCdcKeysTest.java
@@ -47,6 +47,7 @@ import org.apache.beam.sdk.coders.RowCoder;
 import org.apache.beam.sdk.io.iceberg.DynamicDestinations;
 import org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig;
 import org.apache.beam.sdk.io.iceberg.IcebergUtils;
+import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
 import org.apache.beam.sdk.metrics.MetricNameFilter;
 import org.apache.beam.sdk.metrics.MetricResult;
 import org.apache.beam.sdk.metrics.MetricsFilter;
@@ -85,7 +86,7 @@ public class AssignCdcKeysTest {
   @Rule public transient TemporaryFolder tmp = new TemporaryFolder();
   @Rule public final TestName testName = new TestName();
 
-  private static final String SEQ_COL = 
CdcWriteConfig.DEFAULT_SEQUENCE_NUMBER_COLUMN;
+  private static final String SEQ_COL = 
IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER;
   private static final int NUM_SHARDS = 8;
 
   private static final org.apache.iceberg.Schema ICEBERG_SCHEMA =
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfigTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfigTest.java
index fb8ab133b49..dc55147eb8b 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfigTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfigTest.java
@@ -27,6 +27,7 @@ import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
+import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
 import org.apache.beam.sdk.util.SerializableUtils;
 import org.apache.beam.sdk.values.ValueKind;
 import org.junit.Test;
@@ -44,7 +45,8 @@ public class CdcWriteConfigTest {
     CdcWriteConfig config = 
CdcWriteConfig.builder().setSinkId(SINK_ID).build();
 
     assertThat(
-        config.getSequenceNumberColumn(), 
equalTo(CdcWriteConfig.DEFAULT_SEQUENCE_NUMBER_COLUMN));
+        config.getSequenceNumberColumn(),
+        equalTo(IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));
     assertThat(config.getNumShards(), 
equalTo(CdcWriteConfig.DEFAULT_NUM_SHARDS));
     // Unset shards_per_partition resolves to num_shards: the resolved int 
carries no cap.
     assertThat(config.getShardsPerPartition(), 
equalTo(CdcWriteConfig.DEFAULT_NUM_SHARDS));
diff --git a/sdks/python/apache_beam/yaml/yaml_io.py 
b/sdks/python/apache_beam/yaml/yaml_io.py
index b7ce98b3599..3fa15044b60 100644
--- a/sdks/python/apache_beam/yaml/yaml_io.py
+++ b/sdks/python/apache_beam/yaml/yaml_io.py
@@ -651,6 +651,20 @@ def write_to_iceberg(
     only: Optional[str] = None,
     distribution_mode: Optional[str] = None,
     autosharding: Optional[bool] = None,
+    mode: Optional[str] = None,
+    sequence_number_column: Optional[str] = None,
+    change_type_column: Optional[str] = None,
+    change_type_map: Optional[Mapping[str, str]] = None,
+    upsert: Optional[bool] = None,
+    equality_columns: Optional[Iterable[str]] = None,
+    num_shards: Optional[int] = None,
+    shards_per_partition: Optional[int] = None,
+    allowed_lateness_seconds: Optional[int] = None,
+    sink_id: Optional[str] = None,
+    token_heartbeat_seconds: Optional[int] = None,
+    snapshot_properties: Optional[Mapping[str, str]] = None,
+    error_handling: Optional[Mapping[str, Any]] = None,
+    sorter_memory_mb: Optional[int] = None,
 ):
   # TODO(robertwb): It'd be nice to derive this list of parameters, along with
   # their types and docs, programmatically from the iceberg (or managed)
@@ -709,6 +723,51 @@ def write_to_iceberg(
       further sub-dividing partitions into multiple shards to prevent
       bottlenecks during high-throughput writes. Only available with 'hash'
       distribution mode.
+    mode: Controls how rows are written. 'append' (default) appends every row
+      as new data. 'merge-on-read' treats each row as a change (INSERT,
+      UPDATE_BEFORE, UPDATE_AFTER, or DELETE) applied to the table by primary
+      key.
+    sequence_number_column: Merge-on-read only. the required column name
+      representing the monotonic sequence number used to order a single key's
+      changes. Default column name is '_commit_snapshot_sequence_number'. This
+      column will be stripped from the data row before writing to Iceberg.
+    change_type_column: Merge-on-read only. The optional column name
+      representing the row's change type (INSERT, UPDATE_BEFORE, UPDATE_AFTER,
+      or DELETE). This column will be stripped from the data row before writing
+      to Iceberg.
+    change_type_map: Merge-on-read only. Optional map from the
+      `change_type_column` value to the canonical change type name (see above).
+    upsert: Merge-on-read only. If true, only the after-image of each change
+      (INSERT/UPDATE_AFTER) is applied as an upsert. UPDATE_BEFORE records
+      are dropped. Default: false.
+    equality_columns: 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
+      'merge-on-read' mode.
+    num_shards: The number of deterministic primary-key-hash shards per
+      destination, i.e. the max write parallelism per destination. Defaults to
+      16. Currently only supported in 'merge-on-read' mode.
+    shards_per_partition: Maximum number of shards a single partition's rows
+      may occupy. Defaults to `num_shards`. Currently only supported in
+      'merge-on-read' mode.
+    allowed_lateness_seconds: How long a late record may lag behind the
+      watermark before it is dropped entirely. Defaults to 21600 (6 hours).
+      Currently only supported in 'merge-on-read' mode.
+    sink_id: A stable identifier namespacing the idempotency tokens written to
+      each commit. Defaults to a unique per-write UUID. Override with a stable
+      ID for exactly-once commits across relaunches of a particular streaming
+      write. Currently only relevant in 'merge-on-read' mode.
+    token_heartbeat_seconds: Streaming only. Refresh each idle
+      destination's commit token every this many seconds. Disabled by default.
+      Currently only relevant in 'merge-on-read' mode.
+    snapshot_properties: Extra properties added to every commit's snapshot
+      summary. Currently only supported in 'merge-on-read' mode.
+    error_handling: Where to send records the sink cannot apply (e.g. an
+      unknown change type) instead of failing the pipeline. Currently only
+      supported in 'merge-on-read' mode.
+    sorter_memory_mb: The in-memory buffer size (MB) for the pre-write sort.
+      Groups larger than this spill to disk. Default: 100MB. Currently only
+      relevant in 'merge-on-read' mode.
   """
   return beam.managed.Write(
       "iceberg",
@@ -724,7 +783,21 @@ def write_to_iceberg(
           drop=drop,
           only=only,
           distribution_mode=distribution_mode,
-          autosharding=autosharding))
+          autosharding=autosharding,
+          mode=mode,
+          sequence_number_column=sequence_number_column,
+          change_type_column=change_type_column,
+          change_type_map=change_type_map,
+          upsert=upsert,
+          equality_columns=equality_columns,
+          num_shards=num_shards,
+          shards_per_partition=shards_per_partition,
+          allowed_lateness_seconds=allowed_lateness_seconds,
+          sink_id=sink_id,
+          token_heartbeat_seconds=token_heartbeat_seconds,
+          snapshot_properties=snapshot_properties,
+          error_handling=error_handling,
+          sorter_memory_mb=sorter_memory_mb))
 
 
 def io_providers():

Reply via email to