MCPcopy Create free account

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

Functions813 in github.com/dataArtisans/flink-dataflow

MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:50
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/SideInputITCase.java:58
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:70
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesEmptyITCase.java:40
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountITCase.java:48
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:47
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:61
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/JoinExamplesITCase.java:74
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/AvroITCase.java:45
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesITCase.java:41
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupByNullKeyTest.java:51
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:55
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/TopWikipediaSessionsITCase.java:57
MethodprocessElement
(DoFn<Void, String>.ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/MaybeEmptyTestITCase.java:55
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/ParDoMultiOutputITCase.java:71
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/SideInputITCase.java:48
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:119
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:134
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:69
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:104
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:119
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/AvroITCase.java:84
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupByNullKeyTest.java:64
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:80
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/TopWikipediaSessionsITCase.java:103
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:74
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:106
MethodprocessElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:124
MethodprocessElement
(final ProcessContext c)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:244
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:236
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:423
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:36
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:83
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:173
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:239
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:306
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:333
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:47
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:66
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/WindowedWordCount.java:54
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/WindowedWordCount.java:65
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/JoinExamples.java:79
MethodprocessElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/JoinExamples.java:113
MethodproducesSortedKeys
(PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:109
MethodputNormalizedKey
(T record, MemorySegment target, int offset, int numBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:166
MethodputNormalizedKey
(KV<K, V> record, MemorySegment target, int offset, int numBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:210
MethodreachedEnd
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:145
Methodread
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:309
Methodread
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:383
Methodread
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:564
MethodreadLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:314
MethodreadLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:388
MethodreadLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:558
MethodreadLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:649
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:52
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:65
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:54
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:64
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SinkOutputFormat.java:116
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSourceWrapper.java:164
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:67
MethodreadObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:71
MethodreadWithKeyDenormalization
(T reuse, DataInputView source)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:191
MethodreadWithKeyDenormalization
(KV<K, V> reuse, DataInputView source)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:237
Methodreduce
(Iterable<KV<K, V>> values, Collector<KV<K, Iterable<V>>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkKeyedListAggregationFunction.java:32
Methodreduce
(Iterable<KV<K, VA>> values, Collector<KV<K, VO>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkReduceFunction.java:42
MethodregisterByteSizeObserver
Notifies ElementByteSizeObserver about the byte size of the encoded value using this coder.
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:123
MethodresetLocal
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:56
MethodrestoreState
(StreamTaskState taskState, long recoveryTimestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:605
MethodrestoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:343
MethodrestoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:431
MethodrestoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:621
MethodrestoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:704
Methodrun
(SourceContext<WindowedValue<T>> ctx)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSourceWrapper.java:76
Methodserialize
(VoidValue record, DataOutputView target)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:64
MethodsetBroker
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:85
MethodsetCurrentInputWatermark
(Instant watermark)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:40
MethodsetCurrentOutputWatermark
(Instant watermark)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:45
MethodsetCurrentTransform
Sets the AppliedPTransform which carries input/output. @param currentTransform
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:74
MethodsetGroup
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:97
MethodsetInput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:101
MethodsetInput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:88
MethodsetKafkaTopic
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:79
MethodsetOutput
(String value)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:143
MethodsetOutput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:106
MethodsetOutput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:92
MethodsetRecursive
(Boolean value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:356
MethodsetReference
(T toCompare)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:64
MethodsetReference
(KV<K, V> toCompare)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:85
MethodsetStableUniqueNames
(CheckEnabled enabled)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:474
MethodsetTimer
(TimerData timerKey)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:386
MethodsetZookeeper
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:91
MethodshouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:320
MethodshouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:417
MethodshouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:593
MethodshouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:676
MethodsideInput
(PCollectionView<T> view, BoundedWindow mainInputWindow)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:92
MethodsideInput
(PCollectionView<T> view, BoundedWindow mainInputWindow)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:542
MethodsideInput
(PCollectionView<T> view)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:152
MethodsideInput
(PCollectionView<T> view)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:530
← previousnext →601–700 of 813, ranked by callers