Packages

abstract class StreamExecutionContext extends ProgressContext

Holds the mutable state and metrics for a single batch for streaming query.

Linear Supertypes
ProgressContext, Logging, AnyRef, Any
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. StreamExecutionContext
  2. ProgressContext
  3. Logging
  4. AnyRef
  5. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. Protected

Instance Constructors

  1. new StreamExecutionContext(id: UUID, runId: UUID, name: String, triggerClock: Clock, sources: Seq[SparkDataStream], sink: Table, progressReporter: ProgressReporter, batchId: Long, sparkSession: SparkSession)

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. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  6. var batchId: Long
  7. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws(classOf[java.lang.CloneNotSupportedException]) @IntrinsicCandidate() @native()
  8. var currentStatus: StreamingQueryStatus
    Definition Classes
    ProgressContext
  9. var currentTriggerStartTimestamp: Long
    Attributes
    protected
    Definition Classes
    ProgressContext
  10. var endOffsets: StreamProgress

    Stores the end offsets for this batch.

    Stores the end offsets for this batch. Only the scheduler thread should modify this field, and only in atomic steps. Other threads should make a shallow copy if they are going to access this field more than once, since the field's value may change at any time.

  11. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  12. def equals(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef → Any
  13. var execStatsOnLatestExecutedBatch: Option[ExecutionStats]
    Attributes
    protected
    Definition Classes
    ProgressContext
  14. var executionPlan: IncrementalExecution
  15. def finishNoExecutionTrigger(lastExecutedEpochId: Long): Unit

    Finalizes the trigger which did not execute a batch.

    Finalizes the trigger which did not execute a batch.

    Definition Classes
    ProgressContext
  16. def finishTrigger(hasNewData: Boolean, lastExecution: IncrementalExecution, lastEpoch: Long): Unit

    Override of finishTrigger to extract the map from IncrementalExecution.

    Override of finishTrigger to extract the map from IncrementalExecution.

    Definition Classes
    ProgressContext
  17. def finishTrigger(hasNewData: Boolean, sourceToNumInputRowsMap: Map[SparkDataStream, Long], lastExecution: IncrementalExecution, lastEpochId: Long): Unit

    Finalizes the query progress and adds it to list of recent status updates.

    Finalizes the query progress and adds it to list of recent status updates.

    hasNewData

    Whether the sources of this stream had new data for this trigger.

    Definition Classes
    ProgressContext
  18. final def getClass(): Class[_ <: AnyRef]
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  19. def getDuration(key: String): Option[Long]

    Retrieve a measured duration

    Retrieve a measured duration

    Definition Classes
    ProgressContext
  20. def hashCode(): Int
    Definition Classes
    AnyRef → Any
    Annotations
    @IntrinsicCandidate() @native()
  21. val id: UUID
  22. def initializeLogIfNecessary(isInterpreter: Boolean, silent: Boolean): Boolean
    Attributes
    protected
    Definition Classes
    Logging
  23. def initializeLogIfNecessary(isInterpreter: Boolean): Unit
    Attributes
    protected
    Definition Classes
    Logging
  24. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  25. def isTraceEnabled(): Boolean
    Attributes
    protected
    Definition Classes
    Logging
  26. var lastTriggerStartTimestamp: Long
    Attributes
    protected
    Definition Classes
    ProgressContext
  27. var latestOffsets: StreamProgress

    Tracks the latest offsets for each input source.

    Tracks the latest offsets for each input source. Only the scheduler thread should modify this field, and only in atomic steps. Other threads should make a shallow copy if they are going to access this field more than once, since the field's value may change at any time.

  28. def log: Logger
    Attributes
    protected
    Definition Classes
    Logging
  29. def logBasedOnLevel(level: Level)(f: => MessageWithContext): Unit
    Attributes
    protected
    Definition Classes
    Logging
  30. def logDebug(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  31. def logDebug(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  32. def logDebug(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  33. def logDebug(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  34. def logError(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  35. def logError(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  36. def logError(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  37. def logError(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  38. def logInfo(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  39. def logInfo(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  40. def logInfo(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  41. def logInfo(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  42. def logName: String
    Attributes
    protected
    Definition Classes
    Logging
  43. def logTrace(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  44. def logTrace(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  45. def logTrace(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  46. def logTrace(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  47. def logWarning(msg: => String, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  48. def logWarning(entry: LogEntry, throwable: Throwable): Unit
    Attributes
    protected
    Definition Classes
    Logging
  49. def logWarning(entry: LogEntry): Unit
    Attributes
    protected
    Definition Classes
    Logging
  50. def logWarning(msg: => String): Unit
    Attributes
    protected
    Definition Classes
    Logging
  51. var metricWarningLogged: Boolean

    Flag that signals whether any error with input metrics have already been logged

    Flag that signals whether any error with input metrics have already been logged

    Attributes
    protected
    Definition Classes
    ProgressContext
  52. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  53. var newData: Map[SparkDataStream, LogicalPlan]

    Holds the most recent input data for each source.

    Holds the most recent input data for each source.

    Definition Classes
    StreamExecutionContextProgressContext
  54. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  55. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @IntrinsicCandidate() @native()
  56. var offsetSeqMetadata: OffsetSeqMetadata

    Metadata associated with the offset seq of a batch in the query.

    Metadata associated with the offset seq of a batch in the query.

    Definition Classes
    StreamExecutionContextProgressContext
  57. def onExecutionComplete(): Unit
  58. def onExecutionFailure(): Unit
  59. def recordEndOffsets(to: StreamProgress): Unit

    Only used by Real-time Mode.

    Only used by Real-time Mode. For other cases, end offsets are determined in the batch planning phase so it is never need to be updated.

    Definition Classes
    ProgressContext
  60. def recordTriggerOffsets(from: StreamProgress, to: StreamProgress, latest: StreamProgress): Unit

    Record the offsets range this trigger will process.

    Record the offsets range this trigger will process. Call this before updating committedOffsets in StreamExecution to make sure that the correct range is recorded.

    Definition Classes
    ProgressContext
  61. def reportTimeTaken(triggerDetailKey: String, timeTakenMs: Long): Unit

    Reports an input duration for a particular detail key in the next query progress update.

    Reports an input duration for a particular detail key in the next query progress update. Can be used directly instead of reportTimeTaken(key)(body) when the duration is measured asynchronously.

    Definition Classes
    ProgressContext
  62. def reportTimeTaken[T](triggerDetailKey: String)(body: => T): T

    Records the duration of running body for the next query progress update.

    Records the duration of running body for the next query progress update.

    Definition Classes
    ProgressContext
  63. var sinkCommitProgress: Option[StreamWriterCommitProgress]
    Definition Classes
    ProgressContext
  64. var startOffsets: StreamProgress

    Stores the start offset for this batch.

    Stores the start offset for this batch. Only the scheduler thread should modify this field, and only in atomic steps. Other threads should make a shallow copy if they are going to access this field more than once, since the field's value may change at any time.

  65. def startTrigger(): Unit

    Begins recording statistics about query progress for a given trigger.

    Begins recording statistics about query progress for a given trigger.

    Definition Classes
    ProgressContext
  66. final def synchronized[T0](arg0: => T0): T0
    Definition Classes
    AnyRef
  67. def toString(): String
    Definition Classes
    AnyRef → Any
  68. def updateStatusMessage(message: String): Unit

    Updates the message returned in status.

    Updates the message returned in status.

    Definition Classes
    ProgressContext
  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 ProgressContext

Inherited from Logging

Inherited from AnyRef

Inherited from Any

Ungrouped