Class EqualityConvertPlanner

java.lang.Object
org.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>
org.apache.iceberg.flink.maintenance.operator.EqualityConvertPlanner
All Implemented Interfaces:
Serializable, org.apache.flink.api.common.state.CheckpointListener, org.apache.flink.streaming.api.operators.Input<Trigger>, org.apache.flink.streaming.api.operators.KeyContext, org.apache.flink.streaming.api.operators.KeyContextHandler, org.apache.flink.streaming.api.operators.OneInputStreamOperator<Trigger,ReadCommand>, org.apache.flink.streaming.api.operators.StreamOperator<ReadCommand>, org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.CheckpointedStreamOperator, org.apache.flink.streaming.api.operators.YieldingOperator<ReadCommand>

@Internal public class EqualityConvertPlanner extends org.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand> implements org.apache.flink.streaming.api.operators.OneInputStreamOperator<Trigger,ReadCommand>
Planner for the equality delete conversion pipeline. For each trigger, it picks the oldest staging snapshot that hasn't been converted yet and emits ReadCommands describing the files its downstream readers and workers must process.

Each trigger runs three steps in order:

  1. refreshStagingCursor(org.apache.iceberg.Snapshot): updates lastStagingSnapshotId from the most recent committer marker on the target branch.
  2. ensureIndexCurrent(org.apache.iceberg.Snapshot, org.apache.iceberg.flink.maintenance.operator.EqualityConvertPlanner.LastCommittedWork, org.apache.iceberg.Snapshot): rebuilds the worker index from main when there is no index yet, when external commits (e.g. compaction) have advanced main, or when the staging snapshot picked for this cycle is the one the previous cycle planned.
  3. processStagingSnapshot(org.apache.iceberg.Snapshot, long, java.lang.Long): resolve the chosen staging snapshot's eq deletes against the (now-current) index, pass through any DV files, and index the snapshot's new data files for the next cycle.
Watermarks separate phases that gate the worker's keyed state. The contract is documented on advancePhase().

An EqualityConvertPlan with the current cycle's metadata is emitted via the METADATA_STREAM side output after the read commands.

Assumes a single equality-field set supplied via the builder; staging eq-deletes with a different equalityFieldIds fail fast in retrieveStagingFiles(org.apache.iceberg.Snapshot). Concurrent writes on the target branch are handled by ensureIndexCurrent(org.apache.iceberg.Snapshot, org.apache.iceberg.flink.maintenance.operator.EqualityConvertPlanner.LastCommittedWork, org.apache.iceberg.Snapshot) reindexing from the new main snapshot; commit-time conflicts are caught by RowDelta.validateFromSnapshot.

