Class EqualityConvertPKIndex
- All Implemented Interfaces:
Serializable,org.apache.flink.api.common.functions.Function,org.apache.flink.api.common.functions.RichFunction
SerializedEqualityValues), each instance handles a subset of PK values.
Existing main data is indexed immediately on arrival. A staging snapshot's new data rows
(IndexCommand.Type.ADD_STAGING_DATA_ROW) and every RESOLVE_DELETE instead register an
event-time timer at their phase timestamp and act only when the watermark reaches it. Flink fires
timers in timestamp order, so a delete resolves before the cycle's own new staging rows (a
strictly later phase) join the index: a re-inserted key survives a same-cycle delete regardless
of the order the parallel readers deliver records. Main data is safe to apply eagerly because its
phase precedes the delete, so the delete's watermark only advances once every main row has
arrived.
On a shared staging/target branch, resolution is sequence-number aware: an equality delete deletes an indexed row only when the row's data sequence number is below the delete's. A re-insert of the same key with a higher sequence survives and stays indexed for a later delete. On a separate target branch the committer reassigns data sequence numbers, so every match is deleted; event-time ordering prevents over-deletion.
Stale-index protection runs on two levels. Each key tracks the generation the commands carrying its stored positions were stamped with. The planner hands out a strictly higher generation on every rebuild of the index, so it serves both the equality test below and the ordering test:
- Lazy (per-key): any keyed command at the top of
processElement(org.apache.iceberg.flink.maintenance.operator.IndexCommand, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.ReadOnlyContext, org.apache.flink.util.Collector<org.apache.iceberg.flink.maintenance.operator.DVPosition>)whoseIndexCommand.indexGeneration()differs from the stored one clears stale state and adopts the command's generation. Equality suffices here. Required because the broadcast and keyed inputs are independent streams with no ordering guarantee; without it, an ADD_DATA_ROW that arrived before the broadcast would be wrongly evicted. - Eager (all keys): a CLEAR_INDEX broadcast iterates all keys on the subtask via
KeyedBroadcastProcessFunction.Context.applyToKeyedState(org.apache.flink.api.common.state.StateDescriptor<S, VS>, org.apache.flink.runtime.state.KeyedStateFunction<KS, S>)and clears any whose stored generation is older than the broadcast's. Staleness is ordered by generation. Bounds state size for PKs that were removed from main by an external CoW commit and won't receive any keyed command next cycle.
- See Also:
-
Nested Class Summary
Nested classes/interfaces inherited from class org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction
org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.Context, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.OnTimerContext, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.ReadOnlyContext -
Field Summary
FieldsModifier and TypeFieldDescription -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionvoidonTimer(long timestamp, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues, IndexCommand, IndexCommand, DVPosition>.org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.OnTimerContext ctx, org.apache.flink.util.Collector<DVPosition> out) voidopen(org.apache.flink.api.common.functions.OpenContext context) voidprocessBroadcastElement(IndexCommand cmd, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues, IndexCommand, IndexCommand, DVPosition>.org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.Context ctx, org.apache.flink.util.Collector<DVPosition> out) voidprocessElement(IndexCommand cmd, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues, IndexCommand, IndexCommand, DVPosition>.org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.ReadOnlyContext ctx, org.apache.flink.util.Collector<DVPosition> out) Methods inherited from class org.apache.flink.api.common.functions.AbstractRichFunction
close, getIterationRuntimeContext, getRuntimeContext, setRuntimeContextMethods 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.functions.RichFunction
useInterruptibleTimers
-
Field Details
-
CLEAR_BROADCAST_DESCRIPTOR
-
-
Constructor Details
-
EqualityConvertPKIndex
public EqualityConvertPKIndex(boolean stagingOnTargetBranch)
-
-
Method Details
-
open
- Specified by:
openin interfaceorg.apache.flink.api.common.functions.RichFunction- Overrides:
openin classorg.apache.flink.api.common.functions.AbstractRichFunction- Throws:
Exception
-
processElement
public void processElement(IndexCommand cmd, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues, IndexCommand, throws ExceptionIndexCommand, DVPosition>.org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.ReadOnlyContext ctx, org.apache.flink.util.Collector<DVPosition> out) - Specified by:
processElementin classorg.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues,IndexCommand, IndexCommand, DVPosition> - Throws:
Exception
-
processBroadcastElement
public void processBroadcastElement(IndexCommand cmd, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues, IndexCommand, IndexCommand, DVPosition>.org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.Context ctx, org.apache.flink.util.Collector<DVPosition> out) - Specified by:
processBroadcastElementin classorg.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues,IndexCommand, IndexCommand, DVPosition>
-
onTimer
public void onTimer(long timestamp, org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues, IndexCommand, IndexCommand, DVPosition>.org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction.OnTimerContext ctx, org.apache.flink.util.Collector<DVPosition> out) - Overrides:
onTimerin classorg.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues,IndexCommand, IndexCommand, DVPosition>
-