Class EqualityConvertDVWriter
java.lang.Object
org.apache.flink.streaming.api.operators.AbstractStreamOperator<DVWriteResult>
org.apache.iceberg.flink.maintenance.operator.EqualityConvertDVWriter
- All Implemented Interfaces:
Serializable,org.apache.flink.api.common.state.CheckpointListener,org.apache.flink.streaming.api.operators.KeyContext,org.apache.flink.streaming.api.operators.KeyContextHandler,org.apache.flink.streaming.api.operators.StreamOperator<DVWriteResult>,org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.CheckpointedStreamOperator,org.apache.flink.streaming.api.operators.TwoInputStreamOperator<DVPosition,,EqualityConvertPlan, DVWriteResult> org.apache.flink.streaming.api.operators.YieldingOperator<DVWriteResult>
@Internal
public class EqualityConvertDVWriter
extends org.apache.flink.streaming.api.operators.AbstractStreamOperator<DVWriteResult>
implements org.apache.flink.streaming.api.operators.TwoInputStreamOperator<DVPosition,EqualityConvertPlan,DVWriteResult>
Keyed parallel resolver that buffers
DVPositions per data-file path, then writes Puffin
DV files directly via BaseDVFileWriter. Plan metadata arrives broadcast on input 2, so
every parallel task sees the cycle's metadata and can validate against the main snapshot.
Each buffered DVPosition carries the data file's specId + encoded partition,
so writing DVs needs no data-manifest scan. Existing DVs are folded into the rewrite (V3 allows
one DV per data file): delete manifests are pruned by partition summary to the cycle's affected
partitions, then filtered to entries referencing the affected data files. No cross-cycle state is
kept; reads are bounded by the pruned manifest set, not the table's full DV history.
Buffered positions are transient per-task. On failure recovery, upstream replay rebuilds them.
- 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 inherited from class org.apache.flink.streaming.api.operators.AbstractStreamOperator
combinedWatermark, config, lastRecordAttributes1, lastRecordAttributes2, latencyStats, metrics, output, processingTimeService, stateHandler, stateKeySelector1, stateKeySelector2, timeServiceManager -
Constructor Summary
ConstructorsConstructorDescriptionEqualityConvertDVWriter(String tableName, String taskName, TableLoader tableLoader, String targetBranch) -
Method Summary
Modifier and TypeMethodDescriptionvoidclose()voidopen()voidprocessElement1(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<DVPosition> record) voidprocessElement2(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<EqualityConvertPlan> record) voidprocessWatermark(org.apache.flink.streaming.api.watermark.Watermark mark) 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, initializeState, isAsyncKeyOrderedProcessingEnabled, isUsingCustomRawKeyedState, notifyCheckpointAborted, notifyCheckpointComplete, prepareSnapshotPreBarrier, processLatencyMarker, processLatencyMarker1, processLatencyMarker2, processRecordAttributes, processRecordAttributes1, processRecordAttributes2, processWatermark, processWatermark1, processWatermark1, processWatermark2, processWatermark2, processWatermarkStatus, processWatermarkStatus, processWatermarkStatus1, processWatermarkStatus2, reportOrForwardLatencyMarker, setCurrentKey, setKeyContextElement1, setKeyContextElement2, setMailboxExecutor, setProcessingTimeService, setup, snapshotState, 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.KeyContext
getCurrentKey, setCurrentKeyMethods inherited from interface org.apache.flink.streaming.api.operators.KeyContextHandler
hasKeyContextMethods inherited from interface org.apache.flink.streaming.api.operators.StreamOperator
finish, getMetricGroup, getOperatorAttributes, getOperatorID, initializeState, prepareSnapshotPreBarrier, setKeyContextElement1, setKeyContextElement2, snapshotStateMethods inherited from interface org.apache.flink.streaming.api.operators.TwoInputStreamOperator
processLatencyMarker1, processLatencyMarker2, processRecordAttributes1, processRecordAttributes2, processWatermark1, processWatermark1, processWatermark2, processWatermark2, processWatermarkStatus1, processWatermarkStatus2
-
Constructor Details
-
EqualityConvertDVWriter
public EqualityConvertDVWriter(String tableName, String taskName, TableLoader tableLoader, String targetBranch)
-
-
Method Details
-
open
- Specified by:
openin interfaceorg.apache.flink.streaming.api.operators.StreamOperator<DVWriteResult>- Overrides:
openin classorg.apache.flink.streaming.api.operators.AbstractStreamOperator<DVWriteResult>- Throws:
Exception
-
processElement1
public void processElement1(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<DVPosition> record) - Specified by:
processElement1in interfaceorg.apache.flink.streaming.api.operators.TwoInputStreamOperator<DVPosition,EqualityConvertPlan, DVWriteResult>
-
processElement2
public void processElement2(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<EqualityConvertPlan> record) - Specified by:
processElement2in interfaceorg.apache.flink.streaming.api.operators.TwoInputStreamOperator<DVPosition,EqualityConvertPlan, DVWriteResult>
-
processWatermark
public void processWatermark(org.apache.flink.streaming.api.watermark.Watermark mark) throws Exception - Overrides:
processWatermarkin classorg.apache.flink.streaming.api.operators.AbstractStreamOperator<DVWriteResult>- Throws:
Exception
-
close
- Specified by:
closein interfaceorg.apache.flink.streaming.api.operators.StreamOperator<DVWriteResult>- Overrides:
closein classorg.apache.flink.streaming.api.operators.AbstractStreamOperator<DVWriteResult>- Throws:
Exception
-