See Also:
  • Nested Class Summary

    Nested classes/interfaces inherited from class org.apache.flink.streaming.api.operators.AbstractStreamOperator

    org.apache.flink.streaming.api.operators.AbstractStreamOperator.OutputAdjustment<OUT extends Object>
  • Field Summary

    Fields
    Modifier and Type
    Field
    Description
    static final org.apache.flink.util.OutputTag<IndexCommand>
     
    static final org.apache.flink.util.OutputTag<EqualityConvertPlan>
     

    Fields inherited from class org.apache.flink.streaming.api.operators.AbstractStreamOperator

    combinedWatermark, config, lastRecordAttributes1, lastRecordAttributes2, latencyStats, metrics, output, processingTimeService, stateHandler, stateKeySelector1, stateKeySelector2, timeServiceManager
  • Constructor Summary

    Constructors
    Constructor
    Description
    EqualityConvertPlanner(String tableName, String taskName, TableLoader tableLoader, String stagingBranch, String targetBranch, Set<Integer> eqFieldIds)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    void
     
    void
    initializeState(org.apache.flink.runtime.state.StateInitializationContext context)
     
    void
     
    void
    processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<Trigger> element)
     
    void
    snapshotState(org.apache.flink.runtime.state.StateSnapshotContext context)
     

    Methods inherited from class org.apache.flink.streaming.api.operators.AbstractStreamOperator

    beforeInitializeStateHandler, finish, getContainingTask, getCurrentKey, getExecutionConfig, getInternalTimerService, getKeyedStateBackend, getKeyedStateStore, getMetricGroup, getOperatorConfig, getOperatorID, getOperatorName, getOperatorStateBackend, getOrCreateKeyedState, getPartitionedState, getPartitionedState, getProcessingTimeService, getRuntimeContext, getStateKeySelector1, getStateKeySelector2, getTimeServiceManager, getUserCodeClassloader, hasKeyContext1, hasKeyContext2, initializeState, isAsyncKeyOrderedProcessingEnabled, isUsingCustomRawKeyedState, notifyCheckpointAborted, notifyCheckpointComplete, prepareSnapshotPreBarrier, processLatencyMarker, processLatencyMarker1, processLatencyMarker2, processRecordAttributes, processRecordAttributes1, processRecordAttributes2, processWatermark, processWatermark, processWatermark1, processWatermark1, processWatermark2, processWatermark2, processWatermarkStatus, processWatermarkStatus, processWatermarkStatus1, processWatermarkStatus2, reportOrForwardLatencyMarker, setCurrentKey, setKeyContextElement1, setKeyContextElement2, setMailboxExecutor, setProcessingTimeService, setup, snapshotState, useInterruptibleTimers

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait

    Methods inherited from interface org.apache.flink.api.common.state.CheckpointListener

    notifyCheckpointAborted, notifyCheckpointComplete

    Methods inherited from interface org.apache.flink.streaming.api.operators.Input

    processLatencyMarker, processRecordAttributes, processWatermark, processWatermark, processWatermarkStatus

    Methods inherited from interface org.apache.flink.streaming.api.operators.KeyContext

    getCurrentKey, setCurrentKey

    Methods inherited from interface org.apache.flink.streaming.api.operators.KeyContextHandler

    hasKeyContext

    Methods inherited from interface org.apache.flink.streaming.api.operators.OneInputStreamOperator

    setKeyContextElement

    Methods inherited from interface org.apache.flink.streaming.api.operators.StreamOperator

    finish, getMetricGroup, getOperatorAttributes, getOperatorID, initializeState, prepareSnapshotPreBarrier, setKeyContextElement1, setKeyContextElement2, snapshotState
  • Field Details

    • METADATA_STREAM

      public static final org.apache.flink.util.OutputTag<EqualityConvertPlan> METADATA_STREAM
    • CLEAR_BROADCAST_STREAM

      public static final org.apache.flink.util.OutputTag<IndexCommand> CLEAR_BROADCAST_STREAM
  • Constructor Details

  • Method Details

    • open

      public void open() throws Exception
      Specified by:
      open in interface org.apache.flink.streaming.api.operators.StreamOperator<ReadCommand>
      Overrides:
      open in class org.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>
      Throws:
      Exception
    • initializeState

      public void initializeState(org.apache.flink.runtime.state.StateInitializationContext context) throws Exception
      Specified by:
      initializeState in interface org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.CheckpointedStreamOperator
      Overrides:
      initializeState in class org.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>
      Throws:
      Exception
    • snapshotState

      public void snapshotState(org.apache.flink.runtime.state.StateSnapshotContext context) throws Exception
      Specified by:
      snapshotState in interface org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.CheckpointedStreamOperator
      Overrides:
      snapshotState in class org.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>
      Throws:
      Exception
    • processElement

      public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<Trigger> element) throws Exception
      Specified by:
      processElement in interface org.apache.flink.streaming.api.operators.Input<Trigger>
      Throws:
      Exception
    • close

      public void close() throws Exception
      Specified by:
      close in interface org.apache.flink.streaming.api.operators.StreamOperator<ReadCommand>
      Overrides:
      close in class org.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>
      Throws:
      Exception