Class StreamingEvaluationContext

java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext
org.apache.beam.runners.spark.structuredstreaming.translation.StreamingEvaluationContext

@Internal public class StreamingEvaluationContext extends EvaluationContext
Starts one Spark Structured Streaming query per leaf dataset and blocks until all of them reach a terminal state. Queries end through stop() or the idle stop listener.

Leaf i checkpoints under <checkpointDir>/i, in pipeline graph order. A changed pipeline needs a new checkpoint directory, as with any Spark streaming query.

  • Method Details

    • evaluate

      public void evaluate()
      Starts one streaming query per leaf dataset and blocks until all queries terminate.
      Overrides:
      evaluate in class EvaluationContext
    • stop

      public void stop()
      Stops all queries started by evaluate(). This method is idempotent and thread safe.
      Overrides:
      stop in class EvaluationContext