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:
refreshStagingCursor(org.apache.iceberg.Snapshot): updateslastStagingSnapshotIdfrom the most recent committer marker on the target branch.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.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.
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
FieldsModifier and TypeFieldDescriptionstatic 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
ConstructorsConstructorDescriptionEqualityConvertPlanner(String tableName, String taskName, TableLoader tableLoader, String stagingBranch, String targetBranch, Set<Integer> eqFieldIds) -
Method Summary
Modifier and TypeMethodDescriptionvoidclose()voidinitializeState(org.apache.flink.runtime.state.StateInitializationContext context) voidopen()voidprocessElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<Trigger> element) voidsnapshotState(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, useInterruptibleTimersMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface org.apache.flink.api.common.state.CheckpointListener
notifyCheckpointAborted, notifyCheckpointCompleteMethods inherited from interface org.apache.flink.streaming.api.operators.Input
processLatencyMarker, processRecordAttributes, processWatermark, processWatermark, processWatermarkStatusMethods inherited from interface org.apache.flink.streaming.api.operators.KeyContext
getCurrentKey, setCurrentKeyMethods inherited from interface org.apache.flink.streaming.api.operators.KeyContextHandler
hasKeyContextMethods inherited from interface org.apache.flink.streaming.api.operators.OneInputStreamOperator
setKeyContextElementMethods inherited from interface org.apache.flink.streaming.api.operators.StreamOperator
finish, getMetricGroup, getOperatorAttributes, getOperatorID, initializeState, prepareSnapshotPreBarrier, setKeyContextElement1, setKeyContextElement2, snapshotState
-
Field Details
-
METADATA_STREAM
-
CLEAR_BROADCAST_STREAM
-
-
Constructor Details
-
EqualityConvertPlanner
-
-
Method Details
-
open
- Specified by:
openin interfaceorg.apache.flink.streaming.api.operators.StreamOperator<ReadCommand>- Overrides:
openin classorg.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>- Throws:
Exception
-
initializeState
public void initializeState(org.apache.flink.runtime.state.StateInitializationContext context) throws Exception - Specified by:
initializeStatein interfaceorg.apache.flink.streaming.api.operators.StreamOperatorStateHandler.CheckpointedStreamOperator- Overrides:
initializeStatein classorg.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>- Throws:
Exception
-
snapshotState
public void snapshotState(org.apache.flink.runtime.state.StateSnapshotContext context) throws Exception - Specified by:
snapshotStatein interfaceorg.apache.flink.streaming.api.operators.StreamOperatorStateHandler.CheckpointedStreamOperator- Overrides:
snapshotStatein classorg.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>- Throws:
Exception
-
processElement
-
close
- Specified by:
closein interfaceorg.apache.flink.streaming.api.operators.StreamOperator<ReadCommand>- Overrides:
closein classorg.apache.flink.streaming.api.operators.AbstractStreamOperator<ReadCommand>- Throws:
Exception
-