Package | Description |
---|---|
org.apache.flink.streaming.api.datastream | |
org.apache.flink.streaming.api.functions.windowing | |
org.apache.flink.streaming.api.scala.function.util | |
org.apache.flink.streaming.examples.windowing | |
org.apache.flink.streaming.runtime.operators.windowing |
This package contains the operators that implement the various window operations
on data streams.
|
org.apache.flink.streaming.runtime.operators.windowing.functions | |
org.apache.flink.table.runtime.aggregate |
Modifier and Type | Method and Description |
---|---|
<R> SingleOutputStreamOperator<R> |
WindowedStream.apply(ReduceFunction<T> reduceFunction,
WindowFunction<T,R,K,W> function)
Deprecated.
|
<R> SingleOutputStreamOperator<R> |
WindowedStream.apply(ReduceFunction<T> reduceFunction,
WindowFunction<T,R,K,W> function,
TypeInformation<R> resultType)
Deprecated.
|
<R> SingleOutputStreamOperator<R> |
WindowedStream.apply(R initialValue,
FoldFunction<T,R> foldFunction,
WindowFunction<R,R,K,W> function)
Deprecated.
Use
#fold(R, FoldFunction, WindowFunction) instead. |
<R> SingleOutputStreamOperator<R> |
WindowedStream.apply(R initialValue,
FoldFunction<T,R> foldFunction,
WindowFunction<R,R,K,W> function,
TypeInformation<R> resultType)
Deprecated.
Use
#fold(R, FoldFunction, WindowFunction, TypeInformation, TypeInformation) instead. |
<R> SingleOutputStreamOperator<R> |
WindowedStream.apply(WindowFunction<T,R,K,W> function)
Applies the given window function to each window.
|
<R> SingleOutputStreamOperator<R> |
WindowedStream.apply(WindowFunction<T,R,K,W> function,
TypeInformation<R> resultType)
Applies the given window function to each window.
|
<ACC,R> SingleOutputStreamOperator<R> |
WindowedStream.fold(ACC initialValue,
FoldFunction<T,ACC> foldFunction,
WindowFunction<ACC,R,K,W> function)
Applies the given window function to each window.
|
<ACC,R> SingleOutputStreamOperator<R> |
WindowedStream.fold(ACC initialValue,
FoldFunction<T,ACC> foldFunction,
WindowFunction<ACC,R,K,W> function,
TypeInformation<ACC> foldAccumulatorType,
TypeInformation<R> resultType)
Applies the given window function to each window.
|
<R> SingleOutputStreamOperator<R> |
WindowedStream.reduce(ReduceFunction<T> reduceFunction,
WindowFunction<T,R,K,W> function)
Applies the given window function to each window.
|
<R> SingleOutputStreamOperator<R> |
WindowedStream.reduce(ReduceFunction<T> reduceFunction,
WindowFunction<T,R,K,W> function,
TypeInformation<R> resultType)
Applies the given window function to each window.
|
Modifier and Type | Class and Description |
---|---|
class |
FoldApplyWindowFunction<K,W extends Window,T,ACC,R> |
class |
PassThroughWindowFunction<K,W extends Window,T> |
class |
ReduceApplyWindowFunction<K,W extends Window,T,R> |
class |
ReduceIterableWindowFunction<K,W extends Window,T> |
class |
RichWindowFunction<IN,OUT,KEY,W extends Window>
Rich variant of the
WindowFunction . |
Constructor and Description |
---|
FoldApplyWindowFunction(ACC initialValue,
FoldFunction<T,ACC> foldFunction,
WindowFunction<ACC,R,K,W> windowFunction,
TypeInformation<ACC> accTypeInformation) |
ReduceApplyWindowFunction(ReduceFunction<T> reduceFunction,
WindowFunction<T,R,K,W> windowFunction) |
Modifier and Type | Class and Description |
---|---|
class |
ScalaWindowFunction<IN,OUT,KEY,W extends Window>
A wrapper function that exposes a Scala Function4 as a Java WindowFunction.
|
class |
ScalaWindowFunctionWrapper<IN,OUT,KEY,W extends Window>
A wrapper function that exposes a Scala WindowFunction as a JavaWindow function.
|
Modifier and Type | Class and Description |
---|---|
static class |
GroupedProcessingTimeWindowExample.SummingWindowFunction |
Constructor and Description |
---|
AccumulatingKeyedTimePanes(KeySelector<Type,Key> keySelector,
WindowFunction<Type,Result,Key,Window> function) |
AccumulatingProcessingTimeWindowOperator(WindowFunction<IN,OUT,KEY,TimeWindow> function,
KeySelector<IN,KEY> keySelector,
TypeSerializer<KEY> keySerializer,
TypeSerializer<IN> valueSerializer,
long windowLength,
long windowSlide)
Deprecated.
|
Constructor and Description |
---|
InternalIterableWindowFunction(WindowFunction<IN,OUT,KEY,W> wrappedFunction) |
InternalSingleValueWindowFunction(WindowFunction<IN,OUT,KEY,W> wrappedFunction) |
Modifier and Type | Class and Description |
---|---|
class |
AggregateTimeWindowFunction |
class |
AggregateWindowFunction<W extends Window> |
class |
IncrementalAggregateTimeWindowFunction
Computes the final aggregate value from incrementally computed aggreagtes.
|
class |
IncrementalAggregateWindowFunction<W extends Window>
Computes the final aggregate value from incrementally computed aggreagtes.
|
Modifier and Type | Method and Description |
---|---|
static WindowFunction<Row,Row,Tuple,Window> |
AggregateUtil.createWindowAggregationFunction(LogicalWindow window,
scala.collection.Seq<org.apache.calcite.util.Pair<org.apache.calcite.rel.core.AggregateCall,String>> namedAggregates,
org.apache.calcite.rel.type.RelDataType inputType,
org.apache.calcite.rel.type.RelDataType outputType,
int[] groupings,
scala.collection.Seq<FlinkRelBuilder.NamedWindowProperty> properties)
Create a
WindowFunction to compute partitioned group window aggregates. |
WindowFunction<Row,Row,Tuple,Window> |
AggregateUtil$.createWindowAggregationFunction(LogicalWindow window,
scala.collection.Seq<org.apache.calcite.util.Pair<org.apache.calcite.rel.core.AggregateCall,String>> namedAggregates,
org.apache.calcite.rel.type.RelDataType inputType,
org.apache.calcite.rel.type.RelDataType outputType,
int[] groupings,
scala.collection.Seq<FlinkRelBuilder.NamedWindowProperty> properties)
Create a
WindowFunction to compute partitioned group window aggregates. |
static WindowFunction<Row,Row,Tuple,Window> |
AggregateUtil.createWindowIncrementalAggregationFunction(LogicalWindow window,
scala.collection.Seq<org.apache.calcite.util.Pair<org.apache.calcite.rel.core.AggregateCall,String>> namedAggregates,
org.apache.calcite.rel.type.RelDataType inputType,
org.apache.calcite.rel.type.RelDataType outputType,
int[] groupings,
scala.collection.Seq<FlinkRelBuilder.NamedWindowProperty> properties)
Create a
WindowFunction to finalize incrementally pre-computed window aggregates. |
WindowFunction<Row,Row,Tuple,Window> |
AggregateUtil$.createWindowIncrementalAggregationFunction(LogicalWindow window,
scala.collection.Seq<org.apache.calcite.util.Pair<org.apache.calcite.rel.core.AggregateCall,String>> namedAggregates,
org.apache.calcite.rel.type.RelDataType inputType,
org.apache.calcite.rel.type.RelDataType outputType,
int[] groupings,
scala.collection.Seq<FlinkRelBuilder.NamedWindowProperty> properties)
Create a
WindowFunction to finalize incrementally pre-computed window aggregates. |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.