java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.PipelineTranslator
Direct Known Subclasses:
PipelineTranslatorCommon

@Internal public abstract class PipelineTranslator extends Object
The pipeline translator translates a Beam Pipeline into a Spark correspondence, that can then be evaluated.

The translation involves traversing the hierarchy of a pipeline multiple times:

  1. Detect if streaming mode is required.
  2. Identify datasets that are repeatedly used as input and should be cached.
  3. And finally, translate each primitive or composite PTransform that is known and supported into its Spark correspondence. If a composite is not supported, it will be expanded further into its parts and translated then.
  • Constructor Details

    • PipelineTranslator

      public PipelineTranslator()
  • Method Details