MCPcopy Create free account

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

Functions813 in github.com/dataArtisans/flink-dataflow

MethodcreateAggregatorInternal
(String name, Combine.CombineFn<AggInputT, ?, AggOutputT> combiner)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:167
MethodcreateComparator
(int[] logicalKeyFields, boolean[] orders, int logicalFieldOffset, ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:49
MethodcreateComparator
(boolean sortOrderAscending, ExecutionConfig executionConfig)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:111
MethodcreateInputSplits
(int numSplits)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:109
MethodcreateInstance
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:72
MethodcreateInstance
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:44
MethodcreateReader
(PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:114
MethodcreateReader
(PipelineOptions options, @Nullable CheckpointMark checkpointMark)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:114
MethodcreateReader
(PipelineOptions options, @Nullable C checkpointMark)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:95
MethodcreateSerializer
(ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:86
MethodcreateSerializer
(ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:69
MethodcreateTypeComparator
(ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:181
MethodcreateTypeComparatorBuilder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:168
MethodcurrentInputWatermarkTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:73
MethodcurrentOutputWatermarkTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:85
MethodcurrentProcessingTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:68
MethodcurrentSynchronizedProcessingTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:78
MethoddeleteTimer
(TimerData timerKey)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:391
Methoddeserialize
(DataInputView source)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:69
Methodduplicate
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:247
Methodduplicate
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:67
Methodduplicate
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:39
Methodelement
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:112
Methodelement
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:216
Methodelement
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:96
Methodelement
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:100
MethodenterCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingPipelineTranslator.java:54
MethodenterCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchPipelineTranslator.java:59
MethodequalToReference
(T candidate)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:74
MethodequalToReference
(KV<K, V> candidate)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:95
Methodequals
(Object o)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:97
Methodequals
(Object o)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:83
Methodequals
(Object obj)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:86
MethodextractKeys
(Object record, Object[] target, int index)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:206
MethodextractKeys
(Object record, Object[] target, int index)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:252
MethodextractOutput
(Object key, int[] accumulator, CombineWithContext.Context c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/StateSerializationTest.java:74
Methodfinalize
(Iterable<String> writerResults, PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:107
MethodflatMap
(WindowedValue<T> value, Collector<String> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:152
MethodflatMap
(WindowedValue<RawUnionValue> value, Collector<WindowedValue<?>> collector)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:373
MethodflatMap
(WindowedValue<IN> value, Collector<WindowedValue<OUTFL>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:69
MethodflatMap
(IN value, Collector<WindowedValue<OUT>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/FlinkStreamingCreateFunction.java:45
MethodflatMap
(RawUnionValue rawUnionValue, Collector<T> collector)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputPruningFunction.java:34
MethodflatMap
(IN value, Collector<OUT> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkCreateFunction.java:43
MethodfromOptions
Construct a runner from the provided options. @param options Properties which configure the runner. @return The newly created runner.
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkPipelineRunner.java:65
MethodgenerateInitialSplits
( int desiredNumSplits, PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:108
MethodgenerateInitialSplits
(int desiredNumSplits, PipelineOptions options)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedFlinkSource.java:43
MethodgenerateInitialSplits
(int desiredNumSplits, PipelineOptions options)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:90
MethodgetAccum
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:550
MethodgetAggregatorValues
(final Aggregator<?, T> aggregator)
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerResult.java:50
MethodgetArity
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:65
MethodgetArity
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:52
MethodgetAverageRecordWidth
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:97
MethodgetCheckpointMark
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:188
MethodgetCheckpointMark
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:221
MethodgetCheckpointMarkCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:119
MethodgetCheckpointMarkCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedFlinkSource.java:53
MethodgetCheckpointMarkCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:100
MethodgetCoderArguments
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:98
MethodgetCombineFn
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:87
MethodgetCombineFn
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SerializableFnAggregatorWrapper.java:76
MethodgetComponents
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:103
MethodgetCurrent
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:145
MethodgetCurrent
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:169
MethodgetCurrentRecordId
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:193
MethodgetCurrentSource
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:155
MethodgetCurrentSource
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:193
MethodgetCurrentSource
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:226
MethodgetCurrentTimestamp
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:174
MethodgetDefaultOutputCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:122
MethodgetDefaultOutputCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:129
MethodgetDefaultOutputCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedFlinkSource.java:63
MethodgetDefaultOutputCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:115
MethodgetDefaultOutputCoder
()
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:166
MethodgetExecutionEnvironment
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:51
MethodgetFieldIndex
(String fieldName)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:150
MethodgetFieldNames
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:145
MethodgetFlatComparators
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:212
MethodgetFlatComparators
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:260
MethodgetFlatFields
(String fieldExpression, int offset, List<FlatFieldDescriptor> result)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:162
MethodgetInput
(PTransform<I, ?> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:78
MethodgetInput
(PTransform<I, ?> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTranslationContext.java:120
MethodgetInput
()
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:98
MethodgetInputSplitAssigner
(final SourceInputSplit[] sourceInputSplits)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:128
MethodgetInputTypeInfo
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTranslationContext.java:112
MethodgetLength
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:59
MethodgetLocalValue
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:51
MethodgetName
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SerializableFnAggregatorWrapper.java:71
MethodgetNextInputSplit
(String host, int taskId)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:133
MethodgetNormalizeKeyLen
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:156
MethodgetNormalizeKeyLen
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:200
MethodgetNumberOfRecords
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:92
MethodgetOutput
(PTransform<?, O> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:83
MethodgetOutput
(PTransform<?, O> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTranslationContext.java:125
MethodgetPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerRegistrar.java:47
MethodgetPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:55
MethodgetPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:441
MethodgetPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:147
MethodgetPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:525
MethodgetPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:156
MethodgetPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:125
← previousnext →401–500 of 813, ranked by callers