Class EvaluationContext
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext
- Direct Known Subclasses:
StreamingEvaluationContext
The
EvaluationContext is the result of a pipeline translation and can be used to evaluate / run the pipeline.
However, in some cases pipeline translation involves the early evaluation of some parts of the
pipeline. For example, this is necessary to materialize side-inputs. The EvaluationContext won't re-evaluate such datasets.
-
Nested Class Summary
Nested Classes -
Constructor Summary
ConstructorsModifierConstructorDescriptionprotectedEvaluationContext(Collection<? extends EvaluationContext.NamedDataset<?>> leaves, org.apache.spark.sql.SparkSession session) -
Method Summary
Modifier and TypeMethodDescriptionThe purpose of this utility is to mark the evaluation of Spark actions, both during Pipeline translation, when evaluation is required, and when finally evaluating the pipeline.voidevaluate()Trigger evaluation of all leaf datasets.static <T> voidThe purpose of this utility is to mark the evaluation of Spark actions, both during Pipeline translation, when evaluation is required, and when finally evaluating the pipeline.org.apache.spark.sql.SparkSessionprotected booleanprotected Collection<? extends EvaluationContext.NamedDataset<?>> leaves()The leaf datasets of the translated pipeline that require evaluation.voidstop()Stops the evaluation after the current leaf dataset.
-
Constructor Details
-
EvaluationContext
protected EvaluationContext(Collection<? extends EvaluationContext.NamedDataset<?>> leaves, org.apache.spark.sql.SparkSession session)
-
-
Method Details
-
leaves
The leaf datasets of the translated pipeline that require evaluation. -
evaluate
public void evaluate()Trigger evaluation of all leaf datasets. Returns early oncestop()was called. -
evaluate
The purpose of this utility is to mark the evaluation of Spark actions, both during Pipeline translation, when evaluation is required, and when finally evaluating the pipeline. -
collect
public static <T extends @NonNull Object> T[] collect(String name, org.apache.spark.sql.Dataset<T> ds) The purpose of this utility is to mark the evaluation of Spark actions, both during Pipeline translation, when evaluation is required, and when finally evaluating the pipeline. -
stop
public void stop()Stops the evaluation after the current leaf dataset. Streaming contexts override this to stop their queries. -
isStopped
protected boolean isStopped() -
getSparkSession
public org.apache.spark.sql.SparkSession getSparkSession()
-