This is an automated email from the ASF dual-hosted git repository.

xiangfu0 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git


The following commit(s) were added to refs/heads/master by this push:
     new 8bfe71bbf47 Move polymorphic aggregate construction into per-family 
providers (#19646)
8bfe71bbf47 is described below

commit 8bfe71bbf472ddf178914bc933d912b6f748a469
Author: Xiang Fu <[email protected]>
AuthorDate: Fri Sep 25 20:32:50 2026 -0700

    Move polymorphic aggregate construction into per-family providers (#19646)
    
    MODE, ANY_VALUE, FIRST/LAST_WITH_TIME, ARRAY_AGG and the ExprMin/Max
    parent/child aggregates are now built by AggregationFunctionProvider
    implementations, discovered once through ServiceLoader. Each provider owns 
its
    family's argument validation and type dispatch. The factory consults the
    registry before its switch. Construction, arguments and error messages are
    unchanged.
---
 .../function/AggregationFunctionFactory.java       | 151 +--------------------
 ...ction.java => AggregationFunctionProvider.java} |  21 +--
 .../AggregationFunctionProviderRegistry.java       |  58 ++++++++
 .../function/AnyValueAggregationFunction.java      |  14 ++
 .../ChildExprMinMaxAggregationFunction.java        |  27 ++++
 ...rstLastWithTimeAggregationFunctionProvider.java | 114 ++++++++++++++++
 .../function/ModeAggregationFunction.java          |  14 ++
 .../ParentExprMinMaxAggregationFunction.java       |  27 ++++
 .../function/array/ArrayAggFunctionProvider.java   | 107 +++++++++++++++
 ...ggregation.function.AggregationFunctionProvider |  26 ++++
 .../AggregationFunctionProviderRegistryTest.java   |  58 ++++++++
 11 files changed, 456 insertions(+), 161 deletions(-)

diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionFactory.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionFactory.java
index fdc93df43e8..46670dba85c 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionFactory.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionFactory.java
@@ -23,20 +23,6 @@ import java.util.List;
 import org.apache.datasketches.tuple.aninteger.IntegerSummary;
 import org.apache.pinot.common.request.context.ExpressionContext;
 import org.apache.pinot.common.request.context.FunctionContext;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggBigDecimalFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggBytesFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDistinctBigDecimalFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDistinctBytesFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDistinctDoubleFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDistinctFloatFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDistinctIntFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDistinctLongFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDistinctStringFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggDoubleFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggFloatFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggIntFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggLongFunction;
-import 
org.apache.pinot.core.query.aggregation.function.array.ArrayAggStringFunction;
 import 
org.apache.pinot.core.query.aggregation.function.array.ListAggDistinctFunction;
 import org.apache.pinot.core.query.aggregation.function.array.ListAggFunction;
 import 
org.apache.pinot.core.query.aggregation.function.array.SumArrayDoubleAggregationFunction;
@@ -48,7 +34,6 @@ import 
org.apache.pinot.core.query.aggregation.function.funnel.window.FunnelMatc
 import 
org.apache.pinot.core.query.aggregation.function.funnel.window.FunnelMaxStepAggregationFunction;
 import 
org.apache.pinot.core.query.aggregation.function.funnel.window.FunnelStepDurationStatsAggregationFunction;
 import org.apache.pinot.segment.spi.AggregationFunctionType;
-import org.apache.pinot.spi.data.FieldSpec.DataType;
 import org.apache.pinot.spi.exception.BadQueryRequestException;
 
 
@@ -217,6 +202,10 @@ public class AggregationFunctionFactory {
         throw new IllegalArgumentException("Invalid percentile function: " + 
function);
       } else {
         AggregationFunctionType functionType = 
AggregationFunctionType.valueOf(upperCaseFunctionName);
+        AggregationFunctionProvider provider = 
AggregationFunctionProviderRegistry.getProvider(functionType);
+        if (provider != null) {
+          return provider.create(function, nullHandlingEnabled);
+        }
         switch (functionType) {
           case COUNT:
             return new CountAggregationFunction(arguments, 
nullHandlingEnabled);
@@ -243,39 +232,6 @@ public class AggregationFunctionFactory {
             return new SumPrecisionAggregationFunction(arguments, 
nullHandlingEnabled);
           case AVG:
             return new AvgAggregationFunction(arguments, nullHandlingEnabled);
-          case MODE:
-            return new ModeAggregationFunction(arguments, nullHandlingEnabled);
-          case ANYVALUE:
-            return new AnyValueAggregationFunction(arguments, 
nullHandlingEnabled);
-          case FIRSTWITHTIME: {
-            Preconditions.checkArgument(numArguments == 3,
-                "FIRST_WITH_TIME expects 3 arguments, got: %s. The function 
can be used as "
-                    + "firstWithTime(dataColumn, timeColumn, 'dataType')", 
numArguments);
-            ExpressionContext timeCol = arguments.get(1);
-            ExpressionContext dataTypeExp = arguments.get(2);
-            Preconditions.checkArgument(dataTypeExp.getType() == 
ExpressionContext.Type.LITERAL,
-                "FIRST_WITH_TIME expects the 3rd argument to be literal, got: 
%s. The function can be used as "
-                    + "firstWithTime(dataColumn, timeColumn, 'dataType')", 
dataTypeExp.getType());
-            DataType dataType = 
DataType.valueOf(dataTypeExp.getLiteral().getStringValue().toUpperCase());
-            switch (dataType) {
-              case BOOLEAN:
-                return new 
FirstIntValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled,
-                    true);
-              case INT:
-                return new 
FirstIntValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled,
-                    false);
-              case LONG:
-                return new 
FirstLongValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              case FLOAT:
-                return new 
FirstFloatValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              case DOUBLE:
-                return new 
FirstDoubleValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              case STRING:
-                return new 
FirstStringValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              default:
-                throw new IllegalArgumentException("Unsupported data type for 
FIRST_WITH_TIME: " + dataType);
-            }
-          }
           case LISTAGG: {
             Preconditions.checkArgument(numArguments == 2 || numArguments == 3,
                 "LISTAGG expects 2 or 3 arguments, got: %s. The function can 
be used as "
@@ -302,97 +258,6 @@ public class AggregationFunctionFactory {
             return new SumArrayLongAggregationFunction(arguments, 
nullHandlingEnabled);
           case SUMARRAYDOUBLE:
             return new SumArrayDoubleAggregationFunction(arguments, 
nullHandlingEnabled);
-          case ARRAYAGG: {
-            Preconditions.checkArgument(numArguments >= 2,
-                "ARRAY_AGG expects 2 or 3 arguments, got: %s. The function can 
be used as "
-                    + "arrayAgg(dataColumn, 'dataType', ['isDistinct'])", 
numArguments);
-            ExpressionContext dataTypeExp = arguments.get(1);
-            Preconditions.checkArgument(dataTypeExp.getType() == 
ExpressionContext.Type.LITERAL,
-                "ARRAY_AGG expects the 2nd argument to be literal, got: %s. 
The function can be used as "
-                    + "arrayAgg(dataColumn, 'dataType', ['isDistinct'])", 
dataTypeExp.getType());
-            DataType dataType = 
DataType.valueOf(dataTypeExp.getLiteral().getStringValue().toUpperCase());
-            boolean isDistinct = false;
-            if (numArguments == 3) {
-              ExpressionContext isDistinctExp = arguments.get(2);
-              Preconditions.checkArgument(isDistinctExp.getType() == 
ExpressionContext.Type.LITERAL,
-                  "ARRAY_AGG expects the 3rd argument to be literal, got: %s. 
The function can be used as "
-                      + "arrayAgg(dataColumn, 'dataType', ['isDistinct'])", 
isDistinctExp.getType());
-              isDistinct = isDistinctExp.getLiteral().getBooleanValue();
-            }
-            if (isDistinct) {
-              switch (dataType) {
-                case BOOLEAN:
-                case INT:
-                  return new ArrayAggDistinctIntFunction(firstArgument, 
dataType, nullHandlingEnabled);
-                case LONG:
-                case TIMESTAMP:
-                  return new ArrayAggDistinctLongFunction(firstArgument, 
dataType, nullHandlingEnabled);
-                case FLOAT:
-                  return new ArrayAggDistinctFloatFunction(firstArgument, 
nullHandlingEnabled);
-                case DOUBLE:
-                  return new ArrayAggDistinctDoubleFunction(firstArgument, 
nullHandlingEnabled);
-                case BIG_DECIMAL:
-                  return new ArrayAggDistinctBigDecimalFunction(firstArgument, 
nullHandlingEnabled);
-                case STRING:
-                case JSON:
-                  return new ArrayAggDistinctStringFunction(firstArgument, 
nullHandlingEnabled);
-                case BYTES:
-                case UUID:
-                  return new ArrayAggDistinctBytesFunction(firstArgument, 
dataType, nullHandlingEnabled);
-                default:
-                  throw new IllegalArgumentException("Unsupported data type 
for ARRAY_AGG: " + dataType);
-              }
-            }
-            switch (dataType) {
-              case BOOLEAN:
-              case INT:
-                return new ArrayAggIntFunction(firstArgument, dataType, 
nullHandlingEnabled);
-              case LONG:
-              case TIMESTAMP:
-                return new ArrayAggLongFunction(firstArgument, dataType, 
nullHandlingEnabled);
-              case FLOAT:
-                return new ArrayAggFloatFunction(firstArgument, 
nullHandlingEnabled);
-              case DOUBLE:
-                return new ArrayAggDoubleFunction(firstArgument, 
nullHandlingEnabled);
-              case BIG_DECIMAL:
-                return new ArrayAggBigDecimalFunction(firstArgument, 
nullHandlingEnabled);
-              case STRING:
-              case JSON:
-                return new ArrayAggStringFunction(firstArgument, 
nullHandlingEnabled);
-              case BYTES:
-              case UUID:
-                return new ArrayAggBytesFunction(firstArgument, dataType, 
nullHandlingEnabled);
-              default:
-                throw new IllegalArgumentException("Unsupported data type for 
ARRAY_AGG: " + dataType);
-            }
-          }
-          case LASTWITHTIME: {
-            Preconditions.checkArgument(numArguments == 3,
-                "LAST_WITH_TIME expects 3 arguments, got: %s. The function can 
be used as "
-                    + "lastWithTime(dataColumn, timeColumn, 'dataType')", 
numArguments);
-            ExpressionContext timeCol = arguments.get(1);
-            ExpressionContext dataTypeExp = arguments.get(2);
-            Preconditions.checkArgument(dataTypeExp.getType() == 
ExpressionContext.Type.LITERAL,
-                "LAST_WITH_TIME expects the 3rd argument to be literal, got: 
%s. The function can be used as "
-                    + "lastWithTime(dataColumn, timeColumn, 'dataType')", 
dataTypeExp.getType());
-            DataType dataType = 
DataType.valueOf(dataTypeExp.getLiteral().getStringValue().toUpperCase());
-            switch (dataType) {
-              case BOOLEAN:
-                return new 
LastIntValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled, true);
-              case INT:
-                return new 
LastIntValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled, false);
-              case LONG:
-                return new 
LastLongValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              case FLOAT:
-                return new 
LastFloatValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              case DOUBLE:
-                return new 
LastDoubleValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              case STRING:
-                return new 
LastStringValueWithTimeAggregationFunction(firstArgument, timeCol, 
nullHandlingEnabled);
-              default:
-                throw new IllegalArgumentException("Unsupported data type for 
LAST_WITH_TIME: " + dataType);
-            }
-          }
           case MINMAXRANGE:
             return new MinMaxRangeAggregationFunction(arguments, 
nullHandlingEnabled);
           case DISTINCTCOUNT:
@@ -498,14 +363,6 @@ public class AggregationFunctionFactory {
           case AVGVALUEINTEGERSUMTUPLESKETCH:
             return new 
AvgValueIntegerTupleSketchAggregationFunction(arguments, 
IntegerSummary.Mode.Sum,
                 nullHandlingEnabled);
-          case PINOTPARENTAGGEXPRMAX:
-            return new ParentExprMinMaxAggregationFunction(arguments, true, 
nullHandlingEnabled);
-          case PINOTPARENTAGGEXPRMIN:
-            return new ParentExprMinMaxAggregationFunction(arguments, false, 
nullHandlingEnabled);
-          case PINOTCHILDAGGEXPRMAX:
-            return new ChildExprMinMaxAggregationFunction(arguments, true);
-          case PINOTCHILDAGGEXPRMIN:
-            return new ChildExprMinMaxAggregationFunction(arguments, false);
           case EXPRMAX:
           case EXPRMIN:
             throw new IllegalArgumentException(
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ChildExprMinMaxAggregationFunction.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProvider.java
similarity index 63%
copy from 
pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ChildExprMinMaxAggregationFunction.java
copy to 
pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProvider.java
index ec2be01e140..ee654c4eba9 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ChildExprMinMaxAggregationFunction.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProvider.java
@@ -18,22 +18,15 @@
  */
 package org.apache.pinot.core.query.aggregation.function;
 
-import java.util.List;
-import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
 import org.apache.pinot.segment.spi.AggregationFunctionType;
 
 
-public class ChildExprMinMaxAggregationFunction extends 
ChildAggregationFunction {
+/// Constructs one aggregate family from its whole call, so the family owns 
its argument validation and type dispatch
+/// instead of the core factory. Providers are discovered once through 
ServiceLoader, must be stateless and
+/// thread-safe, and return a new query-local function from each call.
+public interface AggregationFunctionProvider {
+  AggregationFunctionType getType();
 
-  private final boolean _isMax;
-
-  public ChildExprMinMaxAggregationFunction(List<ExpressionContext> operands, 
boolean isMax) {
-    super(operands);
-    _isMax = isMax;
-  }
-
-  @Override
-  public AggregationFunctionType getType() {
-    return _isMax ? AggregationFunctionType.EXPRMAX : 
AggregationFunctionType.EXPRMIN;
-  }
+  AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled);
 }
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProviderRegistry.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProviderRegistry.java
new file mode 100644
index 00000000000..cc6e4a57084
--- /dev/null
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProviderRegistry.java
@@ -0,0 +1,58 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.query.aggregation.function;
+
+import com.google.common.annotations.VisibleForTesting;
+import java.util.EnumMap;
+import java.util.Map;
+import java.util.Objects;
+import java.util.ServiceLoader;
+import javax.annotation.Nullable;
+import org.apache.pinot.segment.spi.AggregationFunctionType;
+
+import static com.google.common.base.Preconditions.checkState;
+
+
+/// Immutable aggregate-provider lookup, initialized once from the 
application's service registrations. Conflicting
+/// registrations fail at initialization instead of depending on classpath 
order. Safe for concurrent query creation.
+public final class AggregationFunctionProviderRegistry {
+  private static final Map<AggregationFunctionType, 
AggregationFunctionProvider> PROVIDERS = loadProviders(
+      ServiceLoader.load(AggregationFunctionProvider.class, 
AggregationFunctionProvider.class.getClassLoader()));
+
+  private AggregationFunctionProviderRegistry() {
+  }
+
+  @Nullable
+  public static AggregationFunctionProvider 
getProvider(AggregationFunctionType type) {
+    return PROVIDERS.get(type);
+  }
+
+  @VisibleForTesting
+  static Map<AggregationFunctionType, AggregationFunctionProvider> 
loadProviders(
+      Iterable<AggregationFunctionProvider> providers) {
+    Map<AggregationFunctionType, AggregationFunctionProvider> result = new 
EnumMap<>(AggregationFunctionType.class);
+    for (AggregationFunctionProvider provider : providers) {
+      AggregationFunctionType type = 
Objects.requireNonNull(provider.getType(), "Provider must declare an 
aggregate");
+      AggregationFunctionProvider previous = result.putIfAbsent(type, 
provider);
+      checkState(previous == null, "Duplicate aggregation provider for %s: %s 
and %s", type,
+          previous != null ? previous.getClass().getName() : "", 
provider.getClass().getName());
+    }
+    return Map.copyOf(result);
+  }
+}
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AnyValueAggregationFunction.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AnyValueAggregationFunction.java
index d593327ffc4..e6669a16a3a 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AnyValueAggregationFunction.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AnyValueAggregationFunction.java
@@ -27,6 +27,7 @@ import java.util.function.Consumer;
 import javax.annotation.Nullable;
 import org.apache.pinot.common.CustomObject;
 import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
 import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
 import org.apache.pinot.core.common.BlockValSet;
 import org.apache.pinot.core.query.aggregation.AggregationResultHolder;
@@ -353,4 +354,17 @@ public class AnyValueAggregationFunction extends 
BaseSingleInputAggregationFunct
         throw new IllegalStateException("ANY_VALUE unsupported type: " + 
bvs.getValueType());
     }
   }
+
+  /// Service registration for ANY_VALUE.
+  public static final class Provider implements AggregationFunctionProvider {
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.ANYVALUE;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      return new AnyValueAggregationFunction(function.getArguments(), 
nullHandlingEnabled);
+    }
+  }
 }
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ChildExprMinMaxAggregationFunction.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ChildExprMinMaxAggregationFunction.java
index ec2be01e140..9a5f38597b7 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ChildExprMinMaxAggregationFunction.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ChildExprMinMaxAggregationFunction.java
@@ -20,6 +20,7 @@ package org.apache.pinot.core.query.aggregation.function;
 
 import java.util.List;
 import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
 import org.apache.pinot.segment.spi.AggregationFunctionType;
 
 
@@ -36,4 +37,30 @@ public class ChildExprMinMaxAggregationFunction extends 
ChildAggregationFunction
   public AggregationFunctionType getType() {
     return _isMax ? AggregationFunctionType.EXPRMAX : 
AggregationFunctionType.EXPRMIN;
   }
+
+  /// Service registration for the ExprMin child aggregation.
+  public static final class MinProvider implements AggregationFunctionProvider 
{
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.PINOTCHILDAGGEXPRMIN;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      return new ChildExprMinMaxAggregationFunction(function.getArguments(), 
false);
+    }
+  }
+
+  /// Service registration for the ExprMax child aggregation.
+  public static final class MaxProvider implements AggregationFunctionProvider 
{
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.PINOTCHILDAGGEXPRMAX;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      return new ChildExprMinMaxAggregationFunction(function.getArguments(), 
true);
+    }
+  }
 }
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/FirstLastWithTimeAggregationFunctionProvider.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/FirstLastWithTimeAggregationFunctionProvider.java
new file mode 100644
index 00000000000..35e7ca647e5
--- /dev/null
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/FirstLastWithTimeAggregationFunctionProvider.java
@@ -0,0 +1,114 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.query.aggregation.function;
+
+import com.google.common.base.Preconditions;
+import java.util.List;
+import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
+import org.apache.pinot.segment.spi.AggregationFunctionType;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
+
+
+/// Creates FIRST_WITH_TIME and LAST_WITH_TIME kernels from their explicit 
data type argument. The providers are
+/// stateless and safe for concurrent query construction; each call creates a 
new aggregation function.
+public final class FirstLastWithTimeAggregationFunctionProvider {
+  private FirstLastWithTimeAggregationFunctionProvider() {
+  }
+
+  /// Service registration for FIRST_WITH_TIME.
+  public static final class First implements AggregationFunctionProvider {
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.FIRSTWITHTIME;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      List<ExpressionContext> arguments = function.getArguments();
+      int numArguments = arguments.size();
+      Preconditions.checkArgument(numArguments == 3,
+          "FIRST_WITH_TIME expects 3 arguments, got: %s. The function can be 
used as "
+              + "firstWithTime(dataColumn, timeColumn, 'dataType')", 
numArguments);
+      ExpressionContext dataCol = arguments.get(0);
+      ExpressionContext timeCol = arguments.get(1);
+      ExpressionContext dataTypeExp = arguments.get(2);
+      Preconditions.checkArgument(dataTypeExp.getType() == 
ExpressionContext.Type.LITERAL,
+          "FIRST_WITH_TIME expects the 3rd argument to be literal, got: %s. 
The function can be used as "
+              + "firstWithTime(dataColumn, timeColumn, 'dataType')", 
dataTypeExp.getType());
+      DataType dataType = 
DataType.valueOf(dataTypeExp.getLiteral().getStringValue().toUpperCase());
+      switch (dataType) {
+        case BOOLEAN:
+          return new FirstIntValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled, true);
+        case INT:
+          return new FirstIntValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled, false);
+        case LONG:
+          return new FirstLongValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        case FLOAT:
+          return new FirstFloatValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        case DOUBLE:
+          return new FirstDoubleValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        case STRING:
+          return new FirstStringValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        default:
+          throw new IllegalArgumentException("Unsupported data type for 
FIRST_WITH_TIME: " + dataType);
+      }
+    }
+  }
+
+  /// Service registration for LAST_WITH_TIME.
+  public static final class Last implements AggregationFunctionProvider {
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.LASTWITHTIME;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      List<ExpressionContext> arguments = function.getArguments();
+      int numArguments = arguments.size();
+      Preconditions.checkArgument(numArguments == 3,
+          "LAST_WITH_TIME expects 3 arguments, got: %s. The function can be 
used as "
+              + "lastWithTime(dataColumn, timeColumn, 'dataType')", 
numArguments);
+      ExpressionContext dataCol = arguments.get(0);
+      ExpressionContext timeCol = arguments.get(1);
+      ExpressionContext dataTypeExp = arguments.get(2);
+      Preconditions.checkArgument(dataTypeExp.getType() == 
ExpressionContext.Type.LITERAL,
+          "LAST_WITH_TIME expects the 3rd argument to be literal, got: %s. The 
function can be used as "
+              + "lastWithTime(dataColumn, timeColumn, 'dataType')", 
dataTypeExp.getType());
+      DataType dataType = 
DataType.valueOf(dataTypeExp.getLiteral().getStringValue().toUpperCase());
+      switch (dataType) {
+        case BOOLEAN:
+          return new LastIntValueWithTimeAggregationFunction(dataCol, timeCol, 
nullHandlingEnabled, true);
+        case INT:
+          return new LastIntValueWithTimeAggregationFunction(dataCol, timeCol, 
nullHandlingEnabled, false);
+        case LONG:
+          return new LastLongValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        case FLOAT:
+          return new LastFloatValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        case DOUBLE:
+          return new LastDoubleValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        case STRING:
+          return new LastStringValueWithTimeAggregationFunction(dataCol, 
timeCol, nullHandlingEnabled);
+        default:
+          throw new IllegalArgumentException("Unsupported data type for 
LAST_WITH_TIME: " + dataType);
+      }
+    }
+  }
+}
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ModeAggregationFunction.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ModeAggregationFunction.java
index b17437b1a2f..95ccf22bbe4 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ModeAggregationFunction.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ModeAggregationFunction.java
@@ -36,6 +36,7 @@ import java.util.Map;
 import javax.annotation.Nullable;
 import org.apache.pinot.common.CustomObject;
 import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
 import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
 import org.apache.pinot.core.common.BlockValSet;
 import org.apache.pinot.core.common.ObjectSerDeUtils;
@@ -733,4 +734,17 @@ public class ModeAggregationFunction extends 
BaseSingleInputAggregationFunction<
       _dictIdCountMap = new Int2IntOpenHashMap();
     }
   }
