public class FlinkStreamingPortablePipelineTranslator extends java.lang.Object implements FlinkPortablePipelineTranslator<FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext>
Modifier and Type | Class and Description |
---|---|
static class |
FlinkStreamingPortablePipelineTranslator.IsFlinkNativeTransform
Predicate to determine whether a URN is a Flink native transform.
|
static class |
FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext
Streaming translation context.
|
FlinkPortablePipelineTranslator.Executor, FlinkPortablePipelineTranslator.TranslationContext
Modifier and Type | Method and Description |
---|---|
FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext |
createTranslationContext(JobInfo jobInfo,
FlinkPipelineOptions pipelineOptions,
java.lang.String confDir,
java.util.List<java.lang.String> filesToStage)
Creates a streaming translation context.
|
java.util.Set<java.lang.String> |
knownUrns() |
FlinkPortablePipelineTranslator.Executor |
translate(FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext context,
org.apache.beam.model.pipeline.v1.RunnerApi.Pipeline pipeline)
Translates the given pipeline.
|
public FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext createTranslationContext(JobInfo jobInfo, FlinkPipelineOptions pipelineOptions, java.lang.String confDir, java.util.List<java.lang.String> filesToStage)
StreamExecutionEnvironment
.public java.util.Set<java.lang.String> knownUrns()
knownUrns
in interface FlinkPortablePipelineTranslator<FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext>
public FlinkPortablePipelineTranslator.Executor translate(FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext context, org.apache.beam.model.pipeline.v1.RunnerApi.Pipeline pipeline)
FlinkPortablePipelineTranslator
translate
in interface FlinkPortablePipelineTranslator<FlinkStreamingPortablePipelineTranslator.StreamingTranslationContext>