org.apache.spark.sql.execution.streaming.operators.stateful.join
SymmetricHashJoinStateManager
Companion class SymmetricHashJoinStateManager
object SymmetricHashJoinStateManager
- Alphabetic
- By Inheritance
- SymmetricHashJoinStateManager
- AnyRef
- Any
- Hide All
- Show All
- Public
- Protected
Type Members
- 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.
- case class ValueAndMatchPair(value: UnsafeRow, matched: Boolean) extends Product with Serializable
Helper class for representing data (value, matched).
Value Members
- final def !=(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
- final def ##: Int
- Definition Classes
- AnyRef → Any
- final def ==(arg0: Any): Boolean
- Definition Classes
- AnyRef → Any
- def allStateStoreNames(joinSides: JoinSide*): Seq[String]
- 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
- final def asInstanceOf[T0]: T0
- Definition Classes
- Any
- def clone(): AnyRef
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.CloneNotSupportedException]) @IntrinsicCandidate() @native()
- final def eq(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
- def equals(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef → Any
- final def getClass(): Class[_ <: AnyRef]
- Definition Classes
- AnyRef → Any
- Annotations
- @IntrinsicCandidate() @native()
- def getSchemaForStateStores(joinSide: JoinSide, inputValueAttributes: Seq[Attribute], joinKeys: Seq[Expression], stateFormatVersion: Int): Map[String, (StructType, StructType)]
- 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
- 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
- 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
- def hashCode(): Int
- Definition Classes
- AnyRef → Any
- Annotations
- @IntrinsicCandidate() @native()
- final def isInstanceOf[T0]: Boolean
- Definition Classes
- Any
- val legacyVersion: Int
- 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(). - final def ne(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
- final def notify(): Unit
- Definition Classes
- AnyRef
- Annotations
- @IntrinsicCandidate() @native()
- final def notifyAll(): Unit
- Definition Classes
- AnyRef
- Annotations
- @IntrinsicCandidate() @native()
- val supportedVersions: Seq[Int]
- final def synchronized[T0](arg0: => T0): T0
- Definition Classes
- AnyRef
- def toString(): String
- Definition Classes
- AnyRef → Any
- final def wait(arg0: Long, arg1: Int): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.InterruptedException])
- final def wait(arg0: Long): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.InterruptedException]) @native()
- final def wait(): Unit
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.InterruptedException])
Deprecated Value Members
- def finalize(): Unit
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.Throwable]) @Deprecated
- Deprecated
(Since version 9)