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.
- Alphabetic
- By Inheritance
- SupportsFineGrainedReplay
- AnyRef
- Any
- Hide All
- Show All
- Public
- Protected
Abstract Value Members
- 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)
- 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
- 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
- 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 hashCode(): Int
- Definition Classes
- AnyRef → Any
- Annotations
- @IntrinsicCandidate() @native()
- final def isInstanceOf[T0]: Boolean
- Definition Classes
- Any
- 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 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
- 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)