Class EqualityConvertCommitter
- 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<Trigger>,org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.CheckpointedStreamOperator,org.apache.flink.streaming.api.operators.TwoInputStreamOperator<DVWriteResult,,EqualityConvertPlan, Trigger> org.apache.flink.streaming.api.operators.YieldingOperator<Trigger>
DVWriteResults from parallel
EqualityConvertDVWriter instances (input 1) and an EqualityConvertPlan from the
planner (input 2). Assembles the final file lists and commits using a RowDelta operation
once the plan result and done-timestamp watermark have both arrived.
The commit is gated on the plan's done-timestamp watermark.
Watermarks are forwarded only after the cycle commits, never mid-cycle. The LockRemover releases the maintenance lock once a watermark past the trigger's start epoch
reaches it. The planner emits phase watermarks in the middle of a cycle; forwarding those would
release the lock before this commit, letting the TriggerManager start a concurrent cycle that
re-processes the same uncommitted staging snapshot.
Emits a Trigger after each cycle (commit, no-op, or error) so the downstream TaskResultAggregator can track task completion. This is the sole source of Trigger records for
the Aggregator.
No-op vs error: a no-op cycle (empty plan result from
EqualityConvertPlanner.emitNoOpResult) returns early in commitIfNeeded without writing
anything; the Trigger emit in processWatermark still happens, so the maintenance task
completes cleanly. Errors are reported via TaskResultAggregator.ERROR_STREAM side output;
the Aggregator collects them and surfaces failure on its own watermark.
The committer is intentionally stateless: bufferedResults and planResult are
not checkpointed. The maintenance framework ensures mutual exclusive tasks, and on restart the
planner re-derives its position from COMMITTED_STAGING_SNAPSHOT_PROPERTY on main, so the
committer never receives a plan for an already-committed staging snapshot.
- 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
FieldsFields inherited from class org.apache.flink.streaming.api.operators.AbstractStreamOperator
combinedWatermark, config, lastRecordAttributes1, lastRecordAttributes2, latencyStats, metrics, output, processingTimeService, stateHandler, stateKeySelector1, stateKeySelector2, timeServiceManager -
Constructor Summary
ConstructorsConstructorDescriptionEqualityConvertCommitter(String tableName, String taskName, TableLoader tableLoader, String stagingBranch, String targetBranch) -
Method Summary
Modifier and TypeMethodDescriptionvoidclose()voidopen()voidprocessElement1(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<DVWriteResult> 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
-
Field Details
-
COMMITTED_STAGING_SNAPSHOT_PROPERTY
- See Also:
-
-
Constructor Details
-
EqualityConvertCommitter
public EqualityConvertCommitter(String tableName, String taskName, TableLoader tableLoader, String stagingBranch, String targetBranch)
-
-
Method Details
-
open
-
processElement1
public void processElement1(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<DVWriteResult> record) - Specified by:
processElement1in interfaceorg.apache.flink.streaming.api.operators.TwoInputStreamOperator<DVWriteResult,EqualityConvertPlan, Trigger>
-
processElement2
public void processElement2(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<EqualityConvertPlan> record) - Specified by:
processElement2in interfaceorg.apache.flink.streaming.api.operators.TwoInputStreamOperator<DVWriteResult,EqualityConvertPlan, Trigger>
-
processWatermark
public void processWatermark(org.apache.flink.streaming.api.watermark.Watermark mark) throws Exception -
close
-