Class WriteCdcRows

java.lang.Object
org.apache.beam.sdk.transforms.PTransform<PCollection<Row>,IcebergWriteResult>
org.apache.beam.sdk.io.iceberg.cdc.sink.WriteCdcRows
All Implemented Interfaces:
Serializable, HasDisplayData

@Internal public abstract class WriteCdcRows extends PTransform<PCollection<Row>,IcebergWriteResult>
The top-level CDC sink transform: applies a collection of change records to one or more Iceberg V2+ tables:

 changeRecords.apply(IcebergIO.writeCdcRows(catalogConfig)
     .to(tableId)
     .withSequenceNumberColumn("seq")
     .withTriggeringFrequency(Duration.standardMinutes(1)));
 

Input contract

Each input Row is one change record. Its change kind is its native ValueKind, but can be overridden with a string column value using withChangeTypeColumn(java.lang.String), optionally translated via withChangeTypeMap(java.util.Map<java.lang.String, java.lang.String>). Each row must also carry a per-key monotonic (long) sequence number specified by withSequenceNumberColumn(java.lang.String), which orders a single primary key's changes; the column is required in the input schema as a non-nullable INT64; set defaults upstream. These control columns are stripped from the input rows before writing to the table.

Ordering requirement (not validated): 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).

Semantics

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 @RequiresStableInput. Records whose grouped pane fires late (watermark is past their commit window's end) are diverted to the DLQ output, accessible with IcebergWriteResult.getDeadLetterRows().

Each dead-lettered record nests the data row under record, beside change_type, sequence_number, and destination.

Sink id

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 withSinkId(java.lang.String) 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.

See Also:
  • Field Details

    • DEFAULT_ALLOWED_LATENESS

      public static final Duration DEFAULT_ALLOWED_LATENESS
  • Constructor Details

    • WriteCdcRows

      public WriteCdcRows()
  • Method Details

    • of

      public static WriteCdcRows of(IcebergCatalogConfig catalogConfig)
    • to

      public WriteCdcRows to(org.apache.iceberg.catalog.TableIdentifier tableIdentifier)
      Writes to a single table. Mutually exclusive with to(DynamicDestinations).
    • to

      public WriteCdcRows to(DynamicDestinations destinations)
      Writes to multiple tables. Mutually exclusive with to(TableIdentifier).

      The sink reads the control columns from the raw element and writes DynamicDestinations.getData(org.apache.beam.sdk.values.Row), whose schema must match the destination table's and must exclude control columns.

    • withEqualityColumns

      public WriteCdcRows withEqualityColumns(List<String> columns)
      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 withUpsert(boolean) or a withShardsPerPartition(int) cap.
    • withSequenceNumberColumn

      public WriteCdcRows withSequenceNumberColumn(String column)
      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 INT64 in the input schema. The column is stripped from the written rows. Defaults to "_commit_snapshot_sequence_number".
    • withChangeTypeColumn

      public WriteCdcRows withChangeTypeColumn(String column)
      When set, reads the change kind from this column instead of the element's native ValueKind. Must be declared as a non-nullable STRING in the input schema. The column is stripped from the written rows.
    • withChangeTypeMap

      public WriteCdcRows withChangeTypeMap(Map<String,String> changeTypeMap)
      Mapping from withChangeTypeColumn(java.lang.String) values to ValueKind names (e.g. for Debezium: {"c": "INSERT", "u": "UPDATE_AFTER", "d": "DELETE"}). Requires withChangeTypeColumn(java.lang.String) to also be set.
    • withNumShards

      public WriteCdcRows withNumShards(int numShards)
      The number of deterministic primary-key-hash shards per destination. Controls the sink's write-parallelism knob. Defaults to 16.

      Too low may cause a write bottleneck with a growing commit backlog, too high increases the sink's output file count (num_shards x touched partitions files per commit window).

      On a partitioned table, a commit window writes up to this many files per touched partition. withShardsPerPartition(int) lowers that per-partition count while leaving this number unchanged, so the file count drops without lowering the ceiling on total write parallelism. Each individual partition is then written by at most that many shards.

    • withShardsPerPartition

      public WriteCdcRows withShardsPerPartition(int shardsPerPartition)
      Caps how many shards one partition may occupy (default: num_shards, i.e. all shards may write to all partitions). A lower cap reduces write parallelism per partition but also concentrates the writes in fewer files. A cap of 1 pins each partition to a single writer. Ignored for unpartitioned tables. A cap below num_shards requires every partition source column to be an equality column, making the partition (and with it the shard) a pure function of the primary key.
    • withSorterMemoryMB

      public WriteCdcRows withSorterMemoryMB(int sorterMemoryMB)
      The in-memory buffer size (MB) for the sorter that orders each shard's records by primary key, then sequence number, then change kind, before writing; groups larger than this spill to disk. Must be >= 1. Defaults to 100.
    • withUpsert

      public WriteCdcRows withUpsert(boolean upsert)
      If true, only the after-image of each change (INSERT/UPDATE_AFTER) is required; UPDATE_BEFORE records are dropped and INSERT/UPDATE_AFTER are applied as upserts (equality-delete-then-insert on the primary key). Defaults to false. Requires every partition source column to be an equality column: with before-images dropped, a row that moved partitions could never be deleted from its old one.
    • withTokenHeartbeat

      public WriteCdcRows withTokenHeartbeat(Duration interval)
      Enables a periodic empty token-refresh (heartbeat) commit for each destination with a committed-through token, whether committed in this run or recovered from the table: while idle, the committer re-writes the token into a fresh snapshot every interval, so the token-bearing snapshot stays recent and is less likely to be lost to expire_snapshots before the sink resumes. Disabled by default, and ignored for bounded (batch) input.

      With heartbeat enabled prefer cancel-and-resubmit over drain: the self-re-arming processing-time timer can keep a drain from completing.

    • withErrorHandling

      public WriteCdcRows withErrorHandling()
      Enables per-record error handling. Poison records (unknown change type, missing/null sequence number, null equality value, an unresolvable destination) are diverted to IcebergWriteResult.getFailedRows() (schema failed_row ROW + error_message STRING). Defaults to off where the DoFn throws an error instead. Diversion covers poison detectable in the assignment stage; a failure later in the pipeline fails and retries its bundle as usual.
    • withSnapshotProperties

      public WriteCdcRows withSnapshotProperties(Map<String,String> snapshotProperties)
      Extra user properties to add to every commit's Iceberg snapshot summary. Keys prefixed with beam.cdc. are reserved for the sink's own idempotency/diagnostic tokens and are rejected at construction.
    • withSinkId

      public WriteCdcRows withSinkId(String sinkId)
      A stable identifier for this sink, used to namespace the idempotency tokens written to each commit's Iceberg snapshot summary (beam.cdc.sink-id, beam.cdc.committed-through-ms.<sinkId>). Mostly useful when re-using the same sinkId across subsequent streaming relaunches to maintain exactly-once commit idempotency from the same source. If unset, a fresh UUID is used per launch.

      Also useful in batch to make runs idempotent for the same sinkId. A stable sinkId should only be used once though. In batch, the idempotency token is always global-window end, so re-using it will lead to no-op for subsequent loads (they will be skipped at commit time). Do not carry a batch job's sinkId into a streaming job, as it would skip every commit window forever.

    • withTriggeringFrequency

      public WriteCdcRows withTriggeringFrequency(Duration triggeringFrequency)
      The size of each event-time commit window (streaming only).
    • withAllowedLateness

      public WriteCdcRows withAllowedLateness(Duration allowedLateness)
      How far behind the watermark an element's event-time may be before it is dropped entirely, otherwise routed to IcebergWriteResult.getDeadLetterRows() as a late firing. Defaults to DEFAULT_ALLOWED_LATENESS (6 hours) if unset.
    • populateDisplayData

      public void populateDisplayData(DisplayData.Builder builder)
      Description copied from class: PTransform
      Register display data for the given transform or component.

      populateDisplayData(DisplayData.Builder) is invoked by Pipeline runners to collect display data via DisplayData.from(HasDisplayData). Implementations may call super.populateDisplayData(builder) in order to register display data in the current namespace, but should otherwise use subcomponent.populateDisplayData(builder) to use the namespace of the subcomponent.

      By default, does not register any display data. Implementors may override this method to provide their own display data.

      Specified by:
      populateDisplayData in interface HasDisplayData
      Overrides:
      populateDisplayData in class PTransform<PCollection<Row>,IcebergWriteResult>
      Parameters:
      builder - The builder to populate with display data.
      See Also:
    • expand

      public IcebergWriteResult expand(PCollection<Row> input)
      Description copied from class: PTransform
      Override this method to specify how this PTransform should be expanded on the given InputT.

      NOTE: This method should not be called directly. Instead apply the PTransform should be applied to the InputT using the apply method.

      Composite transforms, which are defined in terms of other transforms, should return the output of one of the composed transforms. Non-composite transforms, which do not apply any transforms internally, should return a new unbound output and register evaluators (via backend-specific registration methods).

      Specified by:
      expand in class PTransform<PCollection<Row>,IcebergWriteResult>