Package | Description |
---|---|
org.apache.flink.streaming.api.utils | |
org.apache.flink.table.planner.typeutils | |
org.apache.flink.table.runtime.operators.python.aggregate |
Modifier and Type | Method and Description |
---|---|
static FlinkFnApi.UserDefinedAggregateFunction |
PythonOperatorUtils.getUserDefinedAggregateFunctionProto(PythonAggregateFunctionInfo pythonFunctionInfo,
DataViewUtils.DataViewSpec[] dataViewSpecs) |
Modifier and Type | Class and Description |
---|---|
static class |
DataViewUtils.DistinctViewSpec
Specification for a special
MapView for deduplication. |
static class |
DataViewUtils.ListViewSpec
Specification for a
ListView . |
static class |
DataViewUtils.MapViewSpec
Specification for a
MapView . |
Modifier and Type | Method and Description |
---|---|
static List<DataViewUtils.DataViewSpec> |
DataViewUtils.extractDataViews(int aggIndex,
DataType accumulatorDataType)
Searches for data views in the data type of an accumulator and extracts them.
|
Constructor and Description |
---|
PythonStreamGroupAggregateOperator(Configuration config,
RowType inputType,
RowType outputType,
PythonAggregateFunctionInfo[] aggregateFunctions,
DataViewUtils.DataViewSpec[][] dataViewSpecs,
int[] grouping,
int indexOfCountStar,
boolean countStarInserted,
boolean generateUpdateBefore,
long minRetentionTime,
long maxRetentionTime) |
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.