Package org.apache.beam.sdk.io.iceberg.cdc.sink


package org.apache.beam.sdk.io.iceberg.cdc.sink
Iceberg CDC sink: applies a stream or batch of inserts, updates and deletes to Iceberg V2+ tables, exposed as 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 snapshotsCreated for 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 withTokenHeartbeat is 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_id matched a previous load and this run wrote nothing at all. Use a unique sink_id per load, or omit it. Never carry a batch load's sink_id into 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 until remove_orphan_files reclaims 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_files runs.
  • 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-id marker was found in the ancestry but no committed-through token was. Strongly implies that expire_snapshots removed 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 enable withTokenHeartbeat so 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 than withAllowedLateness are a different case: see below.
  • failedRecords: poison records diverted to IcebergWriteResult.getFailedRows() when withErrorHandling() 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_BEFORE records 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 about min(shards_per_partition, distinct keys) files per file kind per commit, with its write parallelism capped to match; 1 is 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 TableMetadata holding 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 by orphanFiles; 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.