Class WriteCdcRows
- All Implemented Interfaces:
Serializable,HasDisplayData
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 Summary
FieldsFields inherited from class org.apache.beam.sdk.transforms.PTransform
annotations, displayData, name, resourceHints -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionexpand(PCollection<Row> input) Override this method to specify how thisPTransformshould be expanded on the givenInputT.static WriteCdcRowsof(IcebergCatalogConfig catalogConfig) voidpopulateDisplayData(DisplayData.Builder builder) Register display data for the given transform or component.to(DynamicDestinations destinations) Writes to multiple tables.to(org.apache.iceberg.catalog.TableIdentifier tableIdentifier) Writes to a single table.withAllowedLateness(Duration allowedLateness) How far behind the watermark an element's event-time may be before it is dropped entirely, otherwise routed toIcebergWriteResult.getDeadLetterRows()as a late firing.withChangeTypeColumn(String column) When set, reads the change kind from this column instead of the element's nativeValueKind.withChangeTypeMap(Map<String, String> changeTypeMap) Mapping fromwithChangeTypeColumn(java.lang.String)values toValueKindnames (e.g.withEqualityColumns(List<String> columns) Columns that define a row's identity (the Iceberg equality-delete fields).Enables per-record error handling.withNumShards(int numShards) The number of deterministic primary-key-hash shards per destination.withSequenceNumberColumn(String column) The column holding the per-primary-key monotonic sequence number used to order a single key's changes.withShardsPerPartition(int shardsPerPartition) Caps how many shards one partition may occupy (default:num_shards, i.e.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>).withSnapshotProperties(Map<String, String> snapshotProperties) Extra user properties to add to every commit's Iceberg snapshot summary.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.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 everyinterval, so the token-bearing snapshot stays recent and is less likely to be lost toexpire_snapshotsbefore the sink resumes.withTriggeringFrequency(Duration triggeringFrequency) The size of each event-time commit window (streaming only).withUpsert(boolean upsert) Iftrue, only the after-image of each change (INSERT/UPDATE_AFTER) is required;UPDATE_BEFORErecords are dropped andINSERT/UPDATE_AFTERare applied as upserts (equality-delete-then-insert on the primary key).Methods inherited from class org.apache.beam.sdk.transforms.PTransform
addAnnotation, compose, compose, getAdditionalInputs, getAnnotations, getDefaultOutputCoder, getDefaultOutputCoder, getDefaultOutputCoder, getKindString, getName, getResourceHints, setDisplayData, setResourceHints, toString, validate, validate
-
Field Details
-
DEFAULT_ALLOWED_LATENESS
-
-
Constructor Details
-
WriteCdcRows
public WriteCdcRows()
-
-
Method Details
-
of
-
to
Writes to a single table. Mutually exclusive withto(DynamicDestinations). -
to
Writes to multiple tables. Mutually exclusive withto(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
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 underwithUpsert(boolean)or awithShardsPerPartition(int)cap. -
withSequenceNumberColumn
The column holding the per-primary-key monotonic sequence number used to order a single key's changes. Must be declared as a non-nullableINT64in the input schema. The column is stripped from the written rows. Defaults to "_commit_snapshot_sequence_number". -
withChangeTypeColumn
When set, reads the change kind from this column instead of the element's nativeValueKind. Must be declared as a non-nullableSTRINGin the input schema. The column is stripped from the written rows. -
withChangeTypeMap
Mapping fromwithChangeTypeColumn(java.lang.String)values toValueKindnames (e.g. for Debezium:{"c": "INSERT", "u": "UPDATE_AFTER", "d": "DELETE"}). RequireswithChangeTypeColumn(java.lang.String)to also be set. -
withNumShards
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 partitionsfiles 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
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 of1pins each partition to a single writer. Ignored for unpartitioned tables. A cap belownum_shardsrequires 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
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
Iftrue, only the after-image of each change (INSERT/UPDATE_AFTER) is required;UPDATE_BEFORErecords are dropped andINSERT/UPDATE_AFTERare applied as upserts (equality-delete-then-insert on the primary key). Defaults tofalse. 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
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 everyinterval, so the token-bearing snapshot stays recent and is less likely to be lost toexpire_snapshotsbefore 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
Enables per-record error handling. Poison records (unknown change type, missing/null sequence number, null equality value, an unresolvable destination) are diverted toIcebergWriteResult.getFailedRows()(schemafailed_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
Extra user properties to add to every commit's Iceberg snapshot summary. Keys prefixed withbeam.cdc.are reserved for the sink's own idempotency/diagnostic tokens and are rejected at construction. -
withSinkId
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 samesinkIdacross 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 stablesinkIdshould 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'ssinkIdinto a streaming job, as it would skip every commit window forever. -
withTriggeringFrequency
The size of each event-time commit window (streaming only). -
withAllowedLateness
How far behind the watermark an element's event-time may be before it is dropped entirely, otherwise routed toIcebergWriteResult.getDeadLetterRows()as a late firing. Defaults toDEFAULT_ALLOWED_LATENESS(6 hours) if unset. -
populateDisplayData
Description copied from class:PTransformRegister display data for the given transform or component.populateDisplayData(DisplayData.Builder)is invoked by Pipeline runners to collect display data viaDisplayData.from(HasDisplayData). Implementations may callsuper.populateDisplayData(builder)in order to register display data in the current namespace, but should otherwise usesubcomponent.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:
populateDisplayDatain interfaceHasDisplayData- Overrides:
populateDisplayDatain classPTransform<PCollection<Row>,IcebergWriteResult> - Parameters:
builder- The builder to populate with display data.- See Also:
-
expand
Description copied from class:PTransformOverride this method to specify how thisPTransformshould be expanded on the givenInputT.NOTE: This method should not be called directly. Instead apply the
PTransformshould be applied to theInputTusing theapplymethod.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:
expandin classPTransform<PCollection<Row>,IcebergWriteResult>
-