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
createAggregatorInternal
(String name, Combine.CombineFn<AggInputT, ?, AggOutputT> combiner)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:167
Method
createComparator
(int[] logicalKeyFields, boolean[] orders, int logicalFieldOffset, ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:49
Method
createComparator
(boolean sortOrderAscending, ExecutionConfig executionConfig)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:111
Method
createInputSplits
(int numSplits)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:109
Method
createInstance
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:72
Method
createInstance
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:44
Method
createReader
(PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:114
Method
createReader
(PipelineOptions options, @Nullable CheckpointMark checkpointMark)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:114
Method
createReader
(PipelineOptions options, @Nullable C checkpointMark)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:95
Method
createSerializer
(ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:86
Method
createSerializer
(ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:69
Method
createTypeComparator
(ExecutionConfig config)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:181
Method
createTypeComparatorBuilder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:168
Method
currentInputWatermarkTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:73
Method
currentOutputWatermarkTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:85
Method
currentProcessingTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:68
Method
currentSynchronizedProcessingTime
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/AbstractFlinkTimerInternals.java:78
Method
deleteTimer
(TimerData timerKey)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:391
Method
deserialize
(DataInputView source)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:69
Method
duplicate
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:247
Method
duplicate
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeSerializer.java:67
Method
duplicate
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:39
Method
element
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:112
Method
element
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:216
Method
element
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:96
Method
element
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:100
Method
enterCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingPipelineTranslator.java:54
Method
enterCompositeTransform
(TransformTreeNode node)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchPipelineTranslator.java:59
Method
equalToReference
(T candidate)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:74
Method
equalToReference
(KV<K, V> candidate)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:95
Method
equals
(Object o)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:97
Method
equals
(Object o)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:83
Method
equals
(Object obj)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:86
Method
extractKeys
(Object record, Object[] target, int index)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:206
Method
extractKeys
(Object record, Object[] target, int index)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:252
Method
extractOutput
(Object key, int[] accumulator, CombineWithContext.Context c)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/StateSerializationTest.java:74
Method
finalize
(Iterable<String> writerResults, PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/WriteSinkITCase.java:107
Method
flatMap
(WindowedValue<T> value, Collector<String> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:152
Method
flatMap
(WindowedValue<RawUnionValue> value, Collector<WindowedValue<?>> collector)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTransformTranslators.java:373
Method
flatMap
(WindowedValue<IN> value, Collector<WindowedValue<OUTFL>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:69
Method
flatMap
(IN value, Collector<WindowedValue<OUT>> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/FlinkStreamingCreateFunction.java:45
Method
flatMap
(RawUnionValue rawUnionValue, Collector<T> collector)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputPruningFunction.java:34
Method
flatMap
(IN value, Collector<OUT> out)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkCreateFunction.java:43
Method
fromOptions
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
Method
generateInitialSplits
( int desiredNumSplits, PipelineOptions options)
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:108
Method
generateInitialSplits
(int desiredNumSplits, PipelineOptions options)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedFlinkSource.java:43
Method
generateInitialSplits
(int desiredNumSplits, PipelineOptions options)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:90
Method
getAccum
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:550
Method
getAggregatorValues
(final Aggregator<?, T> aggregator)
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerResult.java:50
Method
getArity
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:65
Method
getArity
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderTypeInformation.java:52
Method
getAverageRecordWidth
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:97
Method
getCheckpointMark
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:188
Method
getCheckpointMark
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:221
Method
getCheckpointMarkCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:119
Method
getCheckpointMarkCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedFlinkSource.java:53
Method
getCheckpointMarkCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:100
Method
getCoderArguments
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:98
Method
getCombineFn
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:87
Method
getCombineFn
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SerializableFnAggregatorWrapper.java:76
Method
getComponents
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/UnionCoder.java:103
Method
getCurrent
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:145
Method
getCurrent
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:169
Method
getCurrentRecordId
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:193
Method
getCurrentSource
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:155
Method
getCurrentSource
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:193
Method
getCurrentSource
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:226
Method
getCurrentTimestamp
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:174
Method
getDefaultOutputCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/ReadSourceITCase.java:122
Method
getDefaultOutputCoder
()
runner/src/test/java/com/dataartisans/flink/dataflow/streaming/UnboundedSourceITCase.java:129
Method
getDefaultOutputCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedFlinkSource.java:63
Method
getDefaultOutputCoder
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/io/UnboundedSocketSource.java:115
Method
getDefaultOutputCoder
()
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:166
Method
getExecutionEnvironment
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:51
Method
getFieldIndex
(String fieldName)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:150
Method
getFieldNames
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:145
Method
getFlatComparators
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:212
Method
getFlatComparators
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:260
Method
getFlatFields
(String fieldExpression, int offset, List<FlatFieldDescriptor> result)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderTypeInformation.java:162
Method
getInput
(PTransform<I, ?> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:78
Method
getInput
(PTransform<I, ?> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTranslationContext.java:120
Method
getInput
()
examples/src/main/java/com/dataartisans/flink/dataflow/examples/TFIDF.java:98
Method
getInputSplitAssigner
(final SourceInputSplit[] sourceInputSplits)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:128
Method
getInputTypeInfo
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTranslationContext.java:112
Method
getLength
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/VoidCoderTypeSerializer.java:59
Method
getLocalValue
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/CombineFnAggregatorWrapper.java:51
Method
getName
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SerializableFnAggregatorWrapper.java:71
Method
getNextInputSplit
(String host, int taskId)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:133
Method
getNormalizeKeyLen
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/CoderComparator.java:156
Method
getNormalizeKeyLen
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/types/KvCoderComperator.java:200
Method
getNumberOfRecords
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/SourceInputFormat.java:92
Method
getOutput
(PTransform<?, O> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:83
Method
getOutput
(PTransform<?, O> transform)
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkBatchTranslationContext.java:125
Method
getPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/FlinkRunnerRegistrar.java:47
Method
getPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/FlinkStreamingTranslationContext.java:55
Method
getPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkGroupAlsoByWindowWrapper.java:441
Method
getPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/FlinkAbstractParDoWrapper.java:147
Method
getPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/wrappers/streaming/state/FlinkStateInternals.java:525
Method
getPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkDoFnFunction.java:156
Method
getPipelineOptions
()
runner/src/main/java/com/dataartisans/flink/dataflow/translation/functions/FlinkMultiOutputDoFnFunction.java:125
← previous
next →
401–500 of 813, ranked by callers