Class SparkSessionFactory
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.SparkSessionFactory
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic classKryoRegistratorfor Spark to serialize broadcast variables used for side-inputs. -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionstatic org.apache.spark.sql.SparkSessionReturns theSparkSessionfor a pipeline, paired withrelease(org.apache.spark.sql.SparkSession).static voidrelease(org.apache.spark.sql.SparkSession session) Releases a session fromacquire(org.apache.beam.runners.spark.structuredstreaming.SparkStructuredStreamingPipelineOptions)and stops it when no longer used.static org.apache.spark.sql.SparkSession.BuildersessionBuilder(String master) Creates Spark session builder with some optimizations for local mode, e.g.
-
Constructor Details
-
SparkSessionFactory
public SparkSessionFactory()
-
-
Method Details
-
acquire
public static org.apache.spark.sql.SparkSession acquire(SparkStructuredStreamingPipelineOptions options) Returns theSparkSessionfor a pipeline, paired withrelease(org.apache.spark.sql.SparkSession). -
release
public static void release(org.apache.spark.sql.SparkSession session) Releases a session fromacquire(org.apache.beam.runners.spark.structuredstreaming.SparkStructuredStreamingPipelineOptions)and stops it when no longer used. The stop runs under the lock, a pipeline starting meanwhile creates a new session. -
sessionBuilder
Creates Spark session builder with some optimizations for local mode, e.g. in tests.
-