c

cloudflow.spark.testkit

TestSparkStreamletContext

class TestSparkStreamletContext extends SparkStreamletContext

An implementation of SparkCtx for unit testing.

readStream reads from a streaming data source (a csv in this case) and prepares a Dataset[In]

writeStream returns a StreamingQuery that pushes the input Dataset[Out] to a MemorySink.

Linear Supertypes
SparkStreamletContext, Serializable, Serializable, Product, Equals, StreamletContext, AnyRef, Any
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. TestSparkStreamletContext
  2. SparkStreamletContext
  3. Serializable
  4. Serializable
  5. Product
  6. Equals
  7. StreamletContext
  8. AnyRef
  9. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. All

Instance Constructors

  1. new TestSparkStreamletContext(streamletRef: String, session: SparkSession, inletTaps: Seq[SparkInletTap[_]], outletTaps: Seq[SparkOutletTap[_]], config: Config = ConfigFactory.empty)

Type Members

  1. case class MountedPathUnavailableException extends Exception with Product with Serializable
    Definition Classes
    StreamletContext

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. val ProcessingTimeInterval: FiniteDuration
  5. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  6. def checkpointDir(dirName: String): String

    Returns the absolute path to a mounted shared storage that can be used to store reliable checkpoints.

    Returns the absolute path to a mounted shared storage that can be used to store reliable checkpoints. Reliable checkpoints lets the Spark application persist its state across restarts and restart from where it last stopped.

    returns

    the absolute path to a mounted shared storage

    Definition Classes
    TestSparkStreamletContextSparkStreamletContext
  7. def clone(): AnyRef
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native() @HotSpotIntrinsicCandidate()
  8. val config: Config
    Definition Classes
    TestSparkStreamletContext → StreamletContext
  9. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  10. def findTopicForPort(port: StreamletPort): Topic
    Definition Classes
    StreamletContext
  11. final def getClass(): Class[_]
    Definition Classes
    AnyRef → Any
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  12. def getMountedPath(volumeMount: VolumeMount): Path
    Definition Classes
    StreamletContext
  13. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  14. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  15. final def notify(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  16. final def notifyAll(): Unit
    Definition Classes
    AnyRef
    Annotations
    @native() @HotSpotIntrinsicCandidate()
  17. def readStream[In](inPort: CodecInlet[In])(implicit encoder: Encoder[In], typeTag: scala.reflect.api.JavaUniverse.TypeTag[In]): Dataset[In]

    Stream from the underlying external storage and return a DataFrame

    Stream from the underlying external storage and return a DataFrame

    inPort

    the inlet port to read from

    returns

    the data read as Dataset[In]

    Definition Classes
    TestSparkStreamletContextSparkStreamletContext
  18. def runtimeBootstrapServers(topic: Topic): String
    Definition Classes
    StreamletContext
  19. val session: SparkSession
    Definition Classes
    SparkStreamletContext
  20. final def streamletConfig: Config
    Definition Classes
    StreamletContext
  21. val streamletRef: String
    Definition Classes
    TestSparkStreamletContext → StreamletContext
  22. final def synchronized[T0](arg0: ⇒ T0): T0
    Definition Classes
    AnyRef
  23. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  24. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... ) @native()
  25. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  26. def writeStream[Out](stream: Dataset[Out], outPort: CodecOutlet[Out], outputMode: OutputMode, optionalTrigger: Option[Trigger] = None)(implicit encoder: Encoder[Out], typeTag: scala.reflect.api.JavaUniverse.TypeTag[Out]): StreamingQuery

    Start the execution of a StreamingQuery that writes the encodedStream to an external storage using the designated portOut

    Start the execution of a StreamingQuery that writes the encodedStream to an external storage using the designated portOut

    stream

    stream used to write the result of execution of the StreamingQuery

    outPort

    the port used to write the result of execution of the StreamingQuery

    outputMode

    the output mode used to write. Valid values Append, Update, Complete

    returns

    the StreamingQuery that starts executing

    Definition Classes
    TestSparkStreamletContextSparkStreamletContext

Deprecated Value Members

  1. def finalize(): Unit
    Attributes
    protected[lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( classOf[java.lang.Throwable] ) @Deprecated
    Deprecated

Inherited from SparkStreamletContext

Inherited from Serializable

Inherited from Serializable

Inherited from Product

Inherited from Equals

Inherited from StreamletContext

Inherited from AnyRef

Inherited from Any

Ungrouped