Packages

object SymmetricHashJoinStateManager

Linear Supertypes
AnyRef, Any
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. SymmetricHashJoinStateManager
  2. AnyRef
  3. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. Protected

Type Members

  1. case class KeyToValuePair(key: UnsafeRow = null, value: UnsafeRow = null, matched: Boolean = false) extends Product with Serializable

    Helper class for representing data key to (value, matched).

    Helper class for representing data key to (value, matched). Designed for object reuse.

  2. case class ValueAndMatchPair(value: UnsafeRow, matched: Boolean) extends Product with Serializable

    Helper class for representing data (value, matched).

Value Members

  1. final def !=(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  2. final def ##: Int
    Definition Classes
    AnyRef → Any
  3. final def ==(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  4. def allStateStoreNames(joinSides: JoinSide*): Seq[String]
  5. def apply(joinSide: JoinSide, inputValueAttributes: Seq[Attribute], joinKeys: Seq[Expression], stateInfo: Option[StatefulOperatorStateInfo], storeConf: StateStoreConf, hadoopConf: Configuration, partitionId: Int, keyToNumValuesStateStoreCkptId: Option[String], keyWithIndexToValueStateStoreCkptId: Option[String], stateFormatVersion: Int, skippedNullValueCount: Option[SQLMetric] = None, useStateStoreCoordinator: Boolean = true, snapshotOptions: Option[SnapshotOptions] = None, joinStoreGenerator: JoinStateManagerStoreGenerator): SymmetricHashJoinStateManager

    Factory method to determines which version of the join state manager should be created

  6. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  7. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.CloneNotSupportedException]) @IntrinsicCandidate() @native()
  8. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  9. def equals(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef → Any
  10. final def getClass(): Class[_ <: AnyRef]
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  11. def getSchemaForStateStores(joinSide: JoinSide, inputValueAttributes: Seq[Attribute], joinKeys: Seq[Expression], stateFormatVersion: Int): Map[String, (StructType, StructType)]
  12. def getSchemasForStateStoreWithColFamily(joinSide: JoinSide, inputValueAttributes: Seq[Attribute], joinKeys: Seq[Expression], stateFormatVersion: Int): Map[String, StateStoreColFamilySchema]

    Retrieves the schemas used for join operator state stores that use column families

  13. def getStateStoreCheckpointId(storeName: String, partitionId: Int, stateStoreCkptIds: Option[Array[Array[String]]], useColumnFamiliesForJoins: Boolean = false): Option[String]

    Stream-stream join has 4 state stores instead of one.

    Stream-stream join has 4 state stores instead of one. So it will generate 4 different checkpoint IDs when not using virtual column families. This function is used to get the checkpoint ID for a specific state store by the name of the store, partition ID and the stateStoreCkptIds array. The expected names for the stores are generated by getStateStoreName(). If useColumnFamiliesForJoins is true, then it will always return the first checkpoint ID.

    storeName

    the name of the state store

    partitionId

    the partition ID of the state store

    stateStoreCkptIds

    the array of checkpoint IDs for all the state stores

    useColumnFamiliesForJoins

    whether virtual column families are used for the join

    returns

    the checkpoint ID for the specific state store, or None if not found

  14. def getStateStoreCheckpointIds(partitionId: Int, stateStoreCkptIds: Option[Array[Array[String]]], useColumnFamiliesForJoins: Boolean): JoinStateStoreCheckpointId

    Stream-stream join has 4 state stores instead of one.

    Stream-stream join has 4 state stores instead of one. So it will generate 4 different checkpoint IDs using stateStoreCkptIds. They are translated from each joiners' state store into an array through mergeStateStoreCheckpointInfo(). This function is used to read it back into individual state store checkpoint IDs for each store. If useColumnFamiliesForJoins is true, then it will always return the first checkpoint ID.

    partitionId

    the partition ID of the state store

    stateStoreCkptIds

    the array of checkpoint IDs for all the state stores

    useColumnFamiliesForJoins

    whether virtual column families are used for the join

    returns

    the checkpoint IDs for all state stores used by this joiner

  15. def hashCode(): Int
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  16. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  17. val legacyVersion: Int
  18. def mergeStateStoreCheckpointInfo(joinCkptInfo: JoinStateStoreCkptInfo): StatefulOpStateStoreCheckpointInfo

    Stream-stream join has 4 state stores instead of one.

    Stream-stream join has 4 state stores instead of one. So it will generate 4 different checkpoint IDs. The approach we take here is to merge them into one array in the checkpointing path. The driver will process this single checkpointID. When it is passed back to the executors, they will split it back into 4 IDs and use them to load the state. This function is used to merge two checkpoint IDs (each in the form of an array of 1) into one array. The merged array is expected to read back by getStateStoreCheckpointIds().

  19. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  20. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  21. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  22. val supportedVersions: Seq[Int]
  23. final def synchronized[T0](arg0: => T0): T0
    Definition Classes
    AnyRef
  24. def toString(): String
    Definition Classes
    AnyRef → Any
  25. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])
  26. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException]) @native()
  27. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])

Deprecated Value Members

  1. def finalize(): Unit
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.Throwable]) @Deprecated
    Deprecated

    (Since version 9)

Inherited from AnyRef

Inherited from Any

Ungrouped