java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext
Direct Known Subclasses:
StreamingEvaluationContext

@Internal public class EvaluationContext extends Object
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
    Modifier and Type
    Class
    Description
    static interface 
     
  • Constructor Summary

    Constructors
    Modifier
    Constructor
    Description
    protected
    EvaluationContext(Collection<? extends EvaluationContext.NamedDataset<?>> leaves, org.apache.spark.sql.SparkSession session)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    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.
    void
    Trigger evaluation of all leaf datasets.
    static <T> void
    evaluate(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.
    org.apache.spark.sql.SparkSession
     
    protected boolean
     
    The leaf datasets of the translated pipeline that require evaluation.
    void
    Stops the evaluation after the current leaf dataset.

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Constructor Details

  • Method Details

    • leaves

      protected Collection<? extends EvaluationContext.NamedDataset<?>> leaves()
      The leaf datasets of the translated pipeline that require evaluation.
    • evaluate

      public void evaluate()
      Trigger evaluation of all leaf datasets. Returns early once stop() was called.
    • evaluate

      public static <T> void evaluate(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.
    • 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()