Class EqualityConvertPKIndex

java.lang.Object
org.apache.flink.api.common.functions.AbstractRichFunction
org.apache.flink.streaming.api.functions.co.BaseBroadcastProcessFunction
org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues,IndexCommand,IndexCommand,DVPosition>
org.apache.iceberg.flink.maintenance.operator.EqualityConvertPKIndex
All Implemented Interfaces:
Serializable, org.apache.flink.api.common.functions.Function, org.apache.flink.api.common.functions.RichFunction

@Internal public class EqualityConvertPKIndex extends org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction<SerializedEqualityValues,IndexCommand,IndexCommand,DVPosition>
Parallel operator that maintains a shard of the primary key index. Keyed by the full serialized primary key (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:

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

    Fields
    Modifier and Type
    Field
    Description
    static final org.apache.flink.api.common.state.MapStateDescriptor<Void,Void>
     
  • Constructor Summary

    Constructors
    Constructor
    Description
    EqualityConvertPKIndex(boolean stagingOnTargetBranch)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    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)
     
    void
    open(org.apache.flink.api.common.functions.OpenContext context)
     
    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)
     
    void
    processElement(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, setRuntimeContext

    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.functions.RichFunction

    useInterruptibleTimers