Class PipelineTranslatorStreaming
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.PipelineTranslator
org.apache.beam.runners.spark.structuredstreaming.translation.batch.PipelineTranslatorCommon
org.apache.beam.runners.spark.structuredstreaming.translation.PipelineTranslatorStreaming
Pipeline translator for streaming pipelines on Spark 4. It extends the common registry to reuse
the stateless single output ParDo, Window.Assign, Flatten and Reshuffle translators, which are
safe on a streaming Dataset. Every other primitive fails at translation, the batch translators
for them persist or collect the Dataset, which Spark rejects for streaming plans.
-
Nested Class Summary
Nested classes/interfaces inherited from class org.apache.beam.runners.spark.structuredstreaming.translation.PipelineTranslator
PipelineTranslator.TranslationState, PipelineTranslator.UnresolvedTranslation<InT,T> -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionprotected EvaluationContextcreateEvaluationContext(Collection<? extends EvaluationContext.NamedDataset<?>> leaves, org.apache.spark.sql.SparkSession session, SparkCommonPipelineOptions options) Creates theEvaluationContextfor the translated pipeline.protected <InT extends PInput,OutT extends POutput, TransformT extends PTransform<InT, OutT>>
@Nullable TransformTranslator<InT, OutT, TransformT> getTransformTranslator(TransformT transform) Returns aTransformTranslatorfor the givenPTransformif known.Methods inherited from class org.apache.beam.runners.spark.structuredstreaming.translation.PipelineTranslator
detectStreamingMode, replaceTransforms, translate
-
Constructor Details
-
PipelineTranslatorStreaming
public PipelineTranslatorStreaming()
-
-
Method Details
-
getTransformTranslator
protected <InT extends PInput,OutT extends POutput, @Nullable TransformTranslator<InT,TransformT extends PTransform<InT, OutT>> OutT, getTransformTranslatorTransformT> (TransformT transform) Returns aTransformTranslatorfor the givenPTransformif known.- Overrides:
getTransformTranslatorin classPipelineTranslatorCommon
-
createEvaluationContext
protected EvaluationContext createEvaluationContext(Collection<? extends EvaluationContext.NamedDataset<?>> leaves, org.apache.spark.sql.SparkSession session, SparkCommonPipelineOptions options) Description copied from class:PipelineTranslatorCreates theEvaluationContextfor the translated pipeline.Subclasses may override this to return a specialized context, e.g. to evaluate streaming pipelines.
- Overrides:
createEvaluationContextin classPipelineTranslator
-