Modifier and Type | Method and Description |
---|---|
static FlinkFnApi.UserDefinedAggregateFunction |
PythonOperatorUtils.getUserDefinedAggregateFunctionProto(PythonAggregateFunctionInfo pythonFunctionInfo,
DataViewUtils.DataViewSpec[] dataViewSpecs) |
Modifier and Type | Method and Description |
---|---|
static DataViewUtils.DataViewSpec[] |
CommonPythonUtil.extractDataViewSpecs(int index,
DataType accType) |
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 |
---|
AbstractPythonStreamAggregateOperator(Configuration config,
RowType inputType,
RowType outputType,
PythonAggregateFunctionInfo[] aggregateFunctions,
DataViewUtils.DataViewSpec[][] dataViewSpecs,
int[] grouping,
int indexOfCountStar,
boolean generateUpdateBefore,
String coderUrn,
FlinkFnApi.CoderParam.OutputMode outputMode) |
AbstractPythonStreamGroupAggregateOperator(Configuration config,
RowType inputType,
RowType outputType,
PythonAggregateFunctionInfo[] aggregateFunctions,
DataViewUtils.DataViewSpec[][] dataViewSpecs,
int[] grouping,
int indexOfCountStar,
boolean generateUpdateBefore,
long minRetentionTime,
long maxRetentionTime,
FlinkFnApi.CoderParam.OutputMode outputMode) |
PythonStreamGroupAggregateOperator(Configuration config,
RowType inputType,
RowType outputType,
PythonAggregateFunctionInfo[] aggregateFunctions,
DataViewUtils.DataViewSpec[][] dataViewSpecs,
int[] grouping,
int indexOfCountStar,
boolean countStarInserted,
boolean generateUpdateBefore,
long minRetentionTime,
long maxRetentionTime) |
PythonStreamGroupTableAggregateOperator(Configuration config,
RowType inputType,
RowType outputType,
PythonAggregateFunctionInfo[] aggregateFunctions,
DataViewUtils.DataViewSpec[][] dataViewSpecs,
int[] grouping,
int indexOfCountStar,
boolean generateUpdateBefore,
long minRetentionTime,
long maxRetentionTime) |
PythonStreamGroupWindowAggregateOperator(Configuration config,
RowType inputType,
RowType outputType,
PythonAggregateFunctionInfo[] aggregateFunctions,
DataViewUtils.DataViewSpec[][] dataViewSpecs,
int[] grouping,
int indexOfCountStar,
boolean generateUpdateBefore,
boolean countStarInserted,
int inputTimeFieldIndex,
WindowAssigner<W> windowAssigner,
org.apache.flink.table.planner.plan.logical.LogicalWindow window,
long allowedLateness,
PlannerNamedWindowProperty[] namedProperties,
java.time.ZoneId shiftTimeZone) |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.