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