ahmedabu98 commented on code in PR #40161: URL: https://github.com/apache/beam/pull/40161#discussion_r4075027390
########## sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/WriteCdcRows.java: ########## @@ -0,0 +1,521 @@ +/* + * 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 com.google.auto.value.AutoValue; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; +import org.apache.beam.sdk.annotations.Internal; +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.schemas.Schema; +import org.apache.beam.sdk.transforms.PTransform; +import org.apache.beam.sdk.values.KV; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.PCollectionTuple; +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.base.Preconditions; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Predicates; +import org.apache.iceberg.catalog.TableIdentifier; +import org.checkerframework.checker.nullness.qual.Nullable; +import org.joda.time.Duration; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * The top-level CDC sink transform: applies a collection of change records to one or more Iceberg + * V2+ tables: + * + * <pre>{@code + * changeRecords.apply(IcebergIO.writeCdcRows(catalogConfig) + * .to(tableId) + * .withSequenceNumberColumn("seq") + * .withTriggeringFrequency(Duration.standardMinutes(1))); + * }</pre> + * + * <h3>Input contract</h3> + * + * <p>Each input {@link Row} is one change record. Its change kind is its native {@link ValueKind}, + * but can be overridden with a string column value using {@link #withChangeTypeColumn}, optionally + * translated via {@link #withChangeTypeMap}. Each row <b>must</b> also carry a per-key monotonic + * (long) sequence number specified by {@link #withSequenceNumberColumn}, which orders a single + * primary key's changes; the column is required in the input schema as a non-nullable {@code + * INT64}; set defaults upstream. These control columns are stripped from the input rows before + * writing to the table. + * + * <p><b>Ordering requirement (not validated):</b> for a given primary key, the input element's + * event-time must be non-decreasing with its (sequence number, kind rank): a higher-sequence change + * never carries an earlier event time, and equal-sequence records (an update's before and after + * images) carry equal event times, so neither half lands in an earlier commit window. A violation + * can corrupt final table state (a lower-sequence equality delete deleting a higher-sequence row + * committed in a later snapshot). + * + * <h3>Semantics</h3> + * + * <p>The sink commits one new snapshot to the destination table per commit window, in ascending + * window order, with an idempotency token written to each snapshot's summary. A retried or + * restarted commit finds the token and skips already-committed windows, making the commit + * effectively-once on runners that honor {@code @RequiresStableInput}. Records whose grouped pane + * fires late (watermark is past their commit window's end) are diverted to the DLQ output, + * accessible with {@link IcebergWriteResult#getDeadLetterRows()}. + * + * <p>Each dead-lettered record nests the data row under {@code record}, beside {@code change_type}, + * {@code sequence_number}, and {@code destination}. + * + * <h3>Sink id</h3> + * + * <p>The sink creates a fresh unique id by default. All commits in a single pipeline run share the + * same id. Commits get stamped with the sink id and a window-end millis token. Set {@link + * #withSinkId} explicitly (and keep it stable) for cross-relaunch idempotency. Don't reuse sink ids + * for batch runs though because all commits fall under a single global window, so a second load + * with the same sink id will recognize the same window-end millis token and skip the commit. A + * batch load's sink id must likewise not be carried into a streaming continuation. + */ +@Internal +@AutoValue +public abstract class WriteCdcRows extends PTransform<PCollection<Row>, IcebergWriteResult> { + + private static final Logger LOG = LoggerFactory.getLogger(WriteCdcRows.class); + + public static final Duration DEFAULT_ALLOWED_LATENESS = Duration.standardHours(6); + + abstract IcebergCatalogConfig getCatalogConfig(); + + abstract @Nullable TableIdentifier getTableIdentifier(); + + abstract @Nullable DynamicDestinations getDynamicDestinations(); + + abstract @Nullable List<String> getEqualityColumns(); + + abstract String getSequenceNumberColumn(); + + abstract @Nullable String getChangeTypeColumn(); + + abstract @Nullable Map<String, String> getChangeTypeMap(); + + abstract int getNumShards(); + + abstract @Nullable Integer getShardsPerPartition(); + + abstract int getSorterMemoryMB(); + + abstract boolean getUpsert(); + + abstract @Nullable Long getTokenHeartbeatMillis(); + + abstract boolean getErrorHandlingEnabled(); + + abstract @Nullable Map<String, String> getSnapshotProperties(); + + /** The sink id namespacing the {@code beam.cdc.} snapshot-summary tokens. */ + abstract String getSinkId(); + + abstract @Nullable Duration getTriggeringFrequency(); + + abstract @Nullable Duration getAllowedLateness(); + + abstract Builder toBuilder(); + + @AutoValue.Builder + abstract static class Builder { + abstract Builder setCatalogConfig(IcebergCatalogConfig catalogConfig); + + abstract Builder setTableIdentifier(TableIdentifier tableIdentifier); + + abstract Builder setDynamicDestinations(DynamicDestinations destinations); + + abstract Builder setEqualityColumns(List<String> equalityColumns); + + abstract Builder setSequenceNumberColumn(String sequenceNumberColumn); + + abstract Builder setChangeTypeColumn(String changeTypeColumn); + + abstract Builder setChangeTypeMap(Map<String, String> changeTypeMap); + + abstract Builder setNumShards(int numShards); + + abstract Builder setShardsPerPartition(Integer shardsPerPartition); + + abstract Builder setSorterMemoryMB(int sorterMemoryMB); + + abstract Builder setUpsert(boolean upsert); + + abstract Builder setTokenHeartbeatMillis(@Nullable Long tokenHeartbeatMillis); + + abstract Builder setErrorHandlingEnabled(boolean errorHandlingEnabled); + + abstract Builder setSnapshotProperties(Map<String, String> snapshotProperties); + + abstract Builder setSinkId(String sinkId); + + abstract Builder setTriggeringFrequency(Duration triggeringFrequency); + + abstract Builder setAllowedLateness(Duration allowedLateness); + + abstract WriteCdcRows build(); + } + + public static WriteCdcRows of(IcebergCatalogConfig catalogConfig) { + return new AutoValue_WriteCdcRows.Builder() + .setCatalogConfig(catalogConfig) + .setSequenceNumberColumn(CdcWriteConfig.DEFAULT_SEQUENCE_NUMBER_COLUMN) + .setNumShards(CdcWriteConfig.DEFAULT_NUM_SHARDS) + .setSorterMemoryMB(CdcWriteConfig.DEFAULT_SORTER_MEMORY_MB) + .setUpsert(false) + .setErrorHandlingEnabled(false) + .setSinkId(UUID.randomUUID().toString()) + .build(); + } + + /** Writes to a single table. Mutually exclusive with {@link #to(DynamicDestinations)}. */ + public WriteCdcRows to(TableIdentifier tableIdentifier) { + return toBuilder().setTableIdentifier(tableIdentifier).build(); + } + + /** + * Writes to multiple tables. Mutually exclusive with {@link #to(TableIdentifier)}. + * + * <p>The sink reads the control columns from the raw element and writes {@link + * DynamicDestinations#getData}, whose schema must match the destination table's and must exclude + * control columns. + */ + public WriteCdcRows to(DynamicDestinations destinations) { + return toBuilder().setDynamicDestinations(destinations).build(); + } + + /** + * Columns that define a row's identity (the Iceberg equality-delete fields). Defaults to the + * destination table's identifier (primary-key) fields. Tables may be partitioned on non-key + * columns; partition source columns must be equality columns only under {@link #withUpsert} or a + * {@link #withShardsPerPartition} cap. + */ + public WriteCdcRows withEqualityColumns(List<String> columns) { + return toBuilder().setEqualityColumns(columns).build(); + } + + /** + * 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}. + */ + public WriteCdcRows withSequenceNumberColumn(String column) { + return toBuilder().setSequenceNumberColumn(column).build(); + } + + /** + * When set, reads the change kind from this column instead of the element's native {@link + * ValueKind}. Must be declared as a non-nullable {@code STRING} in the input schema. The column + * is stripped from the written rows. + */ + public WriteCdcRows withChangeTypeColumn(String column) { + return toBuilder().setChangeTypeColumn(column).build(); + } + + /** + * Mapping from {@link #withChangeTypeColumn} values to {@link ValueKind} names (e.g. for + * Debezium: {@code {"c": "INSERT", "u": "UPDATE_AFTER", "d": "DELETE"}}). Requires {@link + * #withChangeTypeColumn} to also be set. + */ + public WriteCdcRows withChangeTypeMap(Map<String, String> changeTypeMap) { + return toBuilder().setChangeTypeMap(changeTypeMap).build(); Review Comment: Done -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
