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 f9064240c35 [Iceberg CDC sink] table setup (#40007)
f9064240c35 is described below

commit f9064240c35d6dcedb3eba409ce2ce2715e7db18
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Mon Sep 14 20:28:26 2026 +0200

    [Iceberg CDC sink] table setup (#40007)
    
    * table setup
    
    * spotless
    
    * add unit tests for CommitToken and PartitionShardPlan
---
 .../beam/sdk/io/iceberg/cdc/sink/CommitToken.java  |  314 ++++++
 .../io/iceberg/cdc/sink/PartitionShardPlan.java    |  128 +++
 .../beam/sdk/io/iceberg/cdc/sink/TableSetup.java   |  724 ++++++++++++++
 .../sdk/io/iceberg/cdc/sink/CommitTokenTest.java   |  173 ++++
 .../iceberg/cdc/sink/PartitionShardPlanTest.java   |  120 +++
 .../sdk/io/iceberg/cdc/sink/TableSetupTest.java    | 1028 ++++++++++++++++++++
 6 files changed, 2487 insertions(+)

diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitToken.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitToken.java
new file mode 100644
index 00000000000..793b7639da1
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitToken.java
@@ -0,0 +1,314 @@
+/*
+ * 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.cdc.sink;
+
+import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
+
+import java.io.Serializable;
+import java.util.Map;
+import org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig;
+import org.apache.beam.sdk.io.iceberg.TableCache;
+import org.apache.beam.sdk.metrics.Counter;
+import org.apache.iceberg.Snapshot;
+import org.apache.iceberg.SnapshotUpdate;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.exceptions.NoSuchTableException;
+import org.apache.iceberg.util.SnapshotUtil;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * The CDC sink's idempotency-token contract: the sink-id-namespaced 
snapshot-summary keys that make
+ * commits idempotent, plus every token-keyed ancestry walk the committer 
performs.
+ */
+final class CommitToken implements Serializable {
+
+  private static final Logger LOG = LoggerFactory.getLogger(CommitToken.class);
+
+  /** Marks a snapshot as committed by the CDC sink instance named by its 
value. */
+  static final String SINK_ID_KEY = "beam.cdc.sink-id";
+
+  /** Prefix of the committed-through window-end token key ({@code + sinkId}). 
*/
+  static final String COMMITTED_THROUGH_MS_PREFIX = 
"beam.cdc.committed-through-ms.";
+
+  /** Prefix of the max-committed source-sequence key ({@code + sinkId}). */
+  static final String MAX_COMMITTED_SEQ_PREFIX = "beam.cdc.max-committed-seq.";
+
+  /** Prefix of the run-spec stamp key ({@code + sinkId}); value {@code 
<runId>:<specId>}. */
+  static final String RUN_SPEC_PREFIX = "beam.cdc.run-spec.";
+
+  private final String sinkId;
+  private final String runId;
+  private final Counter tokenParseFailures;
+  private final Counter suspectedTokenExpiry;
+
+  /**
+   * @param runId the run runId stamped into the run-spec key
+   * @param tokenParseFailures counts unparseable token/max-seq summary values 
met during recovery
+   * @param suspectedTokenExpiry counts recoveries where the sink-id marker 
survives but no token
+   *     does (the token-bearing snapshots were likely expired away)
+   */
+  CommitToken(
+      String sinkId, String runId, Counter tokenParseFailures, Counter 
suspectedTokenExpiry) {
+    this.sinkId = sinkId;
+    this.runId = runId;
+    this.tokenParseFailures = tokenParseFailures;
+    this.suspectedTokenExpiry = suspectedTokenExpiry;
+  }
+
+  /** Writes the three token keys and the {@code pinnedSpecId} onto a pending 
snapshot operation. */
+  void writeTo(
+      SnapshotUpdate<?> op,
+      long committedThroughMs,
+      long maxCommittedSeq,
+      @Nullable Integer pinnedSpecId) {
+    op.set(COMMITTED_THROUGH_MS_PREFIX + sinkId, 
Long.toString(committedThroughMs));
+    op.set(MAX_COMMITTED_SEQ_PREFIX + sinkId, Long.toString(maxCommittedSeq));
+    if (pinnedSpecId != null) {
+      op.set(RUN_SPEC_PREFIX + sinkId, runId + ":" + pinnedSpecId);
+    }
+    op.set(SINK_ID_KEY, sinkId);
+  }
+
+  /**
+   * Writes the token keys for an idle token-refresh (heartbeat) commit. 
Unlike {@link #writeTo},
+   * the max-committed-seq key is omitted when unknown ({@code MIN}, meaning 
recovery found a token
+   * whose snapshot carried no parseable max-seq).
+   */
+  void writeHeartbeatTo(
+      SnapshotUpdate<?> op,
+      long committedThroughMs,
+      long maxCommittedSeq,
+      @Nullable Integer pinnedSpecId) {
+    op.set(COMMITTED_THROUGH_MS_PREFIX + sinkId, 
Long.toString(committedThroughMs));
+    if (maxCommittedSeq != Long.MIN_VALUE) {
+      op.set(MAX_COMMITTED_SEQ_PREFIX + sinkId, 
Long.toString(maxCommittedSeq));
+    }
+    if (pinnedSpecId != null) {
+      op.set(RUN_SPEC_PREFIX + sinkId, runId + ":" + pinnedSpecId);
+    }
+    op.set(SINK_ID_KEY, sinkId);
+  }
+
+  /**
+   * Returns the spec id stamped for {@code sinkId} under {@code runId}, read 
from the most recent
+   * stamp-bearing snapshot on {@code table}'s current branch. {@code null} 
when that stamp is
+   * absent, unparseable, or another run's (the caller falls back to the 
current spec).
+   */
+  static @Nullable Integer readRunSpec(Table table, String sinkId, String 
runId) {
+    Snapshot current = table.currentSnapshot();
+    if (current == null) {
+      return null;
+    }
+    String key = RUN_SPEC_PREFIX + sinkId;
+    String wantedPrefix = runId + ":";
+    for (Snapshot s : SnapshotUtil.ancestorsOf(current.snapshotId(), 
table::snapshot)) {
+      Map<String, String> summary = s.summary();
+      if (summary == null) {
+        continue;
+      }
+      String value = summary.get(key);
+      if (value == null) {
+        continue;
+      }
+      // Only the newest stamp counts; a foreign runId or garbage value reads 
as no stamp.
+      if (!value.startsWith(wantedPrefix)) {
+        return null;
+      }
+      try {
+        return Integer.parseInt(value.substring(wantedPrefix.length()));
+      } catch (NumberFormatException e) {
+        return null;
+      }
+    }
+    return null;
+  }
+
+  /**
+   * The state recovered from a table's ancestry: the committed-through-ms 
window token and the
+   * max-committed sequence from the same snapshot. The sequence seeds the 
cross-window inversion
+   * detector when committer state is empty, which is a relaunch or the first 
time a destination is
+   * seen. {@link #FRESH_START} means neither was found.
+   */
+  static final class Recovered {
+    static final Recovered FRESH_START = new Recovered(Long.MIN_VALUE, 
Long.MIN_VALUE);
+
+    final long committedThroughMs;
+    final long maxCommittedSeq;
+
+    private Recovered(long committedThroughMs, long maxCommittedSeq) {
+      this.committedThroughMs = committedThroughMs;
+      this.maxCommittedSeq = maxCommittedSeq;
+    }
+  }
+
+  /**
+   * Loads {@code dest} (forcing a refresh) and recovers this sink's token 
from its ancestry. The
+   * table may not exist yet, so a missing table is tolerated as a fresh start.
+   */
+  Recovered recoverFromTable(IcebergCatalogConfig catalogConfig, String dest) {
+    Table table;
+    try {
+      table = TableCache.getRefreshed(catalogConfig, dest);
+    } catch (RuntimeException e) {
+      if (hasCause(e, NoSuchTableException.class)) {
+        return Recovered.FRESH_START;
+      }
+      throw e;
+    }
+    return recoverFrom(table, dest);
+  }
+
+  /**
+   * Recovers this sink's committed-through-ms token (and corresponding 
max-committed sequence) by
+   * scanning a table's snapshot ancestry and returning the first {@code
+   * beam.cdc.committed-through-ms.<sinkId>} found, else {@link 
Long#MIN_VALUE}.
+   */
+  Recovered recoverFrom(Table table, String dest) {
+    Snapshot current = table.currentSnapshot();
+    if (current == null) {
+      return Recovered.FRESH_START;
+    }
+    String tokenKey = COMMITTED_THROUGH_MS_PREFIX + sinkId;
+    String maxSeqKey = MAX_COMMITTED_SEQ_PREFIX + sinkId;
+    boolean sawSinkMarker = false;
+    for (Snapshot s : SnapshotUtil.ancestorsOf(current.snapshotId(), 
table::snapshot)) {
+      Map<String, String> summary = s.summary();
+      if (summary == null) {
+        continue;
+      }
+      if (sinkId.equals(summary.get(SINK_ID_KEY))) {
+        sawSinkMarker = true;
+      }
+      String tokenValue = summary.get(tokenKey);
+      if (tokenValue == null) {
+        continue;
+      }
+      long committedThroughMs;
+      try {
+        committedThroughMs = Long.parseLong(tokenValue);
+      } catch (NumberFormatException e) {
+        // An older intact token is better than crash-looping
+        tokenParseFailures.inc();
+        LOG.error(
+            "CDC sink '{}' found an unparseable committed-through token '{}' 
in snapshot {} "
+                + "of table '{}'; ignoring it and scanning older ancestors.",
+            sinkId,
+            tokenValue,
+            s.snapshotId(),
+            dest);
+        continue;
+      }
+      // Both values come from this snapshot: the pair must describe one 
commit.
+      return new Recovered(committedThroughMs, 
parseMaxSeq(summary.get(maxSeqKey), s, dest));
+    }
+    if (sawSinkMarker) {
+      // This sink has committed to the table before, yet no token survived 
the ancestry scan.
+      // Rare but can happen if expire_snapshots removes the token-bearing 
snapshots.
+      suspectedTokenExpiry.inc();
+      LOG.warn(
+          "CDC sink '{}' found its sink-id marker in table '{}' ancestry but 
no "
+              + "committed-through token; the token-bearing snapshot(s) may 
have been expired. "
+              + "Falling back to MIN, which may replay retained windows.",
+          sinkId,
+          dest);
+    }
+    return Recovered.FRESH_START;
+  }
+
+  /** {@link Long#MIN_VALUE} when the max-committed-seq is absent or 
unparseable. */
+  private long parseMaxSeq(@Nullable String value, Snapshot s, String dest) {
+    if (value == null) {
+      return Long.MIN_VALUE;
+    }
+    try {
+      return Long.parseLong(value);
+    } catch (NumberFormatException e) {
+      tokenParseFailures.inc();
+      LOG.error(
+          "CDC sink '{}' found an unparseable max-committed-seq '{}' in 
snapshot {} of "
+              + "table '{}'; ignoring it.",
+          sinkId,
+          value,
+          s.snapshotId(),
+          dest);
+      return Long.MIN_VALUE;
+    }
+  }
+
+  /**
+   * Whether an idle destination should emit an empty token-refresh 
(heartbeat) commit: {@code true}
+   * iff the most recent table snapshot bearing this sink's committed-through 
token is older than
+   * {@code intervalMillis} relative to {@code nowMs}.
+   */
+  boolean shouldHeartbeat(Table table, long intervalMillis, long nowMs) {
+    @Nullable Snapshot current = table.currentSnapshot();
+    if (current == null) {
+      return false;
+    }
+    String tokenKey = COMMITTED_THROUGH_MS_PREFIX + sinkId;
+    for (Snapshot s : SnapshotUtil.ancestorsOf(current.snapshotId(), 
table::snapshot)) {
+      Map<String, String> summary = s.summary();
+      if (summary != null && summary.get(tokenKey) != null) {
+        return s.timestampMillis() < nowMs - intervalMillis;
+      }
+    }
+    return false;
+  }
+
+  /**
+   * Finds and returns the snapshot corresponding to a just-committed window 
by looking for the
+   * specified {@code windowEndMs}. Expects that the caller has just committed 
the window, so throws
+   * if no such snapshot exists.
+   */
+  Snapshot findRecentlyCommittedTokenSnapshot(Table table, String dest, long 
windowEndMs) {
+    table.refresh();
+    Snapshot current =
+        checkStateNotNull(
+            table.currentSnapshot(),
+            "table '%s' has no current snapshot right after a commit",
+            dest);
+    String tokenKey = COMMITTED_THROUGH_MS_PREFIX + sinkId;
+    String wanted = Long.toString(windowEndMs);
+    for (Snapshot s : SnapshotUtil.ancestorsOf(current.snapshotId(), 
table::snapshot)) {
+      Map<String, String> summary = s.summary();
+      if (summary != null && wanted.equals(summary.get(tokenKey))) {
+        return s;
+      }
+    }
+    throw new IllegalStateException(
+        "CDC sink '"
+            + sinkId
+            + "' committed window-end "
+            + windowEndMs
+            + " ms to table '"
+            + dest
+            + "' but found no snapshot carrying its committed-through token in 
the refreshed "
+            + "ancestry.");
+  }
+
+  private static boolean hasCause(Throwable t, Class<? extends Throwable> 
type) {
+    for (Throwable cause = t; cause != null; cause = cause.getCause()) {
+      if (type.isInstance(cause)) {
+        return true;
+      }
+    }
+    return false;
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlan.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlan.java
new file mode 100644
index 00000000000..73fbd334ab6
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlan.java
@@ -0,0 +1,128 @@
+/*
+ * 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.cdc.sink;
+
+import java.util.ArrayList;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Set;
+import org.apache.beam.sdk.io.iceberg.IcebergUtils;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.values.Row;
+import org.apache.iceberg.PartitionField;
+import org.apache.iceberg.PartitionKey;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.StructLike;
+import org.apache.iceberg.data.InternalRecordWrapper;
+import org.apache.iceberg.types.JavaHash;
+import org.apache.iceberg.types.TypeUtil;
+import org.apache.iceberg.types.Types;
+import org.checkerframework.checker.nullness.qual.Nullable;
+
+/**
+ * Derives a record's write shard from its Iceberg partition tuple. Each 
partition owns a block of
+ * {@code shards_per_partition} consecutive shards, and the caller's {@code 
offset} (derived from
+ * the primary-key hash) selects one of the shards. Used when {@code 
shards_per_partition <
+ * num_shards}.
+ *
+ * <p>Correctness rests on one property: this plan exists only under a {@code 
shards_per_partition}
+ * cap below {@code num_shards}, where {@link TableSetup#validatePartitioning} 
still requires
+ * partition source columns to be equality columns, so the shard is a pure 
function of the primary
+ * key and one key's records never split across shards.
+ */
+final class PartitionShardPlan {
+
+  /** Beam schema of the partition source columns, in the projected Iceberg 
schema's order. */
+  private final Schema sourceSchema;
+
+  /** For each {@link #sourceSchema} field, its position in the CDC data 
schema. */
+  private final int[] sourcePositions;
+
+  /** Iceberg schema of the partition source columns. */
+  private final org.apache.iceberg.Schema sourceIcebergSchema;
+
+  /** Adapts a converted record to the internal representation the transforms 
expect. */
+  private final InternalRecordWrapper wrapper;
+
+  /** The spec's bound transforms over a reused partition tuple. */
+  private final PartitionKey partitionKey;
+
+  /** Type-aware, JVM-stable hash of the partition tuple. */
+  private final JavaHash<StructLike> partitionHash;
+
+  private PartitionShardPlan(
+      Schema sourceSchema,
+      int[] sourcePositions,
+      org.apache.iceberg.Schema sourceIcebergSchema,
+      InternalRecordWrapper wrapper,
+      PartitionKey partitionKey,
+      JavaHash<StructLike> partitionHash) {
+    this.sourceSchema = sourceSchema;
+    this.sourcePositions = sourcePositions;
+    this.sourceIcebergSchema = sourceIcebergSchema;
+    this.wrapper = wrapper;
+    this.partitionKey = partitionKey;
+    this.partitionHash = partitionHash;
+  }
+
+  /** Builds the plan for a partitioned spec. Converts only the partition 
source columns. */
+  static PartitionShardPlan of(
+      PartitionSpec spec, org.apache.iceberg.Schema tableSchema, Schema 
cdcDataSchema) {
+    // Find distinct source ids since one column can feed several partition 
fields
+    Set<Integer> sourceIds = new LinkedHashSet<>();
+    for (PartitionField field : spec.fields()) {
+      sourceIds.add(field.sourceId());
+    }
+    org.apache.iceberg.Schema sourceIcebergSchema = 
TypeUtil.select(tableSchema, sourceIds);
+
+    List<Types.NestedField> sourceColumns = sourceIcebergSchema.columns();
+    Schema.Builder sourceBeamSchemaBuilder = Schema.builder();
+    int[] sourcePositions = new int[sourceColumns.size()];
+    // convert to a Beam schema using input data schema fields
+    for (int i = 0; i < sourceColumns.size(); i++) {
+      String name = sourceColumns.get(i).name();
+      sourceBeamSchemaBuilder.addField(cdcDataSchema.getField(name));
+      sourcePositions[i] = cdcDataSchema.indexOf(name);
+    }
+    Schema sourceBeamSchema = sourceBeamSchemaBuilder.build();
+
+    return new PartitionShardPlan(
+        sourceBeamSchema,
+        sourcePositions,
+        sourceIcebergSchema,
+        new InternalRecordWrapper(sourceIcebergSchema.asStruct()),
+        new PartitionKey(spec, sourceIcebergSchema),
+        JavaHash.forType(spec.partitionType()));
+  }
+
+  /**
+   * Computes the shard for {@code data}: the partition tuple's hash picks the 
block base, and
+   * {@code offset} (in {@code [0, shardsPerPartition)}) selects the shard 
within the block.
+   */
+  int shardFor(Row data, int offset, int numShards) {
+    List<@Nullable Object> values = new ArrayList<>(sourcePositions.length);
+    for (int position : sourcePositions) {
+      values.add(data.getValue(position));
+    }
+    Row sourceRow = Row.withSchema(sourceSchema).attachValues(values);
+    partitionKey.partition(
+        wrapper.wrap(IcebergUtils.beamRowToIcebergRecord(sourceIcebergSchema, 
sourceRow)));
+    int base = TableSetup.shardForHash(partitionHash.hash(partitionKey), 
numShards);
+    return Math.floorMod(base + offset, numShards);
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetup.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetup.java
new file mode 100644
index 00000000000..a0775c0bb43
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetup.java
@@ -0,0 +1,724 @@
+/*
+ * 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.cdc.sink;
+
+import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
+
+import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
+import org.apache.beam.sdk.coders.Coder;
+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.IcebergDestination;
+import org.apache.beam.sdk.io.iceberg.IcebergTableCreateConfig;
+import org.apache.beam.sdk.io.iceberg.IcebergUtils;
+import org.apache.beam.sdk.io.iceberg.TableCache;
+import org.apache.beam.sdk.schemas.Schema;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing;
+import org.apache.iceberg.PartitionField;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.SortOrder;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.TableUtil;
+import org.apache.iceberg.catalog.Catalog;
+import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.SupportsNamespaces;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.exceptions.AlreadyExistsException;
+import org.apache.iceberg.exceptions.NoSuchTableException;
+import org.apache.iceberg.types.Types;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Per-destination table resolution for the CDC sink: loads (or auto-creates) 
the destination {@link
+ * Table}, validates that it can accept CDC writes, and precomputes the 
per-destination artifacts
+ * the write path needs. All catalog I/O and table-level validation lives 
here; failures are thrown
+ * as {@link TableConfigException}.
+ *
+ * <p>One instance lives inside each worker {@code DoFn}, single-owner and not 
thread-safe. Results
+ * are memoized per destination string; the memo is what pins a destination's 
resolution, including
+ * its {@link Dest#spec()}, for the worker's lifetime.
+ */
+final class TableSetup implements Serializable {
+
+  private static final Logger LOG = LoggerFactory.getLogger(TableSetup.class);
+
+  private final IcebergCatalogConfig catalogConfig;
+  private final CdcWriteConfig config;
+  private final DynamicDestinations destinations;
+
+  /** The run runId; a run-spec stamp carrying it names the spec {@link 
#resolve} pins to. */
+  private final String runId;
+
+  /**
+   * Matches {@link TableCache}'s bound so this memo (whose {@link Dest}s 
strongly reference their
+   * {@link Table}s) can never pin more table metadata than that cache would 
hold. Exceeding it
+   * costs a re-resolve, not a failure.
+   */
+  private static final int MAX_MEMOIZED_DESTS = 1000;
+
+  /** Per-destination memo, lazily initialized (never serialized). */
+  private transient @Nullable Map<String, Dest> dests;
+
+  TableSetup(
+      IcebergCatalogConfig catalogConfig,
+      CdcWriteConfig config,
+      DynamicDestinations destinations,
+      String runId) {
+    this.catalogConfig = catalogConfig;
+    this.config = config;
+    this.destinations = destinations;
+    this.runId = runId;
+  }
+
+  /**
+   * Returns the resolved, validated {@link Dest} for {@code destString}, 
memoized per destination
+   * string. Under block sharding a memoized {@link Dest} is re-checked 
against the table's live
+   * partition spec before it is handed back ({@link 
#requireResolvedPartitionSpec}).
+   *
+   * @throws TableConfigException for any table-level problem (including 
catalog failures)
+   */
+  Dest get(String destString, Schema dataSchema) {
+    Map<String, Dest> memo = dests;
+    if (memo == null) {
+      // Access-ordered LRU: this class is single-threaded by contract, so a 
LinkedHashMap is the
+      // whole mechanism needed to keep the memo bounded.
+      memo =
+          new LinkedHashMap<String, Dest>(16, 0.75f, true) {
+            @Override
+            protected boolean removeEldestEntry(Map.Entry<String, Dest> 
eldest) {
+              return size() > MAX_MEMOIZED_DESTS;
+            }
+          };
+      dests = memo;
+    }
+    @Nullable Dest existing = memo.get(destString);
+    if (existing != null) {
+      requireResolvedPartitionSpec(destString, existing);
+      return existing;
+    }
+    Dest dest;
+    try {
+      dest = resolve(destString, dataSchema);
+    } catch (TableConfigException e) {
+      throw e;
+    } catch (RuntimeException e) {
+      throw new TableConfigException("Failed resolving destination table '" + 
destString + "'.", e);
+    }
+    memo.put(destString, dest);
+    return dest;
+  }
+
+  /** Type-stable murmur3 hash of encoded primary-key bytes. Stable across 
workers and JVMs. */
+  static int pkHash(byte[] pkBytes) {
+    return Hashing.murmur3_32_fixed().hashBytes(pkBytes).asInt();
+  }
+
+  /**
+   * The deterministic shard for {@code pkBytes}: {@code 
floorMod(murmur3_32(pkBytes), numShards)}.
+   */
+  static int shardFor(byte[] pkBytes, int numShards) {
+    return Math.floorMod(pkHash(pkBytes), numShards);
+  }
+
+  /**
+   * The deterministic shard for an already-computed value hash, used by 
{@link PartitionShardPlan}
+   * for the partition tuple's block base: {@code floorMod(murmur3_32(hash), 
numShards)}.
+   */
+  static int shardForHash(int hash, int numShards) {
+    return Math.floorMod(Hashing.murmur3_32_fixed().hashInt(hash).asInt(), 
numShards);
+  }
+
+  private Dest resolve(String destString, Schema dataSchema) {
+    TableIdentifier identifier = IcebergUtils.parseTableIdentifier(destString);
+    Table table =
+        TableCache.getAndRefreshIfStale(
+            catalogConfig, identifier, () -> loadOrCreateTable(identifier, 
destString, dataSchema));
+
+    int formatVersion = TableUtil.formatVersion(table);
+    if (formatVersion < 2) {
+      throw new TableConfigException(
+          "CDC sink requires an Iceberg format-version >= 2 table, but '"
+              + destString
+              + "' is format-version "
+              + formatVersion
+              + ". Use IcebergIO.writeRows (the append sink) for V1 tables.");
+    }
+
+    org.apache.iceberg.Schema tableSchema = table.schema();
+    Schema cdcDataSchema = cdcDataSchema(destString, dataSchema, tableSchema);
+    Set<Integer> equalityFieldIds = equalityFieldIds(destString, tableSchema);
+    requireNonNullableEqualityFields(destString, tableSchema, 
equalityFieldIds);
+    Schema pkSchema = pkSchema(tableSchema, cdcDataSchema, equalityFieldIds);
+
+    RowCoder pkCoder = RowCoder.of(pkSchema);
+    try {
+      pkCoder.verifyDeterministic();
+    } catch (Coder.NonDeterministicException e) {
+      throw new TableConfigException(
+          "Primary-key coder for table '"
+              + destString
+              + "' (schema "
+              + pkSchema
+              + ") is not deterministic; Iceberg identifier fields should be 
primitive and "
+              + "required.",
+          e);
+    }
+
+    PartitionSpec spec = runSpec(table);
+    validatePartitioning(destString, spec, tableSchema, equalityFieldIds);
+
+    int[] pkFieldPositions = new int[pkSchema.getFieldCount()];
+    for (int i = 0; i < pkSchema.getFieldCount(); i++) {
+      pkFieldPositions[i] = 
cdcDataSchema.indexOf(pkSchema.getField(i).getName());
+    }
+
+    return new Dest(
+        table,
+        spec,
+        ImmutableSet.copyOf(equalityFieldIds),
+        pkSchema,
+        pkCoder,
+        pkFieldPositions,
+        cdcDataSchema,
+        partitionShardPlan(destString, spec, tableSchema, cdcDataSchema));
+  }
+
+  /**
+   * The spec to resolve and validate against. Resolves to the most recently 
stamped spec in this
+   * run. Otherwise, falls back to the live {@code table.spec()}.
+   */
+  private PartitionSpec runSpec(Table table) {
+    @Nullable Integer stamped = CommitToken.readRunSpec(table, 
config.getSinkId(), runId);
+    if (stamped != null) {
+      @Nullable PartitionSpec pinned = table.specs().get(stamped);
+      if (pinned != null) {
+        return pinned;
+      }
+    }
+    return table.spec();
+  }
+
+  /**
+   * The destination's {@link PartitionShardPlan} when {@code 
shards_per_partition} is below {@code
+   * num_shards} AND the table is partitioned, else {@code null} (plain 
primary-key sharding). A cap
+   * on an unpartitioned destination is a no-op with a WARN, not a rejection: 
rejecting would fail a
+   * whole dynamic-destinations pipeline over one table, and the fallback is 
what the operator wants
+   * anyway.
+   */
+  private @Nullable PartitionShardPlan partitionShardPlan(
+      String destString,
+      PartitionSpec spec,
+      org.apache.iceberg.Schema tableSchema,
+      Schema cdcDataSchema) {
+    if (config.getShardsPerPartition() >= config.getNumShards()) {
+      return null;
+    }
+    if (spec.isUnpartitioned()) {
+      LOG.warn(
+          "shards_per_partition ({}) is below num_shards ({}) but destination 
'{}' is "
+              + "unpartitioned, so there is no partition to bound; ignoring 
the cap and sharding "
+              + "by primary key across num_shards shards. The option only 
helps partitioned "
+              + "tables.",
+          config.getShardsPerPartition(),
+          config.getNumShards(),
+          destString);
+      return null;
+    }
+    return PartitionShardPlan.of(spec, tableSchema, cdcDataSchema);
+  }
+
+  /**
+   * Block-sharding-only drift check, run on every memo hit. Under {@code 
shards_per_partition} the
+   * assigners derive each record's shard from the partition tuple under the 
resolved spec, so
+   * workers resolving different specs would split one primary key's window 
across shards: silent
+   * same-commit duplicates nothing downstream detects. The default path needs 
no check: the write
+   * path builds writers from the pinned {@link Dest#spec()}, never the live 
{@code table.spec()}.
+   * Freshness is best-effort: the compared spec is the process-cached 
table's, refreshed only when
+   * something in the process refreshes it.
+   */
+  private static void requireResolvedPartitionSpec(String destString, Dest 
dest) {
+    if (dest.partitionShardPlan() == null) {
+      return;
+    }
+    int currentSpecId = dest.table().spec().specId();
+    if (currentSpecId != dest.spec().specId()) {
+      throw new TableConfigException(
+          "Table '"
+              + destString
+              + "' changed its partition spec while the CDC sink was running 
with a "
+              + "shards_per_partition cap (spec id "
+              + dest.spec().specId()
+              + " when the sink resolved the table, spec id "
+              + currentSpecId
+              + " now). Partition-block sharding derives each record's shard 
from the partition "
+              + "tuple under the resolved spec, so workers resolving different 
specs would split "
+              + "one primary key's window across shards and silently duplicate 
rows within a "
+              + "commit. "
+              + "Drain the pipeline before evolving the partition spec, and 
restart it "
+              + "afterwards.");
+    }
+  }
+
+  /**
+   * Loads the table, auto-creating it (namespace first) if it does not exist. 
Auto-creation
+   * requires configured equality columns (a brand-new table has no identifier 
fields to infer
+   * from); the created schema is the data schema with the equality columns as 
identifier fields,
+   * honoring the destination's {@link IcebergTableCreateConfig} plus a 
format-version 2 default.
+   */
+  private Table loadOrCreateTable(
+      TableIdentifier identifier, String destString, Schema createSchema) {
+    Catalog catalog = catalogConfig.catalog();
+    try {
+      return catalog.loadTable(identifier);
+    } catch (NoSuchTableException e) {
+      // Missing table: fall through to auto-create (parity with the append 
sink).
+    }
+
+    @Nullable List<String> equalityColumns = config.getEqualityColumns();
+    if (equalityColumns == null || equalityColumns.isEmpty()) {
+      throw new TableConfigException(
+          "Table '"
+              + destString
+              + "' does not exist and no equality_columns are configured, so 
its identifier "
+              + "(primary-key) fields cannot be determined for auto-creation. 
Configure "
+              + "equality_columns, or pre-create the table with identifier 
fields.");
+    }
+
+    org.apache.iceberg.Schema base = 
IcebergUtils.beamSchemaToIcebergSchema(createSchema);
+    Set<Integer> identifierFieldIds = new LinkedHashSet<>();
+    for (String column : equalityColumns) {
+      requireTopLevelEqualityColumn(column);
+      Types.NestedField field = base.findField(column);
+      if (field == null) {
+        throw new TableConfigException(
+            "Cannot auto-create table '"
+                + destString
+                + "': equality column '"
+                + column
+                + "' is not present in the input data schema "
+                + createSchema
+                + ".");
+      }
+      // Iceberg refuses an OPTIONAL identifier field with a cryptic error; 
detect the nullable
+      // Beam field here so the message names the input field the user 
controls.
+      if (createSchema.getField(column).getType().getNullable()) {
+        throw new TableConfigException(
+            "Cannot auto-create table '"
+                + destString
+                + "': equality column '"
+                + column
+                + "' must be non-nullable in the input schema (a nullable 
column cannot be an "
+                + "Iceberg identifier field). Make the input field 
non-nullable, or pre-create "
+                + "the table with required identifier fields.");
+      }
+      identifierFieldIds.add(field.fieldId());
+    }
+    org.apache.iceberg.Schema schemaWithIds =
+        new org.apache.iceberg.Schema(base.columns(), identifierFieldIds);
+
+    IcebergDestination destination = 
destinations.instantiateDestination(destString);
+    @Nullable IcebergTableCreateConfig createConfig = 
destination.getTableCreateConfig();
+    PartitionSpec partitionSpec =
+        createConfig != null ? createConfig.getPartitionSpec() : 
PartitionSpec.unpartitioned();
+    SortOrder sortOrder = createConfig != null ? createConfig.getSortOrder() : 
SortOrder.unsorted();
+    Map<String, String> properties = new HashMap<>();
+    if (createConfig != null) {
+      @Nullable Map<String, String> createProperties = 
createConfig.getTableProperties();
+      if (createProperties != null) {
+        properties.putAll(createProperties);
+      }
+    }
+    properties.putIfAbsent(TableProperties.FORMAT_VERSION, "2");
+
+    Namespace namespace = identifier.namespace();
+    if (!namespace.isEmpty() && catalog instanceof SupportsNamespaces) {
+      SupportsNamespaces supportsNamespaces = (SupportsNamespaces) catalog;
+      if (!supportsNamespaces.namespaceExists(namespace)) {
+        try {
+          supportsNamespaces.createNamespace(namespace);
+          LOG.info("Created new namespace '{}'.", namespace);
+        } catch (AlreadyExistsException ignored) {
+          // Race: another worker created the namespace first.
+        }
+      }
+    }
+
+    try {
+      Table table =
+          catalog
+              .buildTable(identifier, schemaWithIds)
+              .withPartitionSpec(partitionSpec)
+              .withSortOrder(sortOrder)
+              .withProperties(properties)
+              .create();
+      LOG.info(
+          "CDC sink auto-created table '{}' with schema {}, partition spec {}, 
sort order {}, "
+              + "properties {}.",
+          identifier,
+          schemaWithIds,
+          partitionSpec,
+          sortOrder,
+          properties);
+      return table;
+    } catch (AlreadyExistsException ignored) {
+      // Race: another worker created the table first.
+      return catalog.loadTable(identifier);
+    }
+  }
+
+  /**
+   * Returns the table schema as a Beam {@link Schema}, validating that the 
data schema matches the
+   * table's top-level column names exactly AND in the same order. Order 
matters: the written rows
+   * and the shuffle coder are built positionally, so a column reorder would 
silently write values
+   * into the wrong columns.
+   */
+  private Schema cdcDataSchema(
+      String destString, Schema dataSchema, org.apache.iceberg.Schema 
tableSchema) {
+    Schema canonical = IcebergUtils.icebergSchemaToBeamSchema(tableSchema);
+    List<String> dataNames = dataSchema.getFieldNames();
+    List<String> canonicalNames = canonical.getFieldNames();
+    if (!dataNames.equals(canonicalNames)) {
+      throw new TableConfigException(schemaMismatchMessage(destString, 
canonicalNames, dataNames));
+    }
+    requireMatchingColumnTypes(destString, canonical, dataSchema);
+    return canonical;
+  }
+
+  /**
+   * Column-by-column type and nullability check behind the name check: rows 
are encoded against the
+   * table-derived schema, so a mismatched type would only fail later as an 
opaque coder error. A
+   * non-null input column on an optional table column is fine; the reverse is 
not.
+   */
+  private static void requireMatchingColumnTypes(
+      String destString, Schema canonical, Schema dataSchema) {
+    for (int i = 0; i < canonical.getFieldCount(); i++) {
+      String name = canonical.getField(i).getName();
+      Schema.FieldType tableType = canonical.getField(i).getType();
+      Schema.FieldType inputType = dataSchema.getField(i).getType();
+      if 
(!tableType.withNullable(false).equals(inputType.withNullable(false))) {
+        throw new TableConfigException(
+            "CDC data schema mismatch for table '"
+                + destString
+                + "': column '"
+                + name
+                + "' is "
+                + inputType
+                + " in the input but "
+                + tableType
+                + " in the table. Align the input schema with the table.");
+      }
+      if (inputType.getNullable() && !tableType.getNullable()) {
+        throw new TableConfigException(
+            "CDC data schema mismatch for table '"
+                + destString
+                + "': column '"
+                + name
+                + "' is nullable in the input but required in the table. Align 
the input schema "
+                + "with the table.");
+      }
+    }
+  }
+
+  /** The mismatch message: unexpected/missing columns, or the order 
difference. */
+  private static String schemaMismatchMessage(
+      String destString, List<String> canonicalNames, List<String> dataNames) {
+    Set<String> unexpected = new LinkedHashSet<>(dataNames);
+    unexpected.removeAll(canonicalNames);
+    Set<String> missing = new LinkedHashSet<>(canonicalNames);
+    missing.removeAll(dataNames);
+    StringBuilder msg =
+        new StringBuilder("CDC data schema mismatch for table 
'").append(destString).append("':");
+    if (!unexpected.isEmpty()) {
+      msg.append(" unexpected columns (in the input, not in the table): ")
+          .append(unexpected)
+          .append(";");
+    }
+    if (!missing.isEmpty()) {
+      msg.append(" missing columns (in the table, not supplied by the input): 
")
+          .append(missing)
+          .append(";");
+    }
+    if (unexpected.isEmpty() && missing.isEmpty()) {
+      msg.append(
+          " the input's data columns match the table's columns but in a 
different order;"
+              + " column order must match the table (rows are projected and 
encoded"
+              + " positionally);");
+    }
+    msg.append(" Table columns: ")
+        .append(canonicalNames)
+        .append("; input data columns: ")
+        .append(dataNames)
+        .append(".");
+    return msg.toString();
+  }
+
+  /**
+   * Equality columns must be top-level: Iceberg resolves dotted paths to 
nested fields, which are
+   * out of scope as identifier columns (and a same-named leaf could silently 
misbind).
+   */
+  private static void requireTopLevelEqualityColumn(String name) {
+    if (name.contains(".")) {
+      throw new TableConfigException(
+          "equality_columns must be top-level columns; got '"
+              + name
+              + "' (nested fields are not supported).");
+    }
+  }
+
+  /**
+   * The Iceberg field ids that define a row's identity: the configured 
equality columns (resolved
+   * by name) when set, else the table's identifier fields.
+   */
+  private Set<Integer> equalityFieldIds(String destString, 
org.apache.iceberg.Schema tableSchema) {
+    @Nullable List<String> override = config.getEqualityColumns();
+    if (override != null) {
+      if (override.isEmpty()) {
+        // An empty override is a misconfiguration, not a request for the 
identifier fields.
+        throw new TableConfigException(
+            "equality_columns must be non-empty or unset (leave unset to use 
the identifier "
+                + "fields of table '"
+                + destString
+                + "').");
+      }
+      ImmutableSet.Builder<Integer> ids = ImmutableSet.builder();
+      for (String name : override) {
+        requireTopLevelEqualityColumn(name);
+        Types.NestedField field = tableSchema.findField(name);
+        if (field == null) {
+          throw new TableConfigException(
+              "Configured equality column '"
+                  + name
+                  + "' does not exist in table '"
+                  + destString
+                  + "'. Table columns: "
+                  + columnNames(tableSchema)
+                  + ".");
+        }
+        ids.add(field.fieldId());
+      }
+      return ids.build();
+    }
+    Set<Integer> identifierFieldIds = tableSchema.identifierFieldIds();
+    if (identifierFieldIds.isEmpty()) {
+      throw new TableConfigException(
+          "Table '"
+              + destString
+              + "' has no identifier (primary-key) fields and no 
equality_columns are "
+              + "configured. Configure equality_columns, or add identifier 
fields to the table.");
+    }
+    return identifierFieldIds;
+  }
+
+  /** Equality columns must be required (non-null): a nullable column cannot 
define row identity. */
+  private static void requireNonNullableEqualityFields(
+      String destString, org.apache.iceberg.Schema tableSchema, Set<Integer> 
equalityFieldIds) {
+    for (int fieldId : equalityFieldIds) {
+      Types.NestedField field = 
checkStateNotNull(tableSchema.findField(fieldId));
+      if (!field.isRequired()) {
+        throw new TableConfigException(
+            "Equality column '"
+                + field.name()
+                + "' (field id "
+                + fieldId
+                + ") of table '"
+                + destString
+                + "' must be required (non-null); a nullable column cannot 
define row identity.");
+      }
+    }
+  }
+
+  /**
+   * The Beam schema of the equality columns, in ascending Iceberg field-id 
order (a stable,
+   * table-derived order independent of how the identifier fields or overrides 
were declared).
+   */
+  private static Schema pkSchema(
+      org.apache.iceberg.Schema tableSchema, Schema cdcDataSchema, 
Set<Integer> equalityFieldIds) {
+    Schema.Builder builder = Schema.builder();
+    for (int fieldId : new TreeSet<>(equalityFieldIds)) {
+      Types.NestedField field = 
checkStateNotNull(tableSchema.findField(fieldId));
+      builder.addField(cdcDataSchema.getField(field.name()));
+    }
+    return builder.build();
+  }
+
+  /**
+   * Tables may be partitioned on any columns; two options additionally 
require every partition
+   * source field to be an equality field, because they need a row's partition 
to be a pure function
+   * of its primary key: {@code upsert} (before-images are dropped, so a moved 
row's equality delete
+   * could only ever route to its new partition) and a {@code 
shards_per_partition} cap (the shard
+   * is derived from the partition tuple).
+   */
+  private void validatePartitioning(
+      String destString,
+      PartitionSpec spec,
+      org.apache.iceberg.Schema tableSchema,
+      Set<Integer> equalityFieldIds) {
+    if (spec.isUnpartitioned()) {
+      return;
+    }
+    String requirement;
+    if (config.getUpsert()) {
+      requirement =
+          "upsert drops before-images, so a row that moved partitions could 
never be deleted "
+              + "from its old partition";
+    } else if (config.getShardsPerPartition() < config.getNumShards()) {
+      requirement =
+          "shards_per_partition ("
+              + config.getShardsPerPartition()
+              + ") is below num_shards ("
+              + config.getNumShards()
+              + ") and derives each record's shard from its partition tuple, 
which must therefore "
+              + "be a pure function of the primary key";
+    } else {
+      return;
+    }
+    List<String> nonKeySources = new ArrayList<>();
+    for (PartitionField field : spec.fields()) {
+      if (!equalityFieldIds.contains(field.sourceId())) {
+        nonKeySources.add("'" + tableSchema.findColumnName(field.sourceId()) + 
"'");
+      }
+    }
+    if (!nonKeySources.isEmpty()) {
+      throw new TableConfigException(
+          "Table '"
+              + destString
+              + "' has partition source columns "
+              + nonKeySources
+              + " that are not equality columns, but "
+              + requirement
+              + ". Partition only on equality columns, or drop the option.");
+    }
+  }
+
+  private static List<String> columnNames(org.apache.iceberg.Schema 
tableSchema) {
+    List<String> names = new ArrayList<>(tableSchema.columns().size());
+    for (Types.NestedField field : tableSchema.columns()) {
+      names.add(field.name());
+    }
+    return names;
+  }
+
+  /** Precomputed per-destination state, fixed at resolution time. */
+  static final class Dest {
+    private final Table table;
+    private final PartitionSpec spec;
+    private final Set<Integer> equalityFieldIds;
+    private final Schema pkSchema;
+    private final RowCoder pkCoder;
+    private final int[] pkFieldPositions;
+    private final Schema cdcDataSchema;
+    private final @Nullable PartitionShardPlan partitionShardPlan;
+
+    private Dest(
+        Table table,
+        PartitionSpec spec,
+        Set<Integer> equalityFieldIds,
+        Schema pkSchema,
+        RowCoder pkCoder,
+        int[] pkFieldPositions,
+        Schema cdcDataSchema,
+        @Nullable PartitionShardPlan partitionShardPlan) {
+      this.table = table;
+      this.spec = spec;
+      this.equalityFieldIds = equalityFieldIds;
+      this.pkSchema = pkSchema;
+      this.pkCoder = pkCoder;
+      this.pkFieldPositions = pkFieldPositions;
+      this.cdcDataSchema = cdcDataSchema;
+      this.partitionShardPlan = partitionShardPlan;
+    }
+
+    /**
+     * The destination table, loaded or auto-created. This is the 
process-shared {@link TableCache}
+     * instance, which may be refreshed in place, so its live metadata can be 
newer than the
+     * memoized schemas held here.
+     */
+    Table table() {
+      return table;
+    }
+
+    /**
+     * The partition spec this destination was resolved and validated against: 
the worker's pin. The
+     * live {@link #table()} can be refreshed onto a newer spec; the write 
path must build writers
+     * from this pinned spec.
+     */
+    PartitionSpec spec() {
+      return spec;
+    }
+
+    /** The Iceberg field ids of the equality (primary-key) columns. */
+    Set<Integer> equalityFieldIds() {
+      return equalityFieldIds;
+    }
+
+    /** The Beam schema of the equality columns, in ascending Iceberg field-id 
order. */
+    Schema pkSchema() {
+      return pkSchema;
+    }
+
+    /** A deterministic coder for {@link #pkSchema()} rows. */
+    RowCoder pkCoder() {
+      return pkCoder;
+    }
+
+    /**
+     * Position of each {@link #pkSchema()} field within {@link 
#cdcDataSchema()}; do not mutate.
+     */
+    int[] pkFieldPositions() {
+      return pkFieldPositions;
+    }
+
+    /** The written-row schema: the table's schema in Beam form. */
+    Schema cdcDataSchema() {
+      return cdcDataSchema;
+    }
+
+    /** The partition-block sharding plan, or {@code null} to shard by plain 
primary-key hash. */
+    @Nullable PartitionShardPlan partitionShardPlan() {
+      return partitionShardPlan;
+    }
+  }
+
+  /**
+   * A table-level (as opposed to record-level) configuration failure: the 
destination table (or the
+   * sink configuration as applied to it) cannot accept CDC writes at all. 
Callers rethrow this
+   * fail-fast, bypassing any per-record poison-record handling.
+   */
+  static final class TableConfigException extends RuntimeException {
+    TableConfigException(String message) {
+      super(message);
+    }
+
+    TableConfigException(String message, Throwable cause) {
+      super(message, cause);
+    }
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitTokenTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitTokenTest.java
new file mode 100644
index 00000000000..78fb5e95f3a
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CommitTokenTest.java
@@ -0,0 +1,173 @@
+/*
+ * 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.cdc.sink;
+
+import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.nullValue;
+import static org.junit.Assert.assertThrows;
+
+import java.util.Map;
+import org.apache.beam.sdk.metrics.Metrics;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableSet;
+import org.apache.iceberg.AppendFiles;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Snapshot;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.catalog.Catalog;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.types.Types;
+import org.checkerframework.checker.nullness.qual.Nullable;
+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 {@link CommitToken}. */
+@RunWith(JUnit4.class)
+public class CommitTokenTest {
+
+  @Rule public transient TemporaryFolder tmp = new TemporaryFolder();
+
+  private static final String SINK_ID = "sink-a";
+  private static final String RUN_ID = "run-1";
+  private static final String DEST = "db.t";
+  private static final Schema SCHEMA =
+      new Schema(Types.NestedField.required(1, "id", Types.IntegerType.get()));
+
+  private Table table;
+  private CommitToken token;
+
+  @Before
+  public void setUp() {
+    Catalog catalog = CdcSinkTestUtils.hadoopCatalog(tmp.getRoot());
+    table =
+        CdcSinkTestUtils.createTable(
+            catalog,
+            TableIdentifier.of("db", "t" + System.nanoTime()),
+            SCHEMA,
+            ImmutableSet.of(1),
+            2,
+            PartitionSpec.unpartitioned());
+    token =
+        new CommitToken(
+            SINK_ID,
+            RUN_ID,
+            Metrics.counter(CommitTokenTest.class, "tokenParseFailures"),
+            Metrics.counter(CommitTokenTest.class, "suspectedTokenExpiry"));
+  }
+
+  /** Commits an empty snapshot carrying a full token. */
+  private void commitToken(long committedThroughMs, long maxSeq, @Nullable 
Integer specId) {
+    AppendFiles op = table.newAppend();
+    token.writeTo(op, committedThroughMs, maxSeq, specId);
+    op.commit();
+  }
+
+  private Map<String, String> currentSummary() {
+    return checkStateNotNull(table.currentSnapshot()).summary();
+  }
+
+  @Test
+  public void emptyTableIsAFreshStart() {
+    CommitToken.Recovered recovered = token.recoverFrom(table, DEST);
+    assertThat(recovered.committedThroughMs, equalTo(Long.MIN_VALUE));
+    assertThat(recovered.maxCommittedSeq, equalTo(Long.MIN_VALUE));
+    assertThat(CommitToken.readRunSpec(table, SINK_ID, RUN_ID), nullValue());
+  }
+
+  @Test
+  public void newestTokenIsRecoveredWithItsOwnMaxSeq() {
+    commitToken(1_000L, 5L, 3);
+    commitToken(2_000L, 9L, 3);
+
+    CommitToken.Recovered recovered = token.recoverFrom(table, DEST);
+    assertThat(recovered.committedThroughMs, equalTo(2_000L));
+    assertThat(recovered.maxCommittedSeq, equalTo(9L));
+    assertThat(CommitToken.readRunSpec(table, SINK_ID, RUN_ID), equalTo(3));
+    assertThat(currentSummary().get(CommitToken.SINK_ID_KEY), 
equalTo(SINK_ID));
+  }
+
+  @Test
+  public void unparseableTokenIsSkippedWithoutLendingItsMaxSeq() {
+    commitToken(1_000L, 5L, null);
+    AppendFiles corrupt = table.newAppend();
+    corrupt.set(CommitToken.COMMITTED_THROUGH_MS_PREFIX + SINK_ID, "nope");
+    corrupt.set(CommitToken.MAX_COMMITTED_SEQ_PREFIX + SINK_ID, "99");
+    corrupt.commit();
+
+    CommitToken.Recovered recovered = token.recoverFrom(table, DEST);
+    assertThat(recovered.committedThroughMs, equalTo(1_000L));
+    assertThat(recovered.maxCommittedSeq, equalTo(5L));
+  }
+
+  @Test
+  public void otherSinksAndRunsAreIgnored() {
+    table
+        .newAppend()
+        .set(CommitToken.SINK_ID_KEY, "sink-b")
+        .set(CommitToken.COMMITTED_THROUGH_MS_PREFIX + "sink-b", "7")
+        .set(CommitToken.RUN_SPEC_PREFIX + SINK_ID, "run-2:4")
+        .commit();
+
+    assertThat(token.recoverFrom(table, DEST).committedThroughMs, 
equalTo(Long.MIN_VALUE));
+    assertThat(CommitToken.readRunSpec(table, SINK_ID, RUN_ID), nullValue());
+    assertThat(CommitToken.readRunSpec(table, SINK_ID, "run-2"), equalTo(4));
+  }
+
+  @Test
+  public void heartbeatOmitsAnUnknownMaxSeq() {
+    AppendFiles op = table.newAppend();
+    token.writeHeartbeatTo(op, 3_000L, Long.MIN_VALUE, 2);
+    op.commit();
+
+    Map<String, String> summary = currentSummary();
+    assertThat(summary.get(CommitToken.COMMITTED_THROUGH_MS_PREFIX + SINK_ID), 
equalTo("3000"));
+    assertThat(summary.containsKey(CommitToken.MAX_COMMITTED_SEQ_PREFIX + 
SINK_ID), equalTo(false));
+    assertThat(summary.get(CommitToken.RUN_SPEC_PREFIX + SINK_ID), 
equalTo("run-1:2"));
+    assertThat(token.recoverFrom(table, DEST).maxCommittedSeq, 
equalTo(Long.MIN_VALUE));
+  }
+
+  @Test
+  public void heartbeatIsDueOnceTheTokenSnapshotIsOlderThanTheInterval() {
+    assertThat(token.shouldHeartbeat(table, 1_000L, 0L), equalTo(false));
+    commitToken(1_000L, 1L, null);
+    long committedAt = 
checkStateNotNull(table.currentSnapshot()).timestampMillis();
+
+    assertThat(token.shouldHeartbeat(table, 1_000L, committedAt + 500L), 
equalTo(false));
+    assertThat(token.shouldHeartbeat(table, 1_000L, committedAt + 1_001L), 
equalTo(true));
+  }
+
+  @Test
+  public void findsTheSnapshotCarryingAWindowToken() {
+    commitToken(1_000L, 1L, null);
+    Snapshot first = checkStateNotNull(table.currentSnapshot());
+    commitToken(2_000L, 2L, null);
+
+    assertThat(
+        token.findRecentlyCommittedTokenSnapshot(table, DEST, 
1_000L).snapshotId(),
+        equalTo(first.snapshotId()));
+    assertThrows(
+        IllegalStateException.class,
+        () -> token.findRecentlyCommittedTokenSnapshot(table, DEST, 3_000L));
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlanTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlanTest.java
new file mode 100644
index 00000000000..7839cf882f5
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/PartitionShardPlanTest.java
@@ -0,0 +1,120 @@
+/*
+ * 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.cdc.sink;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.greaterThanOrEqualTo;
+import static org.hamcrest.Matchers.lessThan;
+
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.values.Row;
+import org.apache.iceberg.PartitionKey;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.StructLike;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.InternalRecordWrapper;
+import org.apache.iceberg.types.JavaHash;
+import org.apache.iceberg.types.Types;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/** Tests for {@link PartitionShardPlan}. */
+@RunWith(JUnit4.class)
+public class PartitionShardPlanTest {
+
+  private static final int NUM_SHARDS = 8;
+
+  private static final org.apache.iceberg.Schema TABLE_SCHEMA =
+      new org.apache.iceberg.Schema(
+          Types.NestedField.required(1, "id", Types.IntegerType.get()),
+          Types.NestedField.required(2, "region", Types.StringType.get()),
+          Types.NestedField.optional(3, "name", Types.StringType.get()));
+
+  private static final Schema DATA_SCHEMA =
+      Schema.builder()
+          .addInt32Field("id")
+          .addStringField("region")
+          .addNullableField("name", Schema.FieldType.STRING)
+          .build();
+
+  private static final PartitionSpec BY_REGION =
+      PartitionSpec.builderFor(TABLE_SCHEMA).identity("region").build();
+
+  private static Row row(int id, String region, String name) {
+    return Row.withSchema(DATA_SCHEMA).addValues(id, region, name).build();
+  }
+
+  @Test
+  public void offsetSelectsAShardWithinThePartitionsBlock() {
+    PartitionShardPlan plan = PartitionShardPlan.of(BY_REGION, TABLE_SCHEMA, 
DATA_SCHEMA);
+    Row row = row(1, "us", "a");
+    int base = plan.shardFor(row, 0, NUM_SHARDS);
+
+    for (int offset = 0; offset < NUM_SHARDS; offset++) {
+      int shard = plan.shardFor(row, offset, NUM_SHARDS);
+      assertThat(shard, equalTo(Math.floorMod(base + offset, NUM_SHARDS)));
+      assertThat(shard, greaterThanOrEqualTo(0));
+      assertThat(shard, lessThan(NUM_SHARDS));
+    }
+  }
+
+  @Test
+  public void samePartitionTupleSharesABaseRegardlessOfOtherColumns() {
+    PartitionShardPlan plan = PartitionShardPlan.of(BY_REGION, TABLE_SCHEMA, 
DATA_SCHEMA);
+
+    assertThat(
+        plan.shardFor(row(1, "us", "a"), 0, NUM_SHARDS),
+        equalTo(plan.shardFor(row(2, "us", null), 0, NUM_SHARDS)));
+    assertThat(
+        plan.shardFor(row(1, "us", "a"), 3, NUM_SHARDS),
+        equalTo(plan.shardFor(row(1, "us", "a"), 3, NUM_SHARDS)));
+  }
+
+  // The block base is the sink's shard of Iceberg's own hash of the partition 
tuple, so a plan
+  // built from the Beam row agrees with a PartitionKey filled from an Iceberg 
record.
+  @Test
+  public void hashesTheIcebergPartitionTuple() {
+    PartitionShardPlan plan = PartitionShardPlan.of(BY_REGION, TABLE_SCHEMA, 
DATA_SCHEMA);
+
+    GenericRecord record = GenericRecord.create(TABLE_SCHEMA);
+    record.setField("id", 1);
+    record.setField("region", "eu");
+    record.setField("name", "b");
+    PartitionKey key = new PartitionKey(BY_REGION, TABLE_SCHEMA);
+    key.partition(new 
InternalRecordWrapper(TABLE_SCHEMA.asStruct()).wrap(record));
+    JavaHash<StructLike> hash = JavaHash.forType(BY_REGION.partitionType());
+    int expected = TableSetup.shardForHash(hash.hash(key), NUM_SHARDS);
+
+    assertThat(plan.shardFor(row(1, "eu", "b"), 0, NUM_SHARDS), 
equalTo(expected));
+  }
+
+  // One source column feeding two partition fields is converted once.
+  @Test
+  public void oneSourceColumnCanFeedSeveralPartitionFields() {
+    PartitionSpec spec =
+        
PartitionSpec.builderFor(TABLE_SCHEMA).identity("region").truncate("region", 
1).build();
+    PartitionShardPlan plan = PartitionShardPlan.of(spec, TABLE_SCHEMA, 
DATA_SCHEMA);
+
+    int shard = plan.shardFor(row(1, "us", "a"), 0, NUM_SHARDS);
+    assertThat(shard, equalTo(plan.shardFor(row(9, "us", "z"), 0, 
NUM_SHARDS)));
+    assertThat(shard, greaterThanOrEqualTo(0));
+    assertThat(shard, lessThan(NUM_SHARDS));
+  }
+}
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetupTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetupTest.java
new file mode 100644
index 00000000000..63133acc32b
--- /dev/null
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/cdc/sink/TableSetupTest.java
@@ -0,0 +1,1028 @@
+/*
+ * 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.cdc.sink;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.contains;
+import static org.hamcrest.Matchers.containsInAnyOrder;
+import static org.hamcrest.Matchers.containsString;
+import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.greaterThan;
+import static org.hamcrest.Matchers.greaterThanOrEqualTo;
+import static org.hamcrest.Matchers.hasSize;
+import static org.hamcrest.Matchers.lessThan;
+import static org.hamcrest.Matchers.not;
+import static org.hamcrest.Matchers.notNullValue;
+import static org.hamcrest.Matchers.nullValue;
+import static org.hamcrest.Matchers.sameInstance;
+import static org.junit.Assert.assertArrayEquals;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+import java.nio.charset.StandardCharsets;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.beam.sdk.io.iceberg.DynamicDestinations;
+import org.apache.beam.sdk.io.iceberg.IcebergCatalogConfig;
+import org.apache.beam.sdk.io.iceberg.IcebergDestination;
+import org.apache.beam.sdk.io.iceberg.IcebergTableCreateConfig;
+import org.apache.beam.sdk.io.iceberg.IcebergUtils;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.util.SerializableUtils;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.sdk.values.ValueInSingleWindow;
+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.FileFormat;
+import org.apache.iceberg.PartitionField;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableUtil;
+import org.apache.iceberg.catalog.Catalog;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.expressions.Expressions;
+import org.apache.iceberg.types.Type;
+import org.apache.iceberg.types.Types;
+import org.checkerframework.checker.nullness.qual.Nullable;
+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 {@link TableSetup}. */
+@RunWith(JUnit4.class)
+public class TableSetupTest {
+
+  @Rule public TemporaryFolder tmp = new TemporaryFolder();
+
+  /** Canonical test table schema: {@code id INT (required)}, {@code 
name}/{@code 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()));
+
+  /** Data schema for {@link #ICEBERG_SCHEMA}. */
+  private static final Schema DATA_SCHEMA =
+      Schema.builder()
+          .addInt32Field("id")
+          .addNullableField("name", Schema.FieldType.STRING)
+          .addNullableField("data", Schema.FieldType.STRING)
+          .build();
+
+  private Catalog catalog;
+  private IcebergCatalogConfig catalogConfig;
+
+  @Before
+  public void setUp() {
+    catalog = CdcSinkTestUtils.hadoopCatalog(tmp.getRoot());
+    catalogConfig = CdcSinkTestUtils.catalogConfig(tmp.getRoot());
+  }
+
+  private static TableIdentifier uniqueId(String prefix) {
+    return TableIdentifier.of("db", prefix + "_" + System.nanoTime());
+  }
+
+  /** A fresh unpartitioned V2 {@link #ICEBERG_SCHEMA} table (PK {@code id}), 
named from prefix. */
+  private TableIdentifier v2Table(String prefix) {
+    TableIdentifier id = uniqueId(prefix);
+    CdcSinkTestUtils.createTable(
+        catalog, id, ICEBERG_SCHEMA, ImmutableSet.of(1), 2, 
PartitionSpec.unpartitioned());
+    return id;
+  }
+
+  /** {@link #v2Table} partitioned by {@code bucket(column, buckets)}. */
+  private TableIdentifier bucketPartitionedTable(String prefix, String column, 
int buckets) {
+    TableIdentifier id = uniqueId(prefix);
+    PartitionSpec spec = 
PartitionSpec.builderFor(ICEBERG_SCHEMA).bucket(column, buckets).build();
+    CdcSinkTestUtils.createTable(catalog, id, ICEBERG_SCHEMA, 
ImmutableSet.of(1), 2, spec);
+    return id;
+  }
+
+  private static CdcWriteConfig.Builder cfg() {
+    return 
CdcWriteConfig.builder().setSinkId("test-sink").setNumShards(8).setShardsPerPartition(8);
+  }
+
+  private TableSetup tableSetup(CdcWriteConfig config, DynamicDestinations 
destinations) {
+    return new TableSetup(catalogConfig, config, destinations, "test-runId");
+  }
+
+  private TableSetup tableSetup(CdcWriteConfig config) {
+    return tableSetup(config, new TestDestinations(DATA_SCHEMA, null, null));
+  }
+
+  private static Schema dataSchemaFor(org.apache.iceberg.Schema icebergSchema) 
{
+    return IcebergUtils.icebergSchemaToBeamSchema(icebergSchema);
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // Loading and Dest population
+  // 
-------------------------------------------------------------------------------------------
+
+  @Test
+  public void loadsExistingTableAndPopulatesDest() {
+    org.apache.iceberg.Schema schema =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.IntegerType.get()),
+            Types.NestedField.optional(2, "name", Types.StringType.get()),
+            Types.NestedField.required(3, "num", Types.LongType.get()));
+    TableIdentifier id = uniqueId("existing");
+    // Identifier fields deliberately declared out of ascending-field-id 
order: {3, 1}.
+    CdcSinkTestUtils.createTable(
+        catalog, id, schema, ImmutableSet.of(3, 1), 2, 
PartitionSpec.unpartitioned());
+
+    Schema sourceSchema = dataSchemaFor(schema);
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(sourceSchema, null, null));
+
+    TableSetup.Dest dest = setup.get(id.toString(), sourceSchema);
+
+    assertThat(dest.table().name(), containsString(id.name()));
+    assertThat(dest.equalityFieldIds(), containsInAnyOrder(1, 3));
+    // pkSchema is in ascending field-id order even though the identifiers 
were declared {3, 1}.
+    assertThat(dest.pkSchema().getFieldNames(), contains("id", "num"));
+    assertThat(dest.cdcDataSchema(), 
equalTo(IcebergUtils.icebergSchemaToBeamSchema(schema)));
+    assertArrayEquals(new int[] {0, 2}, dest.pkFieldPositions());
+    assertThat(dest.pkCoder(), notNullValue());
+  }
+
+  @Test
+  public void memoizesDestPerDestinationString() {
+    TableIdentifier id = v2Table("memoized");
+    TableSetup setup = tableSetup(cfg().build());
+
+    TableSetup.Dest first = setup.get(id.toString(), DATA_SCHEMA);
+    TableSetup.Dest second = setup.get(id.toString(), DATA_SCHEMA);
+
+    assertThat(second, sameInstance(first));
+  }
+
+  /**
+   * The memo is keyed by destination string: two destinations resolved 
through ONE {@link
+   * TableSetup} get their own {@link TableSetup.Dest} each; a memo ignoring 
the destination would
+   * silently hand every later destination the first table's Dest, and 
single-destination tests
+   * cannot see it.
+   */
+  @Test
+  public void memoizesEachDestinationSeparately() {
+    // Genuinely different tables: different column names, types, and 
identifier-field counts, so a
+    // cross-wired Dest cannot masquerade as the right one.
+    org.apache.iceberg.Schema schemaA =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.IntegerType.get()),
+            Types.NestedField.optional(2, "name", Types.StringType.get()));
+    org.apache.iceberg.Schema schemaB =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "sku", Types.StringType.get()),
+            Types.NestedField.required(2, "region", Types.StringType.get()),
+            Types.NestedField.optional(3, "qty", Types.LongType.get()));
+    TableIdentifier idA = uniqueId("multi_a");
+    TableIdentifier idB = uniqueId("multi_b");
+    CdcSinkTestUtils.createTable(
+        catalog, idA, schemaA, ImmutableSet.of(1), 2, 
PartitionSpec.unpartitioned());
+    CdcSinkTestUtils.createTable(
+        catalog, idB, schemaB, ImmutableSet.of(1, 2), 2, 
PartitionSpec.unpartitioned());
+
+    Schema sourceA = dataSchemaFor(schemaA);
+    Schema sourceB = dataSchemaFor(schemaB);
+    TableSetup setup = tableSetup(cfg().build(), new TestDestinations(sourceA, 
null, null));
+
+    TableSetup.Dest destA = setup.get(idA.toString(), sourceA);
+    TableSetup.Dest destB = setup.get(idB.toString(), sourceB);
+
+    assertThat(destB, not(sameInstance(destA)));
+    assertThat(destA.table().name(), containsString(idA.name()));
+    assertThat(destB.table().name(), containsString(idB.name()));
+    assertThat(destA.equalityFieldIds(), contains(1));
+    assertThat(destB.equalityFieldIds(), containsInAnyOrder(1, 2));
+    assertThat(destA.pkSchema().getFieldNames(), contains("id"));
+    assertThat(destB.pkSchema().getFieldNames(), contains("sku", "region"));
+    assertThat(destA.cdcDataSchema().getFieldNames(), contains("id", "name"));
+    assertThat(destB.cdcDataSchema().getFieldNames(), contains("sku", 
"region", "qty"));
+    assertArrayEquals(new int[] {0}, destA.pkFieldPositions());
+    assertArrayEquals(new int[] {0, 1}, destB.pkFieldPositions());
+
+    // Both entries live in the memo at once: re-getting either returns ITS 
OWN instance.
+    assertThat(setup.get(idA.toString(), sourceA), sameInstance(destA));
+    assertThat(setup.get(idB.toString(), sourceB), sameInstance(destB));
+  }
+
+  @Test
+  public void serializesAndResolvesAfterDeserialization() {
+    TableIdentifier id = v2Table("serializable");
+
+    TableSetup roundTripped = 
SerializableUtils.clone(tableSetup(cfg().build()));
+
+    TableSetup.Dest dest = roundTripped.get(id.toString(), DATA_SCHEMA);
+    assertThat(dest.pkSchema().getFieldNames(), contains("id"));
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // Auto-creation
+  // 
-------------------------------------------------------------------------------------------
+
+  @Test
+  public void autoCreatesMissingTable() {
+    TableIdentifier id = uniqueId("autocreate");
+    CdcWriteConfig config = 
cfg().setEqualityColumns(ImmutableList.of("id")).build();
+    TestDestinations destinations =
+        new TestDestinations(DATA_SCHEMA, ImmutableList.of("id"), 
ImmutableList.of("name"));
+
+    TableSetup.Dest dest = tableSetup(config, destinations).get(id.toString(), 
DATA_SCHEMA);
+
+    Table table = catalog.loadTable(id);
+    assertThat(TableUtil.formatVersion(table), equalTo(2));
+    // Identifier fields are the configured equality columns.
+    int idFieldId = table.schema().findField("id").fieldId();
+    assertThat(table.schema().identifierFieldIds(), contains(idFieldId));
+    // Partition spec and sort order from the destination's create config are 
honored.
+    assertThat(table.spec().fields(), hasSize(1));
+    PartitionField partitionField = table.spec().fields().get(0);
+    assertThat(partitionField.sourceId(), equalTo(idFieldId));
+    assertTrue(partitionField.transform().isIdentity());
+    assertThat(table.sortOrder().fields(), hasSize(1));
+    assertThat(
+        table.sortOrder().fields().get(0).sourceId(),
+        equalTo(table.schema().findField("name").fieldId()));
+
+    assertThat(dest.equalityFieldIds(), contains(idFieldId));
+    assertThat(dest.cdcDataSchema().getFieldNames(), contains("id", "name", 
"data"));
+  }
+
+  /**
+   * Pins the created columns' TYPES and required/optional flags: a wrong type 
or a silently
+   * nullable column is invisible to the name-only assertions everywhere else.
+   */
+  @Test
+  public void autoCreatedColumnsCarryInputTypesAndNullability() {
+    Schema inputSchema =
+        Schema.builder()
+            .addInt32Field("id")
+            .addStringField("code")
+            .addNullableField("name", Schema.FieldType.STRING)
+            .addNullableField("amount", Schema.FieldType.DOUBLE)
+            .addInt64Field("version")
+            .addBooleanField("active")
+            .build();
+    TableIdentifier id = uniqueId("autocreate_types");
+    CdcWriteConfig config = cfg().setEqualityColumns(ImmutableList.of("id", 
"code")).build();
+
+    TableSetup.Dest dest =
+        tableSetup(config, new TestDestinations(inputSchema, null, null))
+            .get(id.toString(), inputSchema);
+
+    org.apache.iceberg.Schema created = catalog.loadTable(id).schema();
+    assertFieldIs(created, "id", Types.IntegerType.get(), /* required= */ 
true);
+    assertFieldIs(created, "code", Types.StringType.get(), /* required= */ 
true);
+    assertFieldIs(created, "name", Types.StringType.get(), /* required= */ 
false);
+    assertFieldIs(created, "amount", Types.DoubleType.get(), /* required= */ 
false);
+    assertFieldIs(created, "version", Types.LongType.get(), /* required= */ 
true);
+    assertFieldIs(created, "active", Types.BooleanType.get(), /* required= */ 
true);
+    // Both equality columns became identifier fields.
+    assertThat(
+        created.identifierFieldIds(),
+        containsInAnyOrder(created.findField("id").fieldId(), 
created.findField("code").fieldId()));
+    // The sink's own resolved view agrees with the table it just created.
+    assertThat(dest.cdcDataSchema(), 
equalTo(IcebergUtils.icebergSchemaToBeamSchema(created)));
+    assertThat(dest.pkSchema().getFieldNames(), contains("id", "code"));
+  }
+
+  /** Asserts one created Iceberg column's type and required/optional flag. */
+  private static void assertFieldIs(
+      org.apache.iceberg.Schema schema, String name, Type type, boolean 
required) {
+    Types.NestedField field = schema.findField(name);
+    assertThat("column '" + name + "' is missing", field, notNullValue());
+    assertThat("column '" + name + "' type", field.type(), equalTo(type));
+    assertThat("column '" + name + "' requiredness", field.isRequired(), 
equalTo(required));
+  }
+
+  /**
+   * Create-config table properties reach the created table, and the sink's 
{@code format-version=2}
+   * default applies only when the user did not ask for one (a plain {@code 
put} would silently
+   * downgrade a requested V3 table).
+   */
+  @Test
+  public void autoCreateHonorsTablePropertiesAndDefaultsToFormatVersion2() {
+    CdcWriteConfig config = 
cfg().setEqualityColumns(ImmutableList.of("id")).build();
+
+    // facet: explicit format-version 3 and a custom property both honored.
+    TableIdentifier propsId = uniqueId("autocreate_props");
+    TestDestinations destinations =
+        new TestDestinations(
+            DATA_SCHEMA,
+            null,
+            null,
+            ImmutableMap.of("format-version", "3", "cdc.test.owner", 
"cdc-team"));
+    tableSetup(config, destinations).get(propsId.toString(), DATA_SCHEMA);
+    Table table = catalog.loadTable(propsId);
+    assertThat(TableUtil.formatVersion(table), equalTo(3));
+    assertThat(table.properties().get("cdc.test.owner"), equalTo("cdc-team"));
+
+    // facet: no create-config properties at all still defaults to V2.
+    TableIdentifier defaultId = uniqueId("autocreate_default_fv");
+    tableSetup(config, new TestDestinations(DATA_SCHEMA, null, null))
+        .get(defaultId.toString(), DATA_SCHEMA);
+    assertThat(TableUtil.formatVersion(catalog.loadTable(defaultId)), 
equalTo(2));
+  }
+
+  @Test
+  public void rejectsAutoCreateWithNullableEqualityColumn() {
+    TableIdentifier id = uniqueId("autocreate_nullable");
+    CdcWriteConfig config = 
cfg().setEqualityColumns(ImmutableList.of("name")).build();
+    TableSetup setup = tableSetup(config);
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString("'name'"));
+    assertThat(error.getMessage(), containsString("non-nullable"));
+  }
+
+  @Test
+  public void rejectsAutoCreateWithoutEqualityColumns() {
+    TableIdentifier id = uniqueId("autocreate_no_eq");
+    TableSetup setup = tableSetup(cfg().build());
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString("does not exist"));
+    assertThat(error.getMessage(), containsString("equality_columns"));
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // Validation rejections
+  // 
-------------------------------------------------------------------------------------------
+
+  @Test
+  public void rejectsFormatVersion1Table() {
+    TableIdentifier id = uniqueId("v1");
+    CdcSinkTestUtils.createTable(
+        catalog, id, ICEBERG_SCHEMA, ImmutableSet.of(1), 1, 
PartitionSpec.unpartitioned());
+    TableSetup setup = tableSetup(cfg().build());
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString(id.toString()));
+    assertThat(error.getMessage(), containsString("append sink"));
+  }
+
+  @Test
+  public void rejectsNullableEqualityColumn() {
+    TableIdentifier id = v2Table("nullable_pk");
+    // 'name' exists in the table but is optional, so it cannot define row 
identity.
+    CdcWriteConfig config = 
cfg().setEqualityColumns(ImmutableList.of("name")).build();
+    TableSetup setup = tableSetup(config);
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString("'name'"));
+    assertThat(error.getMessage(), containsString("must be required"));
+  }
+
+  /** A fresh {@code day(ts)}-partitioned table whose partition source is NOT 
an equality column. */
+  private TableIdentifier nonKeyDayPartitionedTable(String prefix) {
+    org.apache.iceberg.Schema schema = timestampSchema();
+    TableIdentifier id = uniqueId(prefix);
+    CdcSinkTestUtils.createTable(
+        catalog,
+        id,
+        schema,
+        ImmutableSet.of(1),
+        2,
+        PartitionSpec.builderFor(schema).day("ts").build());
+    return id;
+  }
+
+  /**
+   * {@code day(ts)} with PK {@code [id]} resolves at default config: the 
writer routes each
+   * equality delete by its block's opening record, so no 
key-derived-partition requirement applies.
+   */
+  @Test
+  public void acceptsNonKeyPartitionSourceAtDefaultConfig() {
+    TableIdentifier id = nonKeyDayPartitionedTable("nonkey_default");
+    Schema sourceSchema = dataSchemaFor(timestampSchema());
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(sourceSchema, null, null));
+
+    TableSetup.Dest dest = setup.get(id.toString(), sourceSchema);
+
+    assertThat(dest.equalityFieldIds(), contains(1));
+    assertThat(dest.partitionShardPlan(), nullValue());
+  }
+
+  /**
+   * The two options that need the partition to be a pure function of the 
primary key still reject a
+   * non-key partition source, each naming itself and the offending column.
+   */
+  @Test
+  public void upsertAndShardCapStillRequireKeyDerivedPartitions() {
+    Schema sourceSchema = dataSchemaFor(timestampSchema());
+
+    // facet: upsert (before-images are dropped, so a moved row's old 
partition is unreachable).
+    TableIdentifier upsertId = nonKeyDayPartitionedTable("nonkey_upsert");
+    TableSetup upsertSetup =
+        tableSetup(cfg().setUpsert(true).build(), new 
TestDestinations(sourceSchema, null, null));
+    TableSetup.TableConfigException upsertError =
+        assertThrows(
+            TableSetup.TableConfigException.class,
+            () -> upsertSetup.get(upsertId.toString(), sourceSchema));
+    assertThat(upsertError.getMessage(), containsString("upsert"));
+    assertThat(upsertError.getMessage(), containsString("'ts'"));
+
+    // facet: shards_per_partition below num_shards (the shard is derived from 
the partition
+    // tuple, which must therefore follow from the primary key).
+    TableIdentifier cappedId = nonKeyDayPartitionedTable("nonkey_capped");
+    TableSetup cappedSetup =
+        tableSetup(
+            cfg().setShardsPerPartition(2).build(), new 
TestDestinations(sourceSchema, null, null));
+    TableSetup.TableConfigException cappedError =
+        assertThrows(
+            TableSetup.TableConfigException.class,
+            () -> cappedSetup.get(cappedId.toString(), sourceSchema));
+    assertThat(cappedError.getMessage(), 
containsString("shards_per_partition"));
+    assertThat(cappedError.getMessage(), containsString("'ts'"));
+  }
+
+  /**
+   * A memoized {@link TableSetup.Dest} is handed back unchanged after a live 
spec evolution: the
+   * write path pins {@code specId()}, so no drift check runs on the default 
path. Only block
+   * sharding re-checks ({@link #blockShardingStillRefusesSpecDrift}).
+   */
+  @Test
+  public void memoHitUnderEvolvedSpecReturnsPinnedDest() {
+    TableIdentifier id = v2Table("spec_evolution");
+    TableSetup setup = tableSetup(cfg().build());
+
+    TableSetup.Dest dest = setup.get(id.toString(), DATA_SCHEMA);
+    int resolvedSpecId = dest.spec().specId();
+    assertThat(resolvedSpecId, equalTo(dest.table().spec().specId()));
+
+    // An operator evolves the spec mid-run; the sink's shared Table instance 
picks it up.
+    dest.table().updateSpec().addField(Expressions.bucket("id", 4)).commit();
+    dest.table().refresh();
+    assertThat(dest.table().spec().specId(), not(equalTo(resolvedSpecId)));
+
+    TableSetup.Dest again = setup.get(id.toString(), DATA_SCHEMA);
+
+    assertThat(again, sameInstance(dest));
+    assertThat(again.spec().specId(), equalTo(resolvedSpecId));
+  }
+
+  /**
+   * With a {@link PartitionShardPlan} present, memo-hit spec drift must still 
throw naming both
+   * spec ids: workers on different specs would split one key across shards, 
silent duplicates.
+   */
+  @Test
+  public void blockShardingStillRefusesSpecDrift() {
+    TableIdentifier id = bucketPartitionedTable("block_spec_drift", "id", 4);
+    TableSetup setup = tableSetup(cfg().setShardsPerPartition(2).build());
+
+    TableSetup.Dest dest = setup.get(id.toString(), DATA_SCHEMA);
+    assertThat(dest.partitionShardPlan(), notNullValue());
+    int resolvedSpecId = dest.spec().specId();
+
+    dest.table().updateSpec().addField(Expressions.bucket("id", 8)).commit();
+    dest.table().refresh();
+    int evolvedSpecId = dest.table().spec().specId();
+    assertThat(evolvedSpecId, not(equalTo(resolvedSpecId)));
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString(id.toString()));
+    assertThat(error.getMessage(), containsString("spec id " + 
resolvedSpecId));
+    assertThat(error.getMessage(), containsString("spec id " + evolvedSpecId));
+    assertThat(error.getMessage(), containsString("Drain the pipeline"));
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // Run-spec stamp adoption at resolution
+  // 
-------------------------------------------------------------------------------------------
+
+  /**
+   * Commits an empty snapshot carrying the run-spec stamp as the committer 
writes it. Literal
+   * strings on purpose: a contract pin, like the token keys.
+   */
+  private static void stampRunSpec(Table table, String runId, int specId) {
+    table.newAppend().set("beam.cdc.run-spec.test-sink", runId + ":" + 
specId).commit();
+    table.refresh();
+  }
+
+  /** Creates a {@code bucket(id, 4)}-partitioned table and resolves it once 
(warming the cache). */
+  private TableSetup.Dest resolvedBucketDest(TableIdentifier id) {
+    PartitionSpec spec = PartitionSpec.builderFor(ICEBERG_SCHEMA).bucket("id", 
4).build();
+    CdcSinkTestUtils.createTable(catalog, id, ICEBERG_SCHEMA, 
ImmutableSet.of(1), 2, spec);
+    return tableSetup(cfg().build()).get(id.toString(), DATA_SCHEMA);
+  }
+
+  /**
+   * A worker that first resolves a destination after a mid-run spec evolution 
adopts the spec the
+   * committer stamped for this run's runId, not the live current spec, and 
validates against it.
+   */
+  @Test
+  public void joiningWorkerAdoptsStampedSpec() {
+    TableIdentifier id = uniqueId("joining_worker");
+    TableSetup.Dest dest = resolvedBucketDest(id);
+    int stampedSpecId = dest.spec().specId();
+    stampRunSpec(dest.table(), "runId-n", stampedSpecId);
+
+    dest.table().updateSpec().addField(Expressions.bucket("id", 8)).commit();
+    dest.table().refresh();
+    assertThat(dest.table().spec().specId(), not(equalTo(stampedSpecId)));
+
+    TableSetup joining =
+        new TableSetup(
+            catalogConfig, cfg().build(), new TestDestinations(DATA_SCHEMA, 
null, null), "runId-n");
+
+    assertThat(joining.get(id.toString(), DATA_SCHEMA).spec().specId(), 
equalTo(stampedSpecId));
+  }
+
+  /** A stamp from another run's runId is ignored: a fresh run resolves the 
current spec. */
+  @Test
+  public void freshRunAdoptsCurrentSpec() {
+    TableIdentifier id = uniqueId("fresh_run");
+    TableSetup.Dest dest = resolvedBucketDest(id);
+    stampRunSpec(dest.table(), "runId-n", dest.spec().specId());
+
+    dest.table().updateSpec().addField(Expressions.bucket("id", 8)).commit();
+    dest.table().refresh();
+    int currentSpecId = dest.table().spec().specId();
+
+    TableSetup fresh =
+        new TableSetup(
+            catalogConfig, cfg().build(), new TestDestinations(DATA_SCHEMA, 
null, null), "runId-m");
+
+    assertThat(fresh.get(id.toString(), DATA_SCHEMA).spec().specId(), 
equalTo(currentSpecId));
+  }
+
+  /** A stamp naming a spec id the table does not have falls back to the 
current spec, no throw. */
+  @Test
+  public void stampedSpecMissingFallsBackToCurrent() {
+    TableIdentifier id = uniqueId("stamp_missing");
+    TableSetup.Dest dest = resolvedBucketDest(id);
+    int currentSpecId = dest.spec().specId();
+    stampRunSpec(dest.table(), "runId-n", 99);
+
+    TableSetup joining =
+        new TableSetup(
+            catalogConfig, cfg().build(), new TestDestinations(DATA_SCHEMA, 
null, null), "runId-n");
+
+    assertThat(joining.get(id.toString(), DATA_SCHEMA).spec().specId(), 
equalTo(currentSpecId));
+  }
+
+  @Test
+  public void rejectsEqualityOverrideColumnMissingFromTable() {
+    TableIdentifier id = v2Table("missing_override");
+    CdcWriteConfig config = 
cfg().setEqualityColumns(ImmutableList.of("nonexistent")).build();
+    TableSetup setup = tableSetup(config);
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString("'nonexistent'"));
+    assertThat(error.getMessage(), containsString("does not exist"));
+  }
+
+  /**
+   * An EMPTY {@code equality_columns} override must be rejected at resolution 
too (not only by
+   * {@code CdcWriteConfig#validate}): an empty pk schema encodes every row to 
the SAME key: one
+   * shard takes the table and every equality delete matches every row, with 
nothing failing.
+   */
+  @Test
+  public void rejectsEmptyEqualityColumnsOverrideAtResolution() {
+    TableIdentifier id = v2Table("empty_eq_override");
+    CdcWriteConfig config = 
cfg().setEqualityColumns(ImmutableList.<String>of()).build();
+    TableSetup setup = tableSetup(config);
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString("equality_columns must be 
non-empty"));
+    assertThat(error.getMessage(), containsString(id.toString()));
+  }
+
+  @Test
+  public void rejectsMissingEqualityColumnsEverywhere() {
+    TableIdentifier id = uniqueId("no_identifiers");
+    // A V2 table with no identifier fields, and no equality_columns override 
configured.
+    CdcSinkTestUtils.createTable(
+        catalog, id, ICEBERG_SCHEMA, ImmutableSet.of(), 2, 
PartitionSpec.unpartitioned());
+    TableSetup setup = tableSetup(cfg().build());
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString("identifier"));
+    assertThat(error.getMessage(), containsString("equality_columns"));
+  }
+
+  @Test
+  public void rejectsCdcDataSchemaMismatch() {
+    TableIdentifier id = v2Table("mismatch");
+    Schema withExtra =
+        
Schema.builder().addFields(DATA_SCHEMA.getFields()).addStringField("extra").build();
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(withExtra, null, null));
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), withExtra));
+
+    assertThat(error.getMessage(), containsString("unexpected"));
+    assertThat(error.getMessage(), containsString("extra"));
+  }
+
+  @Test
+  public void rejectsDataColumnsInDifferentOrderThanTable() {
+    TableIdentifier id = v2Table("reordered");
+    // Same column names as the table (id, name, data) but in a different 
order: the written rows
+    // and the shuffle coder are built positionally, so order must match, not 
just the name set.
+    Schema reordered =
+        Schema.builder()
+            .addNullableField("name", Schema.FieldType.STRING)
+            .addInt32Field("id")
+            .addNullableField("data", Schema.FieldType.STRING)
+            .build();
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(reordered, null, null));
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), reordered));
+
+    assertThat(error.getMessage(), containsString("order"));
+    assertThat(error.getMessage(), containsString(id.toString()));
+  }
+
+  /** A mismatched column type is rejected naming the column, both types, and 
the remedy. */
+  @Test
+  public void rejectsColumnTypeMismatch() {
+    org.apache.iceberg.Schema schema =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.IntegerType.get()),
+            Types.NestedField.optional(2, "num", Types.LongType.get()));
+    TableIdentifier id = uniqueId("type_mismatch");
+    CdcSinkTestUtils.createTable(
+        catalog, id, schema, ImmutableSet.of(1), 2, 
PartitionSpec.unpartitioned());
+    // 'num' declared INT32 in the input where the table column is a long.
+    Schema mismatched =
+        Schema.builder()
+            .addInt32Field("id")
+            .addNullableField("num", Schema.FieldType.INT32)
+            .build();
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(mismatched, null, null));
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), mismatched));
+
+    assertThat(error.getMessage(), containsString("'num'"));
+    assertThat(error.getMessage(), containsString("INT32"));
+    assertThat(error.getMessage(), containsString("INT64"));
+    assertThat(error.getMessage(), containsString("Align the input schema with 
the table"));
+  }
+
+  /** A nullable-declared input column against a required table column is 
rejected. */
+  @Test
+  public void rejectsNullableInputColumnForRequiredTableColumn() {
+    org.apache.iceberg.Schema schema =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.IntegerType.get()),
+            Types.NestedField.required(2, "name", Types.StringType.get()));
+    TableIdentifier id = uniqueId("nullable_input");
+    CdcSinkTestUtils.createTable(
+        catalog, id, schema, ImmutableSet.of(1), 2, 
PartitionSpec.unpartitioned());
+    Schema nullableName =
+        Schema.builder()
+            .addInt32Field("id")
+            .addNullableField("name", Schema.FieldType.STRING)
+            .build();
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(nullableName, null, null));
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), nullableName));
+
+    assertThat(error.getMessage(), containsString("'name'"));
+    assertThat(error.getMessage(), containsString("nullable in the input"));
+    assertThat(error.getMessage(), containsString("required in the table"));
+  }
+
+  /** Exactly matching types resolve, including a non-null input column on an 
OPTIONAL one. */
+  @Test
+  public void acceptsMatchingTypesAndNonNullInputForOptionalTableColumn() {
+    TableIdentifier id = v2Table("types_ok");
+    // 'name' non-null in the input against the table's optional column: the 
safe direction.
+    Schema nonNullName =
+        Schema.builder()
+            .addInt32Field("id")
+            .addStringField("name")
+            .addNullableField("data", Schema.FieldType.STRING)
+            .build();
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(nonNullName, null, null));
+
+    TableSetup.Dest dest = setup.get(id.toString(), nonNullName);
+
+    assertThat(
+        dest.cdcDataSchema(), 
equalTo(IcebergUtils.icebergSchemaToBeamSchema(ICEBERG_SCHEMA)));
+  }
+
+  /**
+   * A column type with no Beam conversion fails destination resolution rather 
than silently writing
+   * null in every record. Pins the natural failure; the sink has no dedicated 
check.
+   */
+  @Test
+  public void unconvertibleColumnTypeFailsResolution() {
+    org.apache.iceberg.Schema withTimestampNano =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.IntegerType.get()),
+            Types.NestedField.optional(2, "ts_ns", 
Types.TimestampNanoType.withoutZone()));
+    TableIdentifier id = uniqueId("ts_nano");
+    // timestamp_ns is a format-version 3 type.
+    CdcSinkTestUtils.createTable(
+        catalog, id, withTimestampNano, ImmutableSet.of(1), 3, 
PartitionSpec.unpartitioned());
+    TableSetup setup = tableSetup(cfg().build());
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString(id.toString()));
+  }
+
+  @Test
+  public void rejectsNestedEqualityOverrideColumn() {
+    TableIdentifier id = v2Table("nested_eq");
+    CdcWriteConfig config = 
cfg().setEqualityColumns(ImmutableList.of("user.id")).build();
+    TableSetup setup = tableSetup(config);
+
+    TableSetup.TableConfigException error =
+        assertThrows(
+            TableSetup.TableConfigException.class, () -> 
setup.get(id.toString(), DATA_SCHEMA));
+
+    assertThat(error.getMessage(), containsString("top-level"));
+    assertThat(error.getMessage(), containsString("'user.id'"));
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // Partition transforms that must be ACCEPTED (all transforms are legal for 
the CDC sink)
+  // 
-------------------------------------------------------------------------------------------
+
+  private static org.apache.iceberg.Schema timestampSchema() {
+    return new org.apache.iceberg.Schema(
+        Types.NestedField.required(1, "id", Types.IntegerType.get()),
+        Types.NestedField.required(2, "ts", Types.TimestampType.withZone()));
+  }
+
+  private void assertPartitionAccepted(
+      org.apache.iceberg.Schema schema, PartitionSpec spec, String prefix) {
+    TableIdentifier id = uniqueId(prefix);
+    CdcSinkTestUtils.createTable(catalog, id, schema, ImmutableSet.of(1, 2), 
2, spec);
+    Schema sourceSchema = dataSchemaFor(schema);
+    TableSetup setup = tableSetup(cfg().build(), new 
TestDestinations(sourceSchema, null, null));
+
+    TableSetup.Dest dest = setup.get(id.toString(), sourceSchema);
+
+    assertThat(dest, notNullValue());
+    assertThat(dest.equalityFieldIds(), containsInAnyOrder(1, 2));
+  }
+
+  @Test
+  public void acceptsDayHourAndIdentityDatePartitions() {
+    org.apache.iceberg.Schema tsSchema = timestampSchema();
+    assertPartitionAccepted(
+        tsSchema, PartitionSpec.builderFor(tsSchema).day("ts").build(), 
"day_ts");
+    assertPartitionAccepted(
+        tsSchema, PartitionSpec.builderFor(tsSchema).hour("ts").build(), 
"hour_ts");
+    org.apache.iceberg.Schema dateSchema =
+        new org.apache.iceberg.Schema(
+            Types.NestedField.required(1, "id", Types.IntegerType.get()),
+            Types.NestedField.required(2, "d", Types.DateType.get()));
+    assertPartitionAccepted(
+        dateSchema, 
PartitionSpec.builderFor(dateSchema).identity("d").build(), "identity_date");
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // shardFor
+  // 
-------------------------------------------------------------------------------------------
+
+  /** Pins the exact hash function (murmur3_32_fixed + floorMod) against 
accidental change. */
+  @Test
+  public void shardForPinsMurmur3FixedFloorMod() {
+    assertThat(TableSetup.shardFor(new byte[] {0, 0, 0, 1}, 8), equalTo(4));
+    // Hashes to a negative int: floorMod maps it to 7, while abs(hash % n) 
would give 1. Freezing
+    // this prevents a silent resharding of half the keyspace.
+    assertThat(TableSetup.shardFor(new byte[] {0, 0, 0, 2}, 8), equalTo(7));
+  }
+
+  /**
+   * Deterministic, in range, spreading distinct keys, and covering EVERY 
shard of a
+   * non-power-of-two count, which pins {@code floorMod} against the {@code 
hash & (n - 1)}
+   * "optimization" (at 10 shards the mask can only produce 0, 1, 8, 9).
+   */
+  @Test
+  public void shardForIsDeterministicSpreadsKeysAndCoversNonPowerOfTwoCounts() 
{
+    // facet: determinism and spread at 8 shards.
+    int numShards = 8;
+    Set<Integer> shards = new HashSet<>();
+    for (int i = 0; i < 100; i++) {
+      byte[] pk = ("pk-" + i).getBytes(StandardCharsets.UTF_8);
+      int shard = TableSetup.shardFor(pk, numShards);
+      assertThat(TableSetup.shardFor(pk, numShards), equalTo(shard));
+      assertThat(shard, greaterThanOrEqualTo(0));
+      assertThat(shard, lessThan(numShards));
+      shards.add(shard);
+    }
+    assertThat(shards.size(), greaterThan(1));
+
+    // facet: full coverage at the non-power-of-two 10.
+    Set<Integer> tenShards = new HashSet<>();
+    for (int i = 0; i < 500; i++) {
+      byte[] pk = ("pk-" + i).getBytes(StandardCharsets.UTF_8);
+      int shard = TableSetup.shardFor(pk, 10);
+      assertThat(shard, greaterThanOrEqualTo(0));
+      assertThat(shard, lessThan(10));
+      tenShards.add(shard);
+    }
+    assertThat(tenShards, containsInAnyOrder(0, 1, 2, 3, 4, 5, 6, 7, 8, 9));
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // shardForHash (the partition-tuple block base)
+  // 
-------------------------------------------------------------------------------------------
+
+  /**
+   * Pins the partition-tuple shard reduction as {@link 
#shardForPinsMurmur3FixedFloorMod} pins the
+   * primary-key one: a silent change here reshards every partitioned 
destination.
+   */
+  @Test
+  public void shardForHashPinsMurmur3FixedAvalanche() {
+    assertThat(TableSetup.shardForHash(0, 8), equalTo(6));
+    assertThat(TableSetup.shardForHash(1, 8), equalTo(2));
+    // Hashes to a negative int, so this also pins floorMod over abs(hash % n) 
(which gives 4).
+    assertThat(TableSetup.shardForHash(-1, 8), equalTo(0));
+  }
+
+  /**
+   * The avalanche is the point: the tuple hash of a single INTEGER partition 
field IS the value, so
+   * without the mix these 32 values striding by 8 would collapse onto one of 
8 shards.
+   */
+  @Test
+  public void shardForHashSpreadsAStridedValueSpace() {
+    Set<Integer> shards = new HashSet<>();
+    for (int i = 0; i < 32; i++) {
+      int shard = TableSetup.shardForHash(i * 8, 8);
+      assertThat(shard, greaterThanOrEqualTo(0));
+      assertThat(shard, lessThan(8));
+      shards.add(shard);
+    }
+    assertThat(shards, containsInAnyOrder(0, 1, 2, 3, 4, 5, 6, 7));
+  }
+
+  @Test
+  public void shardForHashIsDeterministicAndInRange() {
+    for (int i = -50; i < 50; i++) {
+      int shard = TableSetup.shardForHash(i, 10);
+      assertThat(TableSetup.shardForHash(i, 10), equalTo(shard));
+      assertThat(shard, greaterThanOrEqualTo(0));
+      assertThat(shard, lessThan(10));
+    }
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // partitionShardPlan gate: built iff shards_per_partition < num_shards AND 
spec is partitioned
+  // 
-------------------------------------------------------------------------------------------
+
+  /** Resolves a fresh {@code day(ts)}-partitioned destination under {@code 
config}. */
+  private TableSetup.Dest partitionedDest(CdcWriteConfig config, String 
prefix) {
+    org.apache.iceberg.Schema schema = timestampSchema();
+    TableIdentifier id = uniqueId(prefix);
+    CdcSinkTestUtils.createTable(
+        catalog,
+        id,
+        schema,
+        ImmutableSet.of(1, 2),
+        2,
+        PartitionSpec.builderFor(schema).day("ts").build());
+    Schema sourceSchema = dataSchemaFor(schema);
+    TableSetup setup = tableSetup(config, new TestDestinations(sourceSchema, 
null, null));
+    return setup.get(id.toString(), sourceSchema);
+  }
+
+  /**
+   * The gate matrix: a plan is built iff the cap is below {@code num_shards} 
AND the spec is
+   * partitioned; the default (equal) and an unpartitioned table both bypass 
it.
+   */
+  @Test
+  public void partitionShardPlanBuiltOnlyWhenCappedAndPartitioned() {
+    // facet: cap below num_shards on a partitioned spec => plan.
+    assertThat(
+        partitionedDest(cfg().setShardsPerPartition(4).build(), 
"gate_on").partitionShardPlan(),
+        notNullValue());
+
+    // facet: cap == num_shards (today's default exactly) => no plan.
+    assertThat(
+        partitionedDest(cfg().build(), 
"gate_off_default").partitionShardPlan(), nullValue());
+
+    // facet: unpartitioned table ignores the cap => no plan.
+    TableIdentifier id = v2Table("gate_unpartitioned");
+    TableSetup setup = tableSetup(cfg().setShardsPerPartition(1).build());
+    assertThat(setup.get(id.toString(), DATA_SCHEMA).partitionShardPlan(), 
nullValue());
+  }
+
+  // 
-------------------------------------------------------------------------------------------
+  // Test DynamicDestinations
+  // 
-------------------------------------------------------------------------------------------
+
+  /**
+   * A single-table {@link DynamicDestinations} for tests, with an optional 
create config built from
+   * partition and sort field lists plus table properties (mirroring {@code
+   * OneTableDynamicDestinations}).
+   */
+  private static final class TestDestinations implements DynamicDestinations {
+
+    private final Schema dataSchema;
+    private final @Nullable List<String> partitionFields;
+    private final @Nullable List<String> sortFields;
+    private final @Nullable Map<String, String> tableProperties;
+
+    TestDestinations(
+        Schema dataSchema,
+        @Nullable List<String> partitionFields,
+        @Nullable List<String> sortFields) {
+      this(dataSchema, partitionFields, sortFields, null);
+    }
+
+    TestDestinations(
+        Schema dataSchema,
+        @Nullable List<String> partitionFields,
+        @Nullable List<String> sortFields,
+        @Nullable Map<String, String> tableProperties) {
+      this.dataSchema = dataSchema;
+      this.partitionFields = partitionFields;
+      this.sortFields = sortFields;
+      this.tableProperties = tableProperties;
+    }
+
+    @Override
+    public Schema getDataSchema() {
+      return dataSchema;
+    }
+
+    @Override
+    public Row getData(Row element) {
+      return element;
+    }
+
+    @Override
+    public String getTableStringIdentifier(ValueInSingleWindow<Row> element) {
+      throw new UnsupportedOperationException("not used by TableSetup");
+    }
+
+    @Override
+    public IcebergDestination instantiateDestination(String destination) {
+      @Nullable IcebergTableCreateConfig createConfig = null;
+      if (partitionFields != null || sortFields != null || tableProperties != 
null) {
+        createConfig =
+            IcebergTableCreateConfig.builder()
+                .setSchema(dataSchema)
+                .setPartitionFields(partitionFields)
+                .setSortFields(sortFields)
+                .setTableProperties(tableProperties)
+                .build();
+      }
+      return IcebergDestination.builder()
+          .setTableIdentifier(IcebergUtils.parseTableIdentifier(destination))
+          .setFileFormat(FileFormat.PARQUET)
+          .setTableCreateConfig(createConfig)
+          .build();
+    }
+  }
+}

Reply via email to