Modifier and Type | Class and Description |
---|---|
class |
RichAggregateFunction<IN,ACC,OUT>
Rich variant of the
AggregateFunction . |
Modifier and Type | Method and Description |
---|---|
AggregateFunction<IN,ACC,OUT> |
AggregatingStateDescriptor.getAggregateFunction()
Returns the aggregate function to be used for the state.
|
Constructor and Description |
---|
AggregatingStateDescriptor(String name,
AggregateFunction<IN,ACC,OUT> aggFunction,
Class<ACC> stateType)
Creates a new state descriptor with the given name, function, and type.
|
AggregatingStateDescriptor(String name,
AggregateFunction<IN,ACC,OUT> aggFunction,
TypeInformation<ACC> stateType)
Creates a new
ReducingStateDescriptor with the given name and default value. |
AggregatingStateDescriptor(String name,
AggregateFunction<IN,ACC,OUT> aggFunction,
TypeSerializer<ACC> typeSerializer)
Creates a new
ValueStateDescriptor with the given name and default value. |
Modifier and Type | Method and Description |
---|---|
static <IN,ACC> TypeInformation<ACC> |
TypeExtractor.getAggregateFunctionAccumulatorType(AggregateFunction<IN,ACC,?> function,
TypeInformation<IN> inType,
String functionName,
boolean allowMissing) |
static <IN,OUT> TypeInformation<OUT> |
TypeExtractor.getAggregateFunctionReturnType(AggregateFunction<IN,?,OUT> function,
TypeInformation<IN> inType,
String functionName,
boolean allowMissing) |
Modifier and Type | Method and Description |
---|---|
<ACC,R> SingleOutputStreamOperator<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,R> function)
Applies the given aggregation function to each window.
|
<ACC,R> SingleOutputStreamOperator<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,R> function)
Applies the given
AggregateFunction to each window. |
<ACC,R> SingleOutputStreamOperator<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,R> function,
TypeInformation<ACC> accumulatorType,
TypeInformation<R> resultType)
Applies the given aggregation function to each window.
|
<ACC,R> SingleOutputStreamOperator<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,R> function,
TypeInformation<ACC> accumulatorType,
TypeInformation<R> resultType)
Applies the given
AggregateFunction to each window. |
<ACC,V,R> SingleOutputStreamOperator<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,V> aggFunction,
AllWindowFunction<V,R,W> windowFunction)
Applies the given window function to each window.
|
<ACC,V,R> SingleOutputStreamOperator<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,V> aggregateFunction,
AllWindowFunction<V,R,W> windowFunction,
TypeInformation<ACC> accumulatorType,
TypeInformation<V> aggregateResultType,
TypeInformation<R> resultType)
Applies the given window function to each window.
|
<ACC,V,R> SingleOutputStreamOperator<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,V> aggFunction,
ProcessAllWindowFunction<V,R,W> windowFunction)
Applies the given window function to each window.
|
<ACC,V,R> SingleOutputStreamOperator<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,V> aggregateFunction,
ProcessAllWindowFunction<V,R,W> windowFunction,
TypeInformation<ACC> accumulatorType,
TypeInformation<V> aggregateResultType,
TypeInformation<R> resultType)
Applies the given window function to each window.
|
<ACC,V,R> SingleOutputStreamOperator<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,V> aggFunction,
ProcessWindowFunction<V,R,K,W> windowFunction)
Applies the given window function to each window.
|
<ACC,V,R> SingleOutputStreamOperator<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,V> aggregateFunction,
ProcessWindowFunction<V,R,K,W> windowFunction,
TypeInformation<ACC> accumulatorType,
TypeInformation<V> aggregateResultType,
TypeInformation<R> resultType)
Applies the given window function to each window.
|
<ACC,V,R> SingleOutputStreamOperator<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,V> aggFunction,
WindowFunction<V,R,K,W> windowFunction)
Applies the given window function to each window.
|
<ACC,V,R> SingleOutputStreamOperator<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,V> aggregateFunction,
WindowFunction<V,R,K,W> windowFunction,
TypeInformation<ACC> accumulatorType,
TypeInformation<V> aggregateResultType,
TypeInformation<R> resultType)
Applies the given window function to each window.
|
Constructor and Description |
---|
AggregateApplyAllWindowFunction(AggregateFunction<T,ACC,V> aggFunction,
AllWindowFunction<V,R,W> windowFunction) |
AggregateApplyWindowFunction(AggregateFunction<T,ACC,V> aggFunction,
WindowFunction<V,R,K,W> windowFunction) |
Modifier and Type | Method and Description |
---|---|
<ACC,R> DataStream<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,R> aggregateFunction,
TypeInformation<ACC> evidence$5,
TypeInformation<R> evidence$6)
Applies the given aggregation function to each window.
|
<ACC,R> DataStream<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,R> aggregateFunction,
TypeInformation<ACC> evidence$5,
TypeInformation<R> evidence$6)
Applies the given aggregation function to each window and key.
|
<ACC,V,R> DataStream<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,V> preAggregator,
AllWindowFunction<V,R,W> windowFunction,
TypeInformation<ACC> evidence$7,
TypeInformation<V> evidence$8,
TypeInformation<R> evidence$9)
Applies the given window function to each window.
|
<ACC,V,R> DataStream<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,V> preAggregator,
scala.Function3<W,scala.collection.Iterable<V>,Collector<R>,scala.runtime.BoxedUnit> windowFunction,
TypeInformation<ACC> evidence$13,
TypeInformation<V> evidence$14,
TypeInformation<R> evidence$15)
Applies the given window function to each window.
|
<ACC,V,R> DataStream<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,V> preAggregator,
scala.Function4<K,W,scala.collection.Iterable<V>,Collector<R>,scala.runtime.BoxedUnit> windowFunction,
TypeInformation<ACC> evidence$10,
TypeInformation<V> evidence$11,
TypeInformation<R> evidence$12)
Applies the given window function to each window.
|
<ACC,V,R> DataStream<R> |
AllWindowedStream.aggregate(AggregateFunction<T,ACC,V> preAggregator,
ProcessAllWindowFunction<V,R,W> windowFunction,
TypeInformation<ACC> evidence$10,
TypeInformation<V> evidence$11,
TypeInformation<R> evidence$12)
Applies the given window function to each window.
|
<ACC,V,R> DataStream<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,V> preAggregator,
ProcessWindowFunction<V,R,K,W> windowFunction,
TypeInformation<ACC> evidence$13,
TypeInformation<V> evidence$14,
TypeInformation<R> evidence$15)
Applies the given window function to each window.
|
<ACC,V,R> DataStream<R> |
WindowedStream.aggregate(AggregateFunction<T,ACC,V> preAggregator,
WindowFunction<V,R,K,W> windowFunction,
TypeInformation<ACC> evidence$7,
TypeInformation<V> evidence$8,
TypeInformation<R> evidence$9)
Applies the given window function to each window.
|
Constructor and Description |
---|
InternalAggregateProcessAllWindowFunction(AggregateFunction<T,ACC,V> aggFunction,
ProcessAllWindowFunction<V,R,W> windowFunction) |
InternalAggregateProcessWindowFunction(AggregateFunction<T,ACC,V> aggFunction,
ProcessWindowFunction<V,R,K,W> windowFunction) |
Modifier and Type | Class and Description |
---|---|
class |
AggregateAggFunction
Aggregate Function used for the aggregate operator in
WindowedStream |
Modifier and Type | Method and Description |
---|---|
static scala.Tuple3<AggregateFunction<CRow,Row,Row>,RowTypeInfo,RowTypeInfo> |
AggregateUtil.createDataStreamAggregateFunction(CodeGenerator generator,
scala.collection.Seq<org.apache.calcite.util.Pair<org.apache.calcite.rel.core.AggregateCall,String>> namedAggregates,
org.apache.calcite.rel.type.RelDataType inputType,
scala.collection.Seq<TypeInformation<?>> inputFieldTypeInfo,
org.apache.calcite.rel.type.RelDataType outputType,
int[] groupingKeys,
boolean needMerge) |
scala.Tuple3<AggregateFunction<CRow,Row,Row>,RowTypeInfo,RowTypeInfo> |
AggregateUtil$.createDataStreamAggregateFunction(CodeGenerator generator,
scala.collection.Seq<org.apache.calcite.util.Pair<org.apache.calcite.rel.core.AggregateCall,String>> namedAggregates,
org.apache.calcite.rel.type.RelDataType inputType,
scala.collection.Seq<TypeInformation<?>> inputFieldTypeInfo,
org.apache.calcite.rel.type.RelDataType outputType,
int[] groupingKeys,
boolean needMerge) |
Copyright © 2014–2018 The Apache Software Foundation. All rights reserved.