Package org.apache.beam.sdk.io.iceberg.cdc.sink
IcebergIO.writeCdcRows(catalogConfig). For the transform's API and
semantics see WriteCdcRows.
This is a merge-on-read sink. Within one commit window the writer collapses each primary key's changes to their final state, so a superseded intermediate row is never written at all; a change that reaches back into an already-committed snapshot is written as a PK-only equality delete. No position deletes or deletion vectors are ever written.
Monitoring
All metrics below are queryable through PipelineResult.metrics(); the namespace is the
owning class's fully-qualified name, so filter on
org.apache.beam.sdk.io.iceberg.cdc.sink.*.
Committer volume and latency
The normal throughput picture. None of these indicate a problem by themselves; watch their shape over time.
- snapshotsCreated: CDC snapshots committed. In streaming this should track one per commit window per destination. A flat line while data is flowing means commits have stalled.
- committedDataFiles / committedDeleteFiles: files published. Divide by
snapshotsCreatedfor files per commit. A large ratio points to the small-file problem described below. - committedRecords: rows committed.
- committedEqualityDeleteRecords: equality-delete rows committed.
- committedBytes: total bytes of committed files.
- commitDurationMs (distribution): wall-clock time spent waiting on the Iceberg commit operation. Its tail is your catalog's health: a growing maximum usually means catalog contention, not a Beam problem.
- commitFailures: fires whenever the commit path throws. Iceberg's own optimistic retry runs underneath, so a genuine commit collision is counted only once that is exhausted. The failure is rethrown, so the bundle fails and the runner retries it; occasional increments under heavy concurrent writing are survivable. A sustained rate means the table has too many concurrent writers, and the sink will not make progress past the failing window (commits are strictly ordered).
- heartbeatCommits: empty token-refresh commits emitted while a destination is idle
(only if
withTokenHeartbeatis configured). Each one is a real snapshot. With heartbeat enabled prefer cancel-and-resubmit over drain: the self-re-arming processing-time timer can keep a drain from completing.
Safety tripwires
These should be zero. A nonzero value is not necessarily an outage, but each one means something specific happened that you want to know about.
- alreadyCommittedWindowsSkipped: a commit window was skipped because the table
already carried a committed-through token at or past its end. Expected and benign after a
restart or a retried bundle (we skip to ensure idempotency). Alarming in two cases: (1) in
batch, where the whole load commits under one window, a nonzero value on a fresh
load means a stable
sink_idmatched a previous load and this run wrote nothing at all. Use a uniquesink_idper load, or omit it. Never carry a batch load'ssink_idinto a subsequent streaming run (the batch token is the global-window end, above every streaming window; the streaming committer will refuse a recovered batch token); and (2) in streaming, a persistently nonzero rate means data keeps arriving for windows that already closed, which is a source-lateness problem. The token proves that a window with that end committed, not that the skipped rows are the rows that committed. A logged warning names up to 5 of the window's files plus a count: your handle on the skipped rows untilremove_orphan_filesreclaims them. - orphanFiles: files (data and delete) carried by a skipped window: potential
orphans (a pure redelivery's files are the committed live ones). Wasted storage until
remove_orphan_filesruns. - tokenParseFailures: an unparseable
beam.cdc.*value was found while scanning snapshot ancestry. The committed-through scan continues to older ancestors rather than crash-looping; a bad max-seq or run-spec stamp just reads as absent. Nonzero means something wrote a malformed value into a snapshot summary; investigate before trusting the recovered position. - suspectedTokenExpiry: the sink's own
beam.cdc.sink-idmarker was found in the ancestry but no committed-through token was. Strongly implies thatexpire_snapshotsremoved the token-bearing snapshots while the pipeline was down. Recovery then falls back to the beginning and may re-apply retained windows. If you see this, either lengthen snapshot retention or enablewithTokenHeartbeatso the token stays young. - crossWindowSequenceInversions: a committed window's minimum source sequence number was below an earlier window's committed maximum. This is the detector for a violated ordering contract (see below) and the only tripwire here that can mean silently wrong table contents. It has benign false positives when the two windows touch entirely disjoint primary keys, so treat it as "go and check the source's ordering", not as proof of corruption.
- specMismatchedWindows: signals the partition spec was evolved mid-run. Incremented when a committed window carries equality deletes and uses partition-spec ids different from the run's pinned spec. Those deletes may not reach rows written under the other spec.
Input health
These describe what is arriving, not what the sink did with it.
- deadLetterRecords: records whose grouped pane fired late (the watermark had already
passed their commit window's end), diverted to the replayable dead-letter output (
IcebergWriteResult.getDeadLetterRows()) instead of being applied out of order. Every late pane is diverted, including a window's first (SplitLateData's javadoc explains why). Nonzero means your source is lagging past the commit window. The records are not lost, but nothing consumes them unless you wire that output somewhere. Records later thanwithAllowedLatenessare a different case: see below. - failedRecords: poison records diverted to
IcebergWriteResult.getFailedRows()whenwithErrorHandling()is on: unknown change type, missing sequence number, null equality value, unresolvable destination. Without error handling these fail the pipeline instead. Like the dead letters above, the records are not lost, but you need to wire the output somewhere to avoid dropping them. - upsertUpdateBeforeDropped:
UPDATE_BEFORErecords discarded because the sink is in upsert mode, which needs only after-images. Tells you how much of your input was redundant.
Records more than withAllowedLateness behind the watermark are dropped by the runner
at the GroupByShardKey step. Raising withAllowedLateness is the only way to
capture them (at the cost of more live window state per destination).
Files and maintenance
File count per commit is governed by num_shards × touched partitions, and
total file count by that times the number of commit windows. A table with many partitions and a
short triggering_frequency_seconds produces a lot of small files very quickly, regardless
of the data rate. num_shards (default=16) is the sink's write-parallelism knob. Too low
may lead to a write bottleneck and commit backlog. Too high is fine for the sink but may lead to
slower downstream reads due to small files.
There are three levers that you can use:
triggering_frequency_seconds: increasing it leads to fewer, larger commits, at the cost of end-to-end latency.num_shards: trades throughput for file count. Lowering it reduces file count, but also the write-parallelism ceiling.shards_per_partition(default=num_shards, i.e. no cap): on a partitioned table, cap how many shards one partition's rows may occupy. A touched partition then writes aboutmin(shards_per_partition, distinct keys)files per file kind per commit, with its write parallelism capped to match;1is the pure partition-affine endpoint (one writer, about one data file, per partition).
The trade-off a lowered shards_per_partition makes, stated plainly: effective
write parallelism per destination becomes at most the cap times the number of distinct partitions
receiving data in a window: at 1, one writer per partition. A cap of 1 is right
for a table with many partitions and wrong for a table with three, which it would hold at
three concurrent writers however many workers the pipeline has; such a table wants an
intermediate cap, sized so cap × touched partitions still covers the pipeline's write
parallelism. The cap is moot for an unpartitioned table, where it is simply ignored (with a WARN)
rather than funnelling everything through a single shard. At 1 this is the same
limitation Flink's hash write-distribution mode carries; the dial exists because the
alternative (leaving num_shards as the only lever) forces the same trade with none of the
benefit. It is safe with upsert: both options require every partition source column to be
an equality column (see the partitioning section below), which makes a row's partition a pure
function of its primary key.
Snapshot expiry bounds the sink's own cost, not just read performance
This sink creates one snapshot per commit window per table. The following sink costs grow as snapshot count grows:
- Commit latency: Every commit fire loads the destination's metadata once, then walks its ancestry to find the latest committed-through token. Each commit within that fire publishes a snapshot, which requires refreshing and re-parsing the metadata. Both scale with the number of snapshots in that file.
- Worker heap: Each cached table pins a
TableMetadataholding every retained snapshot, per worker, multiplied by the number of dynamic destinations that worker touches.
Orphan files
The writer stage writes data files and delete files and passes their metadata to the committer stage in the pipeline. Every orphan is an ordinary data or delete file that no snapshot ended up referencing. Orphan files can be produced in three ways:
- A retried or failed writer bundle. Each attempt names its files with a fresh UUID, so whichever attempt's output element loses the retry race leaves its files unreferenced.
- A window skipped because the committed-through token already covers its end. This
happens when a relaunch reuses a stable
sink_id. The incoming element is targeting a window that has already been committed with different files. If the incoming element's files contain different content than the committed ones, they may hold rows the table never received. Those files are left in place so you can recover them. Counted byorphanFiles; the skip WARN names up to 5 of them plus a count. - A failed commit attempt. Iceberg normally deletes the manifests it wrote when a
commit fails, but it cannot when the outcome is unknown (i.e.
CommitStateUnknownException). The commit may in fact have succeeded, and deleting those manifests would corrupt a live snapshot. Only the manifests are left orphaned. The window's own data and delete files are untouched, because the pending bag survives the failure and the retry commits those same files.
All of these are reclaimed by Iceberg's remove_orphan_files. The sink writes
everything through the table's own write path. Give the procedure an age threshold much longer
than your longest in-flight bundle to avoid deleting files a live pipeline is about to reference.
Partition-spec evolution mid-run
Iceberg matches an equality delete to its data files by (spec id, partition), so if a
table's spec changes during streaming writes, the sink will choose to keep writing with the old
spec. Otherwise, it could leave new deletes that never reach old-spec rows, which is data
corruption. If a window happens to still mix specs while carrying equality deletes, it commits
anyway, WARNs, and increments the specMismatchedWindows counter.
If you update a table's spec, run rewrite_data_files first so old-spec rows are
rewritten under the new spec. A relaunched pipeline or in-place update will adopt the new spec.
Partitioning
Tables may be partitioned on any columns, key or not, with any transforms. Two options still
require every partition source column to be an equality column: upsert (before-images are
dropped, so a row that moved partitions could never be deleted from its old one) and a
shards_per_partition cap below num_shards (the shard is derived from the partition
tuple, which must therefore be a pure function of the primary key).
Non-key partitioning sharpens the input contract: every update must carry its
UPDATE_BEFORE, because a moved row whose before-image never arrives leaves a permanent duplicate
in the old partition that no later delete reaches; and UPDATE_BEFORE/DELETE rows
must carry the row's actual old values in the partition source columns, because a nulled non-key
column routes the equality delete to the null partition.
The ordering contract
The sink orders each primary key's changes by the sequence-number column, and commits one snapshot per commit window in ascending window order. That is only sound if the two orderings agree at the source: for a given key, an element's event time must be non-decreasing with its (sequence number, kind rank), so equal-sequence records (an update's before and after images) carry equal event times. If a higher-sequence change lands in an earlier commit window than a lower-sequence one, a stale equality delete can be committed after the row it should not have touched, and the table's final contents are wrong.
This is a source contract. Violations are detected at commit time and visible with the
crossWindowSequenceInversions counter and warning logs, but the commit still proceeds. The fix
is in how event times are assigned upstream. Replaying dead letters is bound by the same
contract: see the replay caveat at IcebergWriteResult#getDeadLetterRows.
-
ClassesClassDescriptionThe shuffle key of one write group: a destination and one of its shards.Holds the serialized metadata for
DataFiles andDeleteFiles of one(destination, shard, window)group, plus the source sequence range the bundle covers.The top-level CDC sink transform: applies a collection of change records to one or more Iceberg V2+ tables: