Package org.apache.beam.runners.spark.structuredstreaming.translation
package org.apache.beam.runners.spark.structuredstreaming.translation
Internal translators for running Beam pipelines on Spark.
-
ClassDescriptionThe
EvaluationContextis the result of a pipelinetranslationand can be used to evaluate / run the pipeline.The pipeline translator translates a BeamPipelineinto a Spark correspondence, that can then be evaluated.Shared, mutable state during the translation of a pipeline and omitted afterwards.Unresolved translation, allowing to optimize the generated Spark DAG.Factory to create thePipelineTranslatormatching the execution mode of the pipeline.Pipeline translator for streaming pipelines on Spark 4.KryoRegistratorfor Spark to serialize broadcast variables used for side-inputs.Starts one Spark Structured Streaming query per leaf dataset and blocks until all of them reach a terminal state.TransformTranslator<InT extends PInput,OutT extends POutput, TransformT extends PTransform<InT, OutT>> ATransformTranslatorprovides the capability to translate a specific primitive or compositePTransforminto its Spark correspondence.