Packages

abstract class MemoryStreamBaseClass[A] extends MemoryStreamBase[A] with MicroBatchStream with SupportsTriggerAvailableNow with Logging

Linear Supertypes
Logging, SupportsTriggerAvailableNow, SupportsAdmissionControl, MicroBatchStream, MemoryStreamBase[A], SparkDataStream, AnyRef, Any
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. MemoryStreamBaseClass
  2. Logging
  3. SupportsTriggerAvailableNow
  4. SupportsAdmissionControl
  5. MicroBatchStream
  6. MemoryStreamBase
  7. SparkDataStream
  8. AnyRef
  9. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. Protected

Instance Constructors

  1. new MemoryStreamBaseClass(id: Int, sparkSession: SparkSession, numPartitions: Option[Int] = None)(implicit arg0: Encoder[A])

Type Members

  1. implicit class LogStringContext extends AnyRef
    Definition Classes
    Logging

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. def MDC(key: LogKey, value: Any): MDC
    Attributes
    protected
    Definition Classes
    Logging
  5. def addData(data: IterableOnce[A]): Offset
  6. def addData(data: A*): connector.read.streaming.Offset
    Definition Classes
    MemoryStreamBase
  7. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  8. val attributes: Seq[AttributeReference]
    Attributes
    protected
    Definition Classes
    MemoryStreamBase
  9. val batches: ListBuffer[Array[UnsafeRow]]

    All batches from lastCommittedOffset + 1 to currentOffset, inclusive.

    All batches from lastCommittedOffset + 1 to currentOffset, inclusive. Stored in a ListBuffer to facilitate removing committed batches.

    Attributes
    protected
  10. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.CloneNotSupportedException]) @IntrinsicCandidate() @native()
  11. def commit(end: connector.read.streaming.Offset): Unit
    Definition Classes
    MemoryStreamBaseClassMemoryStreamBase → SparkDataStream
  12. def createReaderFactory(): PartitionReaderFactory
    Definition Classes
    MemoryStreamBaseClass → MicroBatchStream
  13. var currentOffset: LongOffset
    Attributes
    protected
  14. def deserializeOffset(json: String): connector.read.streaming.Offset
    Definition Classes
    MemoryStreamBaseClassMemoryStreamBase → SparkDataStream
  15. val encoder: ExpressionEncoder[A]
    Definition Classes
    MemoryStreamBase
  16. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  17. def equals(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef → Any
  18. def fullSchema(): StructType
    Definition Classes
    MemoryStreamBase
  19. final def getClass(): Class[_ <: AnyRef]
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  20. def getDefaultReadLimit(): ReadLimit
    Definition Classes
    SupportsAdmissionControl
  21. def hashCode(): Int
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  22. def initialOffset(): connector.read.streaming.Offset
    Definition Classes
    MemoryStreamBaseClassMemoryStreamBase → SparkDataStream
  23. def initializeLogIfNecessary(isInterpreter: Boolean, silent: Boolean): Boolean
    Attributes
    protected
    Definition Classes
    Logging
  24. def initializeLogIfNecessary(isInterpreter: Boolean): Unit
    Attributes
    protected
    Definition Classes
    Logging
  25. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  26. def isTraceEnabled(): Boolean
    Attributes
    protected
    Definition Classes
    Logging
  27. var lastOffsetCommitted: LongOffset

    Last offset that was discarded, or -1 if no commits have occurred.

    Last offset that was discarded, or -1 if no commits have occurred. Note that the value -1 is used in calculations below and isn't just an arbitrary constant.

    Attributes
    protected
  28. def latestOffset(startOffset: connector.read.streaming.Offset, limit: ReadLimit): connector.read.streaming.Offset
    Definition Classes
    MemoryStreamBaseClass → SupportsAdmissionControl
  29. def latestOffset(): connector.read.streaming.Offset
    Definition Classes
    MemoryStreamBaseClass → MicroBatchStream
  30. def log: Logger
    Attributes
    protected
    Definition Classes
    Logging
  31. def logBasedOnLevel(level: Level)(f: => MessageWithContext): Unit
    Attributes
    protected
    Definition Classes
    Logging
  32. def logDebug(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  33. def logDebug(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  34. def logDebug(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  35. def logDebug(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  36. def logError(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  37. def logError(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  38. def logError(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  39. def logError(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  40. def logInfo(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  41. def logInfo(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  42. def logInfo(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  43. def logInfo(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  44. def logName: String
    Attributes
    protected
    Definition Classes
    Logging
  45. def logTrace(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  46. def logTrace(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  47. def logTrace(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  48. def logTrace(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  49. def logWarning(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  50. def logWarning(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  51. def logWarning(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  52. def logWarning(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  53. val logicalPlan: LogicalPlan
    Attributes
    protected
    Definition Classes
    MemoryStreamBase
  54. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  55. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  56. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  57. val output: Seq[Attribute]
    Attributes
    protected
  58. def planInputPartitions(start: connector.read.streaming.Offset, end: connector.read.streaming.Offset): Array[InputPartition]
    Definition Classes
    MemoryStreamBaseClass → MicroBatchStream
  59. def prepareForTriggerAvailableNow(): Unit
    Definition Classes
    MemoryStreamBaseClass → SupportsTriggerAvailableNow
  60. def reportLatestOffset(): connector.read.streaming.Offset
    Definition Classes
    SupportsAdmissionControl
  61. def reset(): Unit
  62. var startOffset: LongOffset
    Attributes
    protected
  63. def stop(): Unit
    Definition Classes
    MemoryStreamBaseClass → SparkDataStream
  64. final def synchronized[T0](arg0: => T0): T0
    Definition Classes
    AnyRef
  65. def toDF(): classic.DataFrame
    Definition Classes
    MemoryStreamBase
  66. def toDS(): classic.Dataset[A]
    Definition Classes
    MemoryStreamBase
  67. lazy val toRow: Serializer[A]
    Attributes
    protected
    Definition Classes
    MemoryStreamBase
  68. def toString(): String
    Definition Classes
    MemoryStreamBaseClass → AnyRef → Any
  69. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])
  70. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException]) @native()
  71. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.InterruptedException])
  72. def withLogContext(context: Map[String, String])(body: => Unit): Unit
    Attributes
    protected
    Definition Classes
    Logging

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 Logging

Inherited from SupportsTriggerAvailableNow

Inherited from SupportsAdmissionControl

Inherited from MicroBatchStream

Inherited from MemoryStreamBase[A]

Inherited from SparkDataStream

Inherited from AnyRef

Inherited from Any

Ungrouped