Packages

t

org.apache.spark.sql.execution.streaming.state

SupportsFineGrainedReplay

trait SupportsFineGrainedReplay extends AnyRef

This is an optional trait to be implemented by StateStoreProviders that can read the change of state store over batches. This is used by State Data Source with additional options like snapshotStartBatchId or readChangeFeed.

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

Abstract Value Members

  1. abstract def getStateStoreChangeDataReader(startVersion: Long, endVersion: Long, colFamilyNameOpt: Option[String] = None): NextIterator[(RecordType.Value, UnsafeRow, UnsafeRow, Long)]

    Return an iterator that reads all the entries of changelogs from startVersion to endVersion.

    Return an iterator that reads all the entries of changelogs from startVersion to endVersion. Each record is represented by a tuple of (recordType: RecordType.Value, key: UnsafeRow, value: UnsafeRow, batchId: Long) A put record is returned as a tuple(recordType, key, value, batchId) A delete record is return as a tuple(recordType, key, null, batchId)

    startVersion

    starting changelog version

    endVersion

    ending changelog version

    colFamilyNameOpt

    optional column family name to read from

    returns

    iterator that gives tuple(recordType: RecordType.Value, nested key: UnsafeRow, nested value: UnsafeRow, batchId: Long)

  2. abstract def replayStateFromSnapshot(snapshotVersion: Long, endVersion: Long): StateStore

    Return an instance of StateStore representing state data of the given version.

    Return an instance of StateStore representing state data of the given version. The State Store will be constructed from the snapshot at snapshotVersion, and applying delta files up to the endVersion. If there is no snapshot file at snapshotVersion, an exception will be thrown.

    snapshotVersion

    checkpoint version of the snapshot to start with

    endVersion

    checkpoint version to end with

Concrete 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. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  5. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.CloneNotSupportedException]) @IntrinsicCandidate() @native()
  6. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  7. def equals(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef → Any
  8. final def getClass(): Class[_ <: AnyRef]
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  9. def hashCode(): Int
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  10. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  11. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  12. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  13. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  14. def replayReadStateFromSnapshot(snapshotVersion: Long, endVersion: Long): ReadStateStore

    Return an instance of ReadStateStore representing state data of the given version.

    Return an instance of ReadStateStore representing state data of the given version. The State Store will be constructed from the snapshot at snapshotVersion, and applying delta files up to the endVersion. If there is no snapshot file at snapshotVersion, an exception will be thrown. Only implement this if there is read-only optimization for the state store.

    snapshotVersion

    checkpoint version of the snapshot to start with

    endVersion

    checkpoint version to end with

  15. final def synchronized[T0](arg0: => T0): T0
    Definition Classes
    AnyRef
  16. def toString(): String
    Definition Classes
    AnyRef → Any
  17. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])
  18. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException]) @native()
  19. 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