Class IncrementalChangelogSource

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

public class IncrementalChangelogSource extends PTransform<PBegin,PCollection<Row>>
An Iceberg source that incrementally reads a table's changelogs, processing one snapshot at a time.

Each snapshot is resolved independently. For a given primary key, this source emits the net change the snapshot produces, and not its intermediate states.

Implications: if a writer batches several transitions for the same PK into one snapshot (for example A → B, then B → C), only the endpoints survive. The intermediate B is dropped. A round-trip within a single snapshot (for example A → B → A) is also dropped.

The streaming path uses WatchForSnapshotsSdf for proper per-snapshot watermarks. The bounded path creates the snapshot range up front.

See Also:
  • Constructor Details

    • IncrementalChangelogSource

      public IncrementalChangelogSource(IcebergScanConfig scanConfig)
  • Method Details

    • expand

      public PCollection<Row> expand(PBegin 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<PBegin,PCollection<Row>>