MCPcopy Create free account

hub / github.com/dataArtisans/flink-dataflow / functions

Functions813 in github.com/dataArtisans/flink-dataflow

MethodsideInput
(PCollectionView<T> view)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:130
MethodsideOutput
(TupleTag<T> tag, T output)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:167
MethodsideOutputWithTimestamp
(TupleTag<T> tag, T output, Instant timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:561
MethodsideOutputWithTimestamp
(TupleTag<T> tag, T output, Instant timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:188
MethodsideOutputWithTimestamp
(TupleTag<T> tag, T output, Instant timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:162
MethodsideOutputWithTimestampHelper
(WindowedValue<IN> inElement, T output, Instant timestamp, Collector<WindowedValue<OUT>> outCollector, TupleTa
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:52
MethodsideOutputWithTimestampHelper
(WindowedValue<IN> inElement, T output, Instant timestamp, Collector<WindowedValue<RawUnionValue>> collector,
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundMultiWrapper.java:56
MethodsnapshotOperatorState
(long checkpointId, long timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:584
Methodstart
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:134
Methodstart
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:145
MethodstartBundle
(Context c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:100
MethodstateInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:62
MethodstateInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:509
MethodstateInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:120
MethodsupportsNormalizedKey
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:146
MethodsupportsNormalizedKey
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:190
MethodsupportsSerializationWithKeyNormalization
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:151
MethodsupportsSerializationWithKeyNormalization
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:195
Methodtest
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/StateSerializationTest.java:287
MethodtestAfterCountProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:289
MethodtestAfterWatermarkProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:267
MethodtestCompoundAccumulatingPanesProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:352
MethodtestCompoundProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:312
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/MaybeEmptyTestITCase.java:47
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/TfIdfITCase.java:53
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/FlattenizeITCase.java:52
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/ParDoMultiOutputITCase.java:48
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:60
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/SideInputITCase.java:35
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:80
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesEmptyITCase.java:50
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountITCase.java:58
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:57
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:71
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/JoinExamplesITCase.java:84
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/AvroITCase.java:57
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesITCase.java:51
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupByNullKeyTest.java:77
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:65
MethodtestProgram
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/TopWikipediaSessionsITCase.java:67
MethodtestSessionWindows
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:142
MethodtestSlidingWindows
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:200
MethodtestWithLateness
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:82
MethodtimerInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:72
MethodtimerInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:522
MethodtimerInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:130
Methodtimestamp
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:436
Methodtimestamp
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:117
Methodtimestamp
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:222
Methodtimestamp
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:102
Methodtimestamp
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:105
MethodtoString
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:113
Methodtranslate
Depending on if the job is a Streaming or a Batch one, this method creates the necessary execution environment and pipeline translator, and translates
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkPipelineExecutionEnvironment.java:120
MethodtranslateNode
(Read.Bounded<T> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:145
MethodtranslateNode
(AvroIO.Read.Bound<T> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:164
MethodtranslateNode
(AvroIO.Write.Bound<T> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:197
MethodtranslateNode
(TextIO.Read.Bound<String> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:235
MethodtranslateNode
(TextIO.Write.Bound<T> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:259
MethodtranslateNode
(ConsoleIO.Write.Bound transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:285
MethodtranslateNode
(Write.Bound<T> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:295
MethodtranslateNode
(GroupByKey.GroupByKeyOnly<K, V> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:307
MethodtranslateNode
(GroupByKey<K, V> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:327
MethodtranslateNode
(Combine.PerKey<K, VI, VO> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:345
MethodtranslateNode
(ParDo.Bound<IN, OUT> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:413
MethodtranslateNode
(ParDo.BoundMulti<IN, OUT> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:433
MethodtranslateNode
(Flatten.FlattenPCollectionList<T> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:483
MethodtranslateNode
(View.CreatePCollectionView<R, T> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:500
MethodtranslateNode
(Create.Values<OUT> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:510
MethodtranslateNode
(CoGroupByKey<K> transform, FlinkBatchTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTransformTranslators.java:556
MethodtranslateNode
(Create.Values<OUT> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:94
MethodtranslateNode
(TextIO.Write.Bound<T> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:135
MethodtranslateNode
(Read.Unbounded<T> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:167
MethodtranslateNode
(ParDo.Bound<IN, OUT> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:192
MethodtranslateNode
(Window.Bound<T> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:217
MethodtranslateNode
(GroupByKey<K, V> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:273
MethodtranslateNode
(Combine.PerKey<K, VIN, VOUT> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:293
MethodtranslateNode
(Flatten.FlattenPCollectionList<T> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:315
MethodtranslateNode
(ParDo.BoundMulti<IN, OUT> transform, FlinkStreamingTranslationContext context)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:331
Methodtrigger
(long timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSourceWrapper.java:132
Methodvalidate
(PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:85
Methodvalidate
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:119
Methodvalidate
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:125
Methodvalidate
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:108
MethodverifyDeterministic
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:144
MethodvisitTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingPipelineTranslator.java:88
MethodvisitTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchPipelineTranslator.java:97
MethodvisitValue
(PValue value, TransformTreeNode producer)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingPipelineTranslator.java:108
MethodvisitValue
(PValue value, TransformTreeNode producer)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchPipelineTranslator.java:117
Methodwindow
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:499
Methodwindow
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:122
Methodwindow
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:110
MethodwindowingInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:505
MethodwindowingInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:142
MethodwindowingInternals
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:120
MethodwindowingInternalsHelper
(final WindowedValue<IN> inElement, final Collector<WindowedValue<OUT>> collector)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:59
MethodwindowingInternalsHelper
(WindowedValue<IN> inElement, Collector<WindowedValue<RawUnionValue>> outCollector)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundMultiWrapper.java:69
Methodwindows
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupAlsoByWindowTest.java:493
Methodwindows
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:258
Methodwindows
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:77
Methodwindows
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:527
← previousnext →701–800 of 813, ranked by callers