Record Class IndexCommand

java.lang.Object
java.lang.Record
org.apache.iceberg.flink.maintenance.operator.IndexCommand
All Implemented Interfaces:
Serializable

@Internal public record IndexCommand(IndexCommand.Type type, Long mainSnapshotId, Long indexGeneration, SerializedEqualityValues key, DVPosition rowPosition, long deleteSequenceNumber, int deleteSpecId) extends Record implements Serializable
Command from the EqualityConvertPlanner to the EqualityConvertPKIndex.

IndexCommand.Type.ADD_DATA_ROW adds an existing main-data row to the worker's PK index shard immediately, so this cycle's delete can remove it. IndexCommand.Type.ADD_STAGING_DATA_ROW adds a staging snapshot's new data row; the index applies it only after this cycle's delete resolves, so a re-inserted key survives a same-cycle delete. Both carry the row location as rowPosition. IndexCommand.Type.RESOLVE_DELETE resolves an equality delete key against the index and emits DVPositions for all matching rows. All three flow through the keyed stream and route via key.

IndexCommand.Type.CLEAR_INDEX is emitted on the broadcast side when the index is rebuilt and the worker must evict keyed entries that won't be re-added by that rebuild (e.g. PKs whose data file was removed by CoW). Has no key or row position. Carries the new indexGeneration as the staleness threshold.

rowPosition is the data row's location, set for the two add types and null otherwise; the data sequence number it carries lets the worker apply a delete only to older rows. deleteSequenceNumber is the equality delete's sequence number, set for IndexCommand.Type.RESOLVE_DELETE and unused (-1) otherwise. deleteSpecId is the partition spec id of the delete file, used to scope a partitioned delete to data rows of the same spec; GLOBAL_DELETE_SPEC_ID marks an unpartitioned delete that applies to every spec. Set for IndexCommand.Type.RESOLVE_DELETE and unused (-1) otherwise.

See Also:
  • Field Details

    • GLOBAL_DELETE_SPEC_ID

      public static final int GLOBAL_DELETE_SPEC_ID
      Spec id sentinel for an unpartitioned equality delete, which applies as a global delete.
      See Also:
  • Constructor Details

    • IndexCommand

      public IndexCommand(IndexCommand.Type type, Long mainSnapshotId, Long indexGeneration, SerializedEqualityValues key, DVPosition rowPosition, long deleteSequenceNumber, int deleteSpecId)
      Creates an instance of a IndexCommand record class.
      Parameters:
      type - the value for the type record component
      mainSnapshotId - the value for the mainSnapshotId record component
      indexGeneration - the value for the indexGeneration record component
      key - the value for the key record component
      rowPosition - the value for the rowPosition record component
      deleteSequenceNumber - the value for the deleteSequenceNumber record component
      deleteSpecId - the value for the deleteSpecId record component
  • Method Details

    • addDataRow

      public static IndexCommand addDataRow(Long mainSnapshotId, Long indexGeneration, SerializedEqualityValues key, String filePath, long position, int specId, byte[] partition, long dataSequenceNumber, boolean staging)
    • resolveDelete

      public static IndexCommand resolveDelete(Long mainSnapshotId, Long indexGeneration, SerializedEqualityValues key, long deleteSequenceNumber, int deleteSpecId)
    • clearBeforeReindex

      public static IndexCommand clearBeforeReindex(long mainSnapshotId, long indexGeneration)
    • toString

      public final String toString()
      Returns a string representation of this record class. The representation contains the name of the class, followed by the name and value of each of the record components.
      Specified by:
      toString in class Record
      Returns:
      a string representation of this object
    • hashCode

      public final int hashCode()
      Returns a hash code value for this object. The value is derived from the hash code of each of the record components.
      Specified by:
      hashCode in class Record
      Returns:
      a hash code value for this object
    • equals

      public final boolean equals(Object o)
      Indicates whether some other object is "equal to" this one. The objects are equal if the other object is of the same class and if all the record components are equal. Reference components are compared with Objects::equals(Object,Object); primitive components are compared with '=='.
      Specified by:
      equals in class Record
      Parameters:
      o - the object with which to compare
      Returns:
      true if this object is the same as the o argument; false otherwise.
    • type

      public IndexCommand.Type type()
      Returns the value of the type record component.
      Returns:
      the value of the type record component
    • mainSnapshotId

      public Long mainSnapshotId()
      Returns the value of the mainSnapshotId record component.
      Returns:
      the value of the mainSnapshotId record component
    • indexGeneration

      public Long indexGeneration()
      Returns the value of the indexGeneration record component.
      Returns:
      the value of the indexGeneration record component
    • key

      Returns the value of the key record component.
      Returns:
      the value of the key record component
    • rowPosition

      public DVPosition rowPosition()
      Returns the value of the rowPosition record component.
      Returns:
      the value of the rowPosition record component
    • deleteSequenceNumber

      public long deleteSequenceNumber()
      Returns the value of the deleteSequenceNumber record component.
      Returns:
      the value of the deleteSequenceNumber record component
    • deleteSpecId

      public int deleteSpecId()
      Returns the value of the deleteSpecId record component.
      Returns:
      the value of the deleteSpecId record component