MCPcopy Create free account

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

Functions813 in github.com/dataArtisans/flink-dataflow

MethodgetPipelineRunners
()
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerRegistrar.java:39
MethodgetProducedType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/utils/FlinkMapperWrapper.java:36
MethodgetProducedType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupByKeyWrapper.java:58
MethodgetRunner
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:459
MethodgetSink
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:117
MethodgetSplitNumber
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputSplit.java:43
MethodgetStableUniqueNames
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:469
MethodgetStatistics
(BaseStatistics baseStatistics)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:80
MethodgetTotalFields
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:92
MethodgetTotalFields
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:78
MethodgetTotalInputSize
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:86
MethodgetTypeAt
(int pos)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:120
MethodgetTypeClass
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:70
MethodgetTypeClass
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:57
MethodgetValuesAtSteps
()
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerResult.java:56
MethodgetWatermark
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:183
MethodgetWriteOperation
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:147
MethodgetWriterResultCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:97
Methodhash
(T record)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:59
Methodhash
(KV<K, V> record)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:75
MethodhashCode
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:108
MethodhashCode
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:94
MethodhashCode
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:101
MethodhashCode
()
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:291
MethodinitializeTypeComparatorBuilder
(int size)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:175
MethodinvertNormalizedKey
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:196
MethodinvertNormalizedKey
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:242
MethodisBasicType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:55
MethodisBasicType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:42
MethodisEmpty
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:380
MethodisEmpty
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:555
MethodisImmutableType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:62
MethodisImmutableType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:34
MethodisKeyType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:81
MethodisKeyType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:64
MethodisNormalizedKeyPrefixOnly
(int keyBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:161
MethodisNormalizedKeyPrefixOnly
(int keyBytes)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:205
MethodisRegisterByteSizeObserverCheap
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
MethodisTupleType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:60
MethodisTupleType
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:47
MethodjsonOf
( @JsonProperty(PropertyNames.COMPONENT_ENCODINGS) List<Coder<?>> elements)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:51
MethodleaveCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingPipelineTranslator.java:69
MethodleaveCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchPipelineTranslator.java:78
Methodmain
(String[] args)
runner/src/test/java/com/dataartisans/flink/dataflow/util/JoinExamples.java:146
Methodmain
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:437
Methodmain
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/WordCount.java:95
Methodmain
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:359
Methodmain
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/KafkaWindowedWordCountExample.java:101
Methodmain
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/WindowedWordCount.java:97
Methodmain
(String[] args)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/JoinExamples.java:125
MethodmapPartition
(Iterable<Tuple2<KI, VI>> values, Collector<Tuple2<KO, VO>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/utils/FlinkMapperWrapper.java:41
MethodmapPartition
(Iterable<IN> values, Collector<OUT> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:74
MethodmapPartition
(Iterable<IN> values, Collector<RawUnionValue> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:79
Methodmerge
(Accumulator<AI, Serializable> other)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:61
Methodmerge
(Accumulator<AI, Serializable> other)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SerializableFnAggregatorWrapper.java:60
MethodmergeAccumulators
(Object key, Iterable<int[]> accumulators, CombineWithContext.Context c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/StateSerializationTest.java:65
MethodmergeAccumulators
(Iterable<AccumT> accumulators)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:577
Methodnamed
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
MethodnextRecord
(T t)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:150
Methodopen
(SourceInputSplit<T> sourceInputSplit)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:74
Methodopen
(int taskNumber, int numTasks)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SinkOutputFormat.java:78
Methodopen
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:217
Methodoutput
(KV<K, VOUT> output)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:482
Methodoutput
(OUTDF output)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:157
MethodoutputWindowedValue
(OUT output, Instant timestamp, Collection<? extends BoundedWindow> windows, PaneInfo pane)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:67
MethodoutputWindowedValue
(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
MethodoutputWithTimestamp
(KV<K, VOUT> output, Instant timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:488
MethodoutputWithTimestamp
(OUT output, Instant timestamp)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:147
MethodoutputWithTimestampHelper
(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
MethodoutputWithTimestampHelper
(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
Methodpane
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkParDoBoundWrapper.java:82
Methodpane
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:494
Methodpane
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:137
MethodpartitionFor
(KV<String, List<CompletionCandidate>> elem, int numPartitions)
examples/src/main/java/com/dataartisans/flink/dataflow/examples/streaming/AutoComplete.java:163
MethodpersistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:140
MethodpersistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:325
MethodpersistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:422
MethodpersistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:598
MethodpersistState
(StateCheckpointWriter checkpointBuilder)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:681
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/MaybeEmptyTestITCase.java:42
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/TfIdfITCase.java:48
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/FlattenizeITCase.java:42
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ParDoMultiOutputITCase.java:43
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:55
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/SideInputITCase.java:63
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin3ITCase.java:75
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesEmptyITCase.java:45
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountITCase.java:53
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:52
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/WordCountJoin2ITCase.java:66
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/JoinExamplesITCase.java:79
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/AvroITCase.java:52
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/RemoveDuplicatesITCase.java:46
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/GroupByNullKeyTest.java:56
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:60
MethodpostSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/TopWikipediaSessionsITCase.java:62
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/MaybeEmptyTestITCase.java:37
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/TfIdfITCase.java:43
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/FlattenizeITCase.java:36
MethodpreSubmit
()
runner/src/test/java/com/dataartisans/flink/dataflow/ParDoMultiOutputITCase.java:38
← previousnext →501–600 of 813, ranked by callers