org.apache.spark.sql.execution.streaming.state
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]
- 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 getStateStoreCheckpointIds(partitionId: Int, stateInfo: StatefulOperatorStateInfo): 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. 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.
- 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)