Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/dataArtisans/flink-dataflow
/ functions
Functions
813 in github.com/dataArtisans/flink-dataflow
⨍
Functions
813
◇
Types & classes
169
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:50
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/SideInputITCase.java:58
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:70
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesEmptyITCase.java:40
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountITCase.java:48
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:47
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:61
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/JoinExamplesITCase.java:74
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/AvroITCase.java:45
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesITCase.java:41
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupByNullKeyTest.java:51
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:55
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/TopWikipediaSessionsITCase.java:57
Method
processElement
(DoFn<Void, String>.ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/MaybeEmptyTestITCase.java:55
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/ParDoMultiOutputITCase.java:71
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/SideInputITCase.java:48
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:119
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:134
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:69
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:104
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:119
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/AvroITCase.java:84
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupByNullKeyTest.java:64
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:80
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/TopWikipediaSessionsITCase.java:103
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:74
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:106
Method
processElement
(ProcessContext c)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:124
Method
processElement
(final ProcessContext c)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:244
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:236
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:423
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:36
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:83
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:173
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:239
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:306
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:333
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:47
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:66
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/WindowedWordCount.java:54
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/WindowedWordCount.java:65
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/JoinExamples.java:79
Method
processElement
(ProcessContext c)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/JoinExamples.java:113
Method
producesSortedKeys
(PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:109
Method
putNormalizedKey
(T record, MemorySegment target, int offset, int numBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:166
Method
putNormalizedKey
(KV<K, V> record, MemorySegment target, int offset, int numBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:210
Method
reachedEnd
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:145
Method
read
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:309
Method
read
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:383
Method
read
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:564
Method
readLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:314
Method
readLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:388
Method
readLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:558
Method
readLater
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:649
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:52
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:65
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:54
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:64
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SinkOutputFormat.java:116
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSourceWrapper.java:164
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:67
Method
readObject
(ObjectInputStream in)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:71
Method
readWithKeyDenormalization
(T reuse, DataInputView source)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:191
Method
readWithKeyDenormalization
(KV<K, V> reuse, DataInputView source)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:237
Method
reduce
(Iterable<KV<K, V>> values, Collector<KV<K, Iterable<V>>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkKeyedListAggregationFunction.java:32
Method
reduce
(Iterable<KV<K, VA>> values, Collector<KV<K, VO>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkReduceFunction.java:42
Method
registerByteSizeObserver
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
Method
resetLocal
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:56
Method
restoreState
(StreamTaskState taskState, long recoveryTimestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:605
Method
restoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:343
Method
restoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:431
Method
restoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:621
Method
restoreState
(StateCheckpointReader checkpointReader)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:704
Method
run
(SourceContext<WindowedValue<T>> ctx)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSourceWrapper.java:76
Method
serialize
(VoidValue record, DataOutputView target)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:64
Method
setBroker
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:85
Method
setCurrentInputWatermark
(Instant watermark)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:40
Method
setCurrentOutputWatermark
(Instant watermark)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:45
Method
setCurrentTransform
Sets the AppliedPTransform which carries input/output. @param currentTransform
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:74
Method
setGroup
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:97
Method
setInput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:101
Method
setInput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:88
Method
setKafkaTopic
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:79
Method
setOutput
(String value)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:143
Method
setOutput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:106
Method
setOutput
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:92
Method
setRecursive
(Boolean value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:356
Method
setReference
(T toCompare)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:64
Method
setReference
(KV<K, V> toCompare)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:85
Method
setStableUniqueNames
(CheckEnabled enabled)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:474
Method
setTimer
(TimerData timerKey)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:386
Method
setZookeeper
(String value)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:91
Method
shouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:320
Method
shouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:417
Method
shouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:593
Method
shouldPersist
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:676
Method
sideInput
(PCollectionView<T> view, BoundedWindow mainInputWindow)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:92
Method
sideInput
(PCollectionView<T> view, BoundedWindow mainInputWindow)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:542
Method
sideInput
(PCollectionView<T> view)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:152
Method
sideInput
(PCollectionView<T> view)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:530
← previous
next →
601–700 of 813, ranked by callers