Class IcebergWriteResult
- All Implemented Interfaces:
POutput
IcebergIO write: the snapshots each destination table committed, plus
the two diversion outputs the CDC sink can produce.
Only getSnapshots() is always present. getDeadLetterRows() (late-but-valid
records) and getFailedRows() (per-record poison rows) are non-null only for results
built by cdc(org.apache.beam.sdk.Pipeline, org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.KV<java.lang.String, org.apache.beam.sdk.io.iceberg.SnapshotInfo>>, org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>, org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>), that is, by IcebergIO.writeCdcRows, and getFailedRows()
additionally only when error handling was enabled. The append-only sink (
IcebergIO.writeRows) leaves both null rather than exposing outputs that can never carry data.
-
Method Summary
Modifier and TypeMethodDescriptionstatic IcebergWriteResultcdc(Pipeline pipeline, PCollection<KV<String, SnapshotInfo>> snapshots, PCollection<Row> deadLetterRows, @Nullable PCollection<Row> failedRows) Returns anIcebergWriteResultfor the CDC sink, exposing the committed-snapshotsnapshots, the replayabledeadLetterRows(seegetDeadLetterRows()), and the optional per-recordfailedRows(seegetFailedRows(),nullwhen error handling is off).expand()voidfinishSpecifyingOutput(String transformName, PInput input, PTransform<?, ?> transform) As part of applying the producingPTransform, finalizes this output to make it ready for being used as an input and for running.The replayable dead-letterRows from the CDC sink: records whose grouped pane fired late, i.e.The per-record poison rows diverted by the CDC sink when error handling is enabled (seeWriteCdcRows.withErrorHandling); schema isfailed_row ROW + error_message STRING(seeErrorHandling).The committed snapshots, keyed by destination.
-
Method Details
-
getSnapshots
The committed snapshots, keyed by destination. A window committed in a commit fire that later fails may not re-emit itsSnapshotInfoon the retry (table state is unaffected). -
getDeadLetterRows
The replayable dead-letterRows from the CDC sink: records whose grouped pane fired late, i.e. after the watermark had passed their commit window's end. Schema isrecord ROW<dataSchema> + change_type STRING + sequence_number INT64 + destination STRING. To replay, unnestrecord, mapchange_type/sequence_numberas the sink's control columns, and route each row bydestination. Replaying is only safe while no newer change for those keys has committed; a stale replay's equality delete removes the newer row.- Returns:
- the dead-letter
PCollection<Row>for results produced byIcebergIO.writeCdcRows/cdc(org.apache.beam.sdk.Pipeline, org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.KV<java.lang.String, org.apache.beam.sdk.io.iceberg.SnapshotInfo>>, org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>, org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>). Throws when attempting to retrieve it for the append-only sink (IcebergIO.writeRows) as has no dead-letter output.
-
getFailedRows
The per-record poison rows diverted by the CDC sink when error handling is enabled (seeWriteCdcRows.withErrorHandling); schema isfailed_row ROW + error_message STRING(seeErrorHandling). Distinct fromgetDeadLetterRows()(which carries late-but-valid records).- Returns:
- the failed-rows
PCollection<Row>if error handling was enabled for the CDC sink. Will throw otherwise.
-
cdc
@Internal public static IcebergWriteResult cdc(Pipeline pipeline, PCollection<KV<String, SnapshotInfo>> snapshots, PCollection<Row> deadLetterRows, @Nullable PCollection<Row> failedRows) Returns anIcebergWriteResultfor the CDC sink, exposing the committed-snapshotsnapshots, the replayabledeadLetterRows(seegetDeadLetterRows()), and the optional per-recordfailedRows(seegetFailedRows(),nullwhen error handling is off). -
getPipeline
Description copied from interface:POutput- Specified by:
getPipelinein interfacePOutput
-
expand
Description copied from interface:POutputExpands thisPOutputinto a list of its component outputPValues.- A
PValueexpands to itself. - A tuple or list of
PValues(such asPCollectionTupleorPCollectionList) expands to its componentPValue PValues.
Not intended to be invoked directly by user code.
- A
-
finishSpecifyingOutput
Description copied from interface:POutputAs part of applying the producingPTransform, finalizes this output to make it ready for being used as an input and for running.This includes ensuring that all
PCollectionshaveCodersspecified or defaulted.Automatically invoked whenever this
POutputis output, afterPOutput.finishSpecifyingOutput(String, PInput, PTransform)has been called on each componentPValuereturned byPOutput.expand().- Specified by:
finishSpecifyingOutputin interfacePOutput
-