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
getPipelineRunners
()
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerRegistrar.java:39
Method
getProducedType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/utils/FlinkMapperWrapper.java:36
Method
getProducedType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupByKeyWrapper.java:58
Method
getRunner
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:459
Method
getSink
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:117
Method
getSplitNumber
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputSplit.java:43
Method
getStableUniqueNames
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:469
Method
getStatistics
(BaseStatistics baseStatistics)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:80
Method
getTotalFields
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:92
Method
getTotalFields
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:78
Method
getTotalInputSize
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:86
Method
getTypeAt
(int pos)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:120
Method
getTypeClass
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:70
Method
getTypeClass
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:57
Method
getValuesAtSteps
()
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerResult.java:56
Method
getWatermark
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:183
Method
getWriteOperation
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:147
Method
getWriterResultCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:97
Method
hash
(T record)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:59
Method
hash
(KV<K, V> record)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:75
Method
hashCode
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:108
Method
hashCode
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:94
Method
hashCode
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:101
Method
hashCode
()
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:291
Method
initializeTypeComparatorBuilder
(int size)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:175
Method
invertNormalizedKey
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:196
Method
invertNormalizedKey
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:242
Method
isBasicType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:55
Method
isBasicType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:42
Method
isEmpty
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:380
Method
isEmpty
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:555
Method
isImmutableType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:62
Method
isImmutableType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:34
Method
isKeyType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:81
Method
isKeyType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:64
Method
isNormalizedKeyPrefixOnly
(int keyBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:161
Method
isNormalizedKeyPrefixOnly
(int keyBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:205
Method
isRegisterByteSizeObserverCheap
Since this coder uses elementCoders.get(index) and coders that are known to run in constant time, we defer the return value to that coder.
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:112
Method
isTupleType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:60
Method
isTupleType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:47
Method
jsonOf
( @JsonProperty(PropertyNames.COMPONENT_ENCODINGS) List<Coder<?>> elements)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:51
Method
leaveCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingPipelineTranslator.java:69
Method
leaveCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchPipelineTranslator.java:78
Method
main
(String[] args)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:146
Method
main
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:437
Method
main
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:95
Method
main
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:359
Method
main
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:101
Method
main
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/WindowedWordCount.java:97
Method
main
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/JoinExamples.java:125
Method
mapPartition
(Iterable<Tuple2<KI, VI>> values, Collector<Tuple2<KO, VO>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/utils/FlinkMapperWrapper.java:41
Method
mapPartition
(Iterable<IN> values, Collector<OUT> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:74
Method
mapPartition
(Iterable<IN> values, Collector<RawUnionValue> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:79
Method
merge
(Accumulator<AI, Serializable> other)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:61
Method
merge
(Accumulator<AI, Serializable> other)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SerializableFnAggregatorWrapper.java:60
Method
mergeAccumulators
(Object key, Iterable<int[]> accumulators, CombineWithContext.Context c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/StateSerializationTest.java:65
Method
mergeAccumulators
(Iterable<AccumT> accumulators)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:577
Method
named
Returns a new ConsoleIO.Write PTransform that's like this one but with the given step name. Does not modify this object.
runner/src/main/java/com/dataartisans/flink/dataflow/io/ConsoleIO.java:69
Method
nextRecord
(T t)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:150
Method
open
(SourceInputSplit<T> sourceInputSplit)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:74
Method
open
(int taskNumber, int numTasks)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SinkOutputFormat.java:78
Method
open
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:217
Method
output
(KV<K, VOUT> output)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:482
Method
output
(OUTDF output)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:157
Method
outputWindowedValue
(OUT output, Instant timestamp, Collection<? extends BoundedWindow> windows, PaneInfo pane)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:67
Method
outputWindowedValue
(KV<K, VOUT> output, Instant timestamp, Collection<? extends BoundedWindow> windows, PaneInfo pane)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:514
Method
outputWithTimestamp
(KV<K, VOUT> output, Instant timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:488
Method
outputWithTimestamp
(OUT output, Instant timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:147
Method
outputWithTimestampHelper
(WindowedValue<IN> inElement, OUT output, Instant timestamp, Collector<WindowedValue<OUT>> collector)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:42
Method
outputWithTimestampHelper
(WindowedValue<IN> inElement, OUT output, Instant timestamp, Collector<WindowedValue<RawUnionValue>> collector
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundMultiWrapper.java:45
Method
pane
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:82
Method
pane
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:494
Method
pane
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:137
Method
partitionFor
(KV<String, List<CompletionCandidate>> elem, int numPartitions)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:163
Method
persistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:140
Method
persistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:325
Method
persistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:422
Method
persistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:598
Method
persistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:681
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/MaybeEmptyTestITCase.java:42
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/TfIdfITCase.java:48
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/FlattenizeITCase.java:42
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ParDoMultiOutputITCase.java:43
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:55
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/SideInputITCase.java:63
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:75
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesEmptyITCase.java:45
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountITCase.java:53
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:52
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:66
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/JoinExamplesITCase.java:79
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/AvroITCase.java:52
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesITCase.java:46
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupByNullKeyTest.java:56
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:60
Method
postSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/TopWikipediaSessionsITCase.java:62
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/MaybeEmptyTestITCase.java:37
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/TfIdfITCase.java:43
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/FlattenizeITCase.java:36
Method
preSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ParDoMultiOutputITCase.java:38
← previous
next →
501–600 of 813, ranked by callers