+
+  /// Service registration for MODE.
+  public static final class Provider implements AggregationFunctionProvider {
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.MODE;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      return new ModeAggregationFunction(function.getArguments(), 
nullHandlingEnabled);
+    }
+  }
 }
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ParentExprMinMaxAggregationFunction.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ParentExprMinMaxAggregationFunction.java
index 4624fc1d794..b806e9bf165 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ParentExprMinMaxAggregationFunction.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/ParentExprMinMaxAggregationFunction.java
@@ -24,6 +24,7 @@ import java.util.Map;
 import javax.annotation.Nullable;
 import org.apache.pinot.common.CustomObject;
 import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
 import org.apache.pinot.common.utils.DataSchema;
 import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
 import org.apache.pinot.common.utils.RoaringBitmapUtils.BatchConsumer;
@@ -357,4 +358,30 @@ public class ParentExprMinMaxAggregationFunction extends 
ParentAggregationFuncti
   public ExprMinMaxObject extractFinalResult(@Nullable ExprMinMaxObject 
exprMinMaxObject) {
     return exprMinMaxObject;
   }
+
+  /// Service registration for the ExprMin parent aggregation.
+  public static final class MinProvider implements AggregationFunctionProvider 
{
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.PINOTPARENTAGGEXPRMIN;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      return new ParentExprMinMaxAggregationFunction(function.getArguments(), 
false, nullHandlingEnabled);
+    }
+  }
+
+  /// Service registration for the ExprMax parent aggregation.
+  public static final class MaxProvider implements AggregationFunctionProvider 
{
+    @Override
+    public AggregationFunctionType getType() {
+      return AggregationFunctionType.PINOTPARENTAGGEXPRMAX;
+    }
+
+    @Override
+    public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+      return new ParentExprMinMaxAggregationFunction(function.getArguments(), 
true, nullHandlingEnabled);
+    }
+  }
 }
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/array/ArrayAggFunctionProvider.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/array/ArrayAggFunctionProvider.java
new file mode 100644
index 00000000000..fdd28293fd9
--- /dev/null
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/array/ArrayAggFunctionProvider.java
@@ -0,0 +1,107 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.query.aggregation.function.array;
+
+import com.google.common.base.Preconditions;
+import java.util.List;
+import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
+import org.apache.pinot.core.query.aggregation.function.AggregationFunction;
+import 
org.apache.pinot.core.query.aggregation.function.AggregationFunctionProvider;
+import org.apache.pinot.segment.spi.AggregationFunctionType;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
+
+
+/// Selects the ARRAY_AGG accumulator implementation from the explicit data 
type and distinct arguments. The provider
+/// is stateless and safe to share; each call creates a new aggregation 
function.
+public final class ArrayAggFunctionProvider implements 
AggregationFunctionProvider {
+  @Override
+  public AggregationFunctionType getType() {
+    return AggregationFunctionType.ARRAYAGG;
+  }
+
+  @Override
+  public AggregationFunction<?, ?> create(FunctionContext function, boolean 
nullHandlingEnabled) {
+    List<ExpressionContext> arguments = function.getArguments();
+    int numArguments = arguments.size();
+    Preconditions.checkArgument(numArguments >= 2,
+        "ARRAY_AGG expects 2 or 3 arguments, got: %s. The function can be used 
as "
+            + "arrayAgg(dataColumn, 'dataType', ['isDistinct'])", 
numArguments);
+    ExpressionContext expression = arguments.get(0);
+    ExpressionContext dataTypeExp = arguments.get(1);
+    Preconditions.checkArgument(dataTypeExp.getType() == 
ExpressionContext.Type.LITERAL,
+        "ARRAY_AGG expects the 2nd argument to be literal, got: %s. The 
function can be used as "
+            + "arrayAgg(dataColumn, 'dataType', ['isDistinct'])", 
dataTypeExp.getType());
+    DataType dataType = 
DataType.valueOf(dataTypeExp.getLiteral().getStringValue().toUpperCase());
+    boolean isDistinct = false;
+    if (numArguments == 3) {
+      ExpressionContext isDistinctExp = arguments.get(2);
+      Preconditions.checkArgument(isDistinctExp.getType() == 
ExpressionContext.Type.LITERAL,
+          "ARRAY_AGG expects the 3rd argument to be literal, got: %s. The 
function can be used as "
+              + "arrayAgg(dataColumn, 'dataType', ['isDistinct'])", 
isDistinctExp.getType());
+      isDistinct = isDistinctExp.getLiteral().getBooleanValue();
+    }
+    if (isDistinct) {
+      switch (dataType) {
+        case BOOLEAN:
+        case INT:
+          return new ArrayAggDistinctIntFunction(expression, dataType, 
nullHandlingEnabled);
+        case LONG:
+        case TIMESTAMP:
+          return new ArrayAggDistinctLongFunction(expression, dataType, 
nullHandlingEnabled);
+        case FLOAT:
+          return new ArrayAggDistinctFloatFunction(expression, 
nullHandlingEnabled);
+        case DOUBLE:
+          return new ArrayAggDistinctDoubleFunction(expression, 
nullHandlingEnabled);
+        case BIG_DECIMAL:
+          return new ArrayAggDistinctBigDecimalFunction(expression, 
nullHandlingEnabled);
+        case STRING:
+        case JSON:
+          return new ArrayAggDistinctStringFunction(expression, 
nullHandlingEnabled);
+        case BYTES:
+        case UUID:
+          return new ArrayAggDistinctBytesFunction(expression, dataType, 
nullHandlingEnabled);
+        default:
+          throw new IllegalArgumentException("Unsupported data type for 
ARRAY_AGG: " + dataType);
+      }
+    }
+    switch (dataType) {
+      case BOOLEAN:
+      case INT:
+        return new ArrayAggIntFunction(expression, dataType, 
nullHandlingEnabled);
+      case LONG:
+      case TIMESTAMP:
+        return new ArrayAggLongFunction(expression, dataType, 
nullHandlingEnabled);
+      case FLOAT:
+        return new ArrayAggFloatFunction(expression, nullHandlingEnabled);
+      case DOUBLE:
+        return new ArrayAggDoubleFunction(expression, nullHandlingEnabled);
+      case BIG_DECIMAL:
+        return new ArrayAggBigDecimalFunction(expression, nullHandlingEnabled);
+      case STRING:
+      case JSON:
+        return new ArrayAggStringFunction(expression, nullHandlingEnabled);
+      case BYTES:
+      case UUID:
+        return new ArrayAggBytesFunction(expression, dataType, 
nullHandlingEnabled);
+      default:
+        throw new IllegalArgumentException("Unsupported data type for 
ARRAY_AGG: " + dataType);
+    }
+  }
+}
diff --git 
a/pinot-core/src/main/resources/META-INF/services/org.apache.pinot.core.query.aggregation.function.AggregationFunctionProvider
 
b/pinot-core/src/main/resources/META-INF/services/org.apache.pinot.core.query.aggregation.function.AggregationFunctionProvider
new file mode 100644
index 00000000000..8021dcda997
--- /dev/null
+++ 
b/pinot-core/src/main/resources/META-INF/services/org.apache.pinot.core.query.aggregation.function.AggregationFunctionProvider
@@ -0,0 +1,26 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+org.apache.pinot.core.query.aggregation.function.ModeAggregationFunction$Provider
+org.apache.pinot.core.query.aggregation.function.FirstLastWithTimeAggregationFunctionProvider$First
+org.apache.pinot.core.query.aggregation.function.FirstLastWithTimeAggregationFunctionProvider$Last
+org.apache.pinot.core.query.aggregation.function.AnyValueAggregationFunction$Provider
+org.apache.pinot.core.query.aggregation.function.array.ArrayAggFunctionProvider
+org.apache.pinot.core.query.aggregation.function.ParentExprMinMaxAggregationFunction$MinProvider
+org.apache.pinot.core.query.aggregation.function.ParentExprMinMaxAggregationFunction$MaxProvider
+org.apache.pinot.core.query.aggregation.function.ChildExprMinMaxAggregationFunction$MinProvider
+org.apache.pinot.core.query.aggregation.function.ChildExprMinMaxAggregationFunction$MaxProvider
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProviderRegistryTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProviderRegistryTest.java
new file mode 100644
index 00000000000..3827a84079f
--- /dev/null
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionProviderRegistryTest.java
@@ -0,0 +1,58 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.core.query.aggregation.function;
+
+import java.util.List;
+import java.util.Map;
+import org.apache.pinot.segment.spi.AggregationFunctionType;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertNull;
+import static org.testng.Assert.assertSame;
+import static org.testng.Assert.expectThrows;
+
+
+/// Verifies classpath provider discovery and deterministic rejection of 
conflicting aggregate registrations.
+public class AggregationFunctionProviderRegistryTest {
+  @Test
+  public void testDeclaredFamiliesAreDiscoverable() {
+    for (AggregationFunctionType type : List.of(AggregationFunctionType.MODE, 
AggregationFunctionType.FIRSTWITHTIME,
+        AggregationFunctionType.LASTWITHTIME, 
AggregationFunctionType.ANYVALUE, AggregationFunctionType.ARRAYAGG,
+        AggregationFunctionType.PINOTPARENTAGGEXPRMIN, 
AggregationFunctionType.PINOTPARENTAGGEXPRMAX,
+        AggregationFunctionType.PINOTCHILDAGGEXPRMIN, 
AggregationFunctionType.PINOTCHILDAGGEXPRMAX)) {
+      AggregationFunctionProvider provider = 
AggregationFunctionProviderRegistry.getProvider(type);
+      assertNotNull(provider, type.name());
+      assertEquals(provider.getType(), type);
+    }
+    
assertNull(AggregationFunctionProviderRegistry.getProvider(AggregationFunctionType.COUNT));
+  }
+
+  @Test
+  public void testDuplicateProvidersFailAndRegistryIsImmutable() {
+    AggregationFunctionProvider provider = new 
ModeAggregationFunction.Provider();
+    Map<AggregationFunctionType, AggregationFunctionProvider> registry =
+        AggregationFunctionProviderRegistry.loadProviders(List.of(provider));
+    assertSame(registry.get(AggregationFunctionType.MODE), provider);
+    expectThrows(UnsupportedOperationException.class, registry::clear);
+    expectThrows(IllegalStateException.class, () -> 
AggregationFunctionProviderRegistry.loadProviders(
+        List.of(provider, new ModeAggregationFunction.Provider())));
+  }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to