class SymmetricHashJoinStateManagerV2 extends SymmetricHashJoinStateManager
Streaming join state manager that uses 1 state store with virtual column families enabled. Instead of creating a new state store per join side and store type, this manager uses column families to distinguish data between the original 4 state stores.
- Alphabetic
- By Inheritance
- SymmetricHashJoinStateManagerV2
- SymmetricHashJoinStateManager
- Logging
- AnyRef
- Any
- Hide All
- Show All
- Public
- Protected
Instance Constructors
- new SymmetricHashJoinStateManagerV2(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)
Type Members
- class KeyToNumValuesStore extends StateStoreHandler
A wrapper around a StateStore that stores [key -> number of values].
A wrapper around a StateStore that stores [key -> number of values].
- Attributes
- protected
- Definition Classes
- SymmetricHashJoinStateManager
- class KeyWithIndexToValueStore extends StateStoreHandler
A wrapper around a StateStore that stores the mapping; the mapping depends on the state format version - please refer implementations of KeyWithIndexToValueRowConverter.
A wrapper around a StateStore that stores the mapping; the mapping depends on the state format version - please refer implementations of KeyWithIndexToValueRowConverter.
- Attributes
- protected
- Definition Classes
- SymmetricHashJoinStateManager
- abstract class StateStoreHandler extends Logging
Helper trait for invoking common functionalities of a state store.
Helper trait for invoking common functionalities of a state store.
- Attributes
- protected
- Definition Classes
- SymmetricHashJoinStateManager
- implicit class LogStringContext extends AnyRef
- Definition Classes
- Logging
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 MDC(key: LogKey, value: Any): MDC
- Attributes
- protected
- Definition Classes
- Logging
- def abortIfNeeded(): Unit
Abort any changes to the state store if needed
Abort any changes to the state store if needed
- Definition Classes
- SymmetricHashJoinStateManagerV2 → SymmetricHashJoinStateManager
- def append(key: UnsafeRow, value: UnsafeRow, matched: Boolean): Unit
Append a new value to the key
Append a new value to the key
- Definition Classes
- SymmetricHashJoinStateManager
- 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()
- def commit(): Unit
Commit all the changes to the state store
Commit all the changes to the state store
- Definition Classes
- SymmetricHashJoinStateManagerV2 → SymmetricHashJoinStateManager
- final def eq(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef
- def equals(arg0: AnyRef): Boolean
- Definition Classes
- AnyRef → Any
- def get(key: UnsafeRow): Iterator[UnsafeRow]
Get all the values of a key
Get all the values of a key
- Definition Classes
- SymmetricHashJoinStateManager
- final def getClass(): Class[_ <: AnyRef]
- Definition Classes
- AnyRef → Any
- Annotations
- @IntrinsicCandidate() @native()
- def getInternalRowOfKeyWithIndex(currentKey: UnsafeRow): InternalRow
Projects the key of unsafe row to internal row for printable log message.
Projects the key of unsafe row to internal row for printable log message.
- Definition Classes
- SymmetricHashJoinStateManager
- def getJoinedRows(key: UnsafeRow, generateJoinedRow: (InternalRow) => JoinedRow, predicate: (JoinedRow) => Boolean, excludeRowsAlreadyMatched: Boolean = false): Iterator[JoinedRow]
Get all the matched values for given join condition, with marking matched.
Get all the matched values for given join condition, with marking matched. This method is designed to mark joined rows properly without exposing internal index of row.
- excludeRowsAlreadyMatched
Do not join with rows already matched previously. This is used for right side of left semi join in StreamingSymmetricHashJoinExec only.
- Definition Classes
- SymmetricHashJoinStateManager
- def getLatestCheckpointInfo(): JoinerStateStoreCkptInfo
Get state store checkpoint information of the state store used for this joiner, after they finished data processing.
Get state store checkpoint information of the state store used for this joiner, after they finished data processing.
For SymmetricHashJoinStateManagerV2, this returns the information of the single store used for the entire joiner operator. Both fields of JoinerStateStoreCkptInfo will be identical.
- Definition Classes
- SymmetricHashJoinStateManagerV2 → SymmetricHashJoinStateManager
- def hashCode(): Int
- Definition Classes
- AnyRef → Any
- Annotations
- @IntrinsicCandidate() @native()
- def initializeLogIfNecessary(isInterpreter: Boolean, silent: Boolean): Boolean
- Attributes
- protected
- Definition Classes
- Logging
- def initializeLogIfNecessary(isInterpreter: Boolean): Unit
- Attributes
- protected
- Definition Classes
- Logging
- final def isInstanceOf[T0]: Boolean
- Definition Classes
- Any
- def isTraceEnabled(): Boolean
- Attributes
- protected
- Definition Classes
- Logging
- def iterator: Iterator[KeyToValuePair]
Perform a full scan to provide all available data.
Perform a full scan to provide all available data.
This produces an iterator over the (key, value, match) tuples. Callers are expected to consume fully to clean up underlying iterators correctly.
- Definition Classes
- SymmetricHashJoinStateManager
- val keyAttributes: Seq[AttributeReference]
- Attributes
- protected
- Definition Classes
- SymmetricHashJoinStateManager
- val keySchema: StructType
- Attributes
- protected
- Definition Classes
- SymmetricHashJoinStateManager
- val keyToNumValues: KeyToNumValuesStore
- Attributes
- protected
- Definition Classes
- SymmetricHashJoinStateManager
- val keyWithIndexToValue: KeyWithIndexToValueStore
- Attributes
- protected
- Definition Classes
- SymmetricHashJoinStateManager
- def log: Logger
- Attributes
- protected
- Definition Classes
- Logging
- def logBasedOnLevel(level: Level)(f: => MessageWithContext): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logDebug(msg: => String, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logDebug(entry: LogEntry, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logDebug(entry: LogEntry): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logDebug(msg: => String): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logError(msg: => String, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logError(entry: LogEntry, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logError(entry: LogEntry): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logError(msg: => String): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logInfo(msg: => String, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logInfo(entry: LogEntry, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logInfo(entry: LogEntry): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logInfo(msg: => String): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logName: String
- Attributes
- protected
- Definition Classes
- Logging
- def logTrace(msg: => String, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logTrace(entry: LogEntry, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logTrace(entry: LogEntry): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logTrace(msg: => String): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logWarning(msg: => String, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logWarning(entry: LogEntry, throwable: Throwable): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logWarning(entry: LogEntry): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def logWarning(msg: => String): Unit
- Attributes
- protected
- Definition Classes
- Logging
- def metrics: StateStoreMetrics
Get the state store metrics from the state store manager
Get the state store metrics from the state store manager
- Definition Classes
- SymmetricHashJoinStateManagerV2 → SymmetricHashJoinStateManager
- 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()
- def removeByKeyCondition(removalCondition: (UnsafeRow) => Boolean): Iterator[KeyToValuePair]
Remove using a predicate on keys.
Remove using a predicate on keys.
This produces an iterator over the (key, value, matched) tuples satisfying condition(key), where the underlying store is updated as a side-effect of producing next.
This implies the iterator must be consumed fully without any other operations on this manager or the underlying store being interleaved.
- Definition Classes
- SymmetricHashJoinStateManager
- def removeByValueCondition(removalCondition: (UnsafeRow) => Boolean): Iterator[KeyToValuePair]
Remove using a predicate on values.
Remove using a predicate on values.
At a high level, this produces an iterator over the (key, value, matched) tuples such that value satisfies the predicate, where producing an element removes the value from the state store and producing all elements with a given key updates it accordingly.
This implies the iterator must be consumed fully without any other operations on this manager or the underlying store being interleaved.
- Definition Classes
- SymmetricHashJoinStateManager
- 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])
- def withLogContext(context: Map[String, String])(body: => Unit): Unit
- Attributes
- protected
- Definition Classes
- Logging
Deprecated Value Members
- def finalize(): Unit
- Attributes
- protected[lang]
- Definition Classes
- AnyRef
- Annotations
- @throws(classOf[java.lang.Throwable]) @Deprecated
- Deprecated
(Since version 9)