Class ReadUnboundedTranslator<T>
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.translation.TransformTranslator<PBegin,PCollection<T>,org.apache.beam.sdk.util.construction.SplittableParDo.PrimitiveUnboundedRead<T>>
org.apache.beam.runners.spark.structuredstreaming.translation.streaming.ReadUnboundedTranslator<T>
public class ReadUnboundedTranslator<T>
extends TransformTranslator<PBegin,PCollection<T>,org.apache.beam.sdk.util.construction.SplittableParDo.PrimitiveUnboundedRead<T>>
Translator for
SplittableParDo.PrimitiveUnboundedRead.
Elements arrive in the global window with the record timestamp. Downstream windowing requires
an explicit Window.Assign.
-
Nested Class Summary
Nested classes/interfaces inherited from class org.apache.beam.runners.spark.structuredstreaming.translation.TransformTranslator
TransformTranslator.Context -
Field Summary
Fields inherited from class org.apache.beam.runners.spark.structuredstreaming.translation.TransformTranslator
complexityFactor -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionprotected voidtranslate(org.apache.beam.sdk.util.construction.SplittableParDo.PrimitiveUnboundedRead<T> transform, TransformTranslator<PBegin, PCollection<T>, org.apache.beam.sdk.util.construction.SplittableParDo.PrimitiveUnboundedRead<T>>.Context cxt) Methods inherited from class org.apache.beam.runners.spark.structuredstreaming.translation.TransformTranslator
canTranslate, windowCoder
-
Constructor Details
-
ReadUnboundedTranslator
public ReadUnboundedTranslator()
-
-
Method Details
-
translate
protected void translate(org.apache.beam.sdk.util.construction.SplittableParDo.PrimitiveUnboundedRead<T> transform, TransformTranslator<PBegin, PCollection<T>, org.apache.beam.sdk.util.construction.SplittableParDo.PrimitiveUnboundedRead<T>>.Context cxt) - Specified by:
translatein classTransformTranslator<PBegin,PCollection<T>, org.apache.beam.sdk.util.construction.SplittableParDo.PrimitiveUnboundedRead<T>>
-