Class StreamingEvaluationContext
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext
org.apache.beam.runners.spark.structuredstreaming.translation.StreamingEvaluationContext
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.
-
Nested Class Summary
Nested classes/interfaces inherited from class org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext
EvaluationContext.NamedDataset<T> -
Method Summary
Modifier and TypeMethodDescriptionvoidevaluate()Starts one streaming query per leaf dataset and blocks until all queries terminate.voidstop()Stops all queries started byevaluate().Methods inherited from class org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext
collect, evaluate, getSparkSession, isStopped, leaves
-
Method Details
-
evaluate
public void evaluate()Starts one streaming query per leaf dataset and blocks until all queries terminate.- Overrides:
evaluatein classEvaluationContext
-
stop
public void stop()Stops all queries started byevaluate(). This method is idempotent and thread safe.- Overrides:
stopin classEvaluationContext
-