Jackie-Jiang commented on code in PR #19595:
URL: https://github.com/apache/pinot/pull/19595#discussion_r4052461947
##########
pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/BaseSingleTreeBuilder.java:
##########
@@ -148,7 +148,16 @@ static class Record {
// Ignore the column for COUNT aggregation function
if (_valueAggregators[index].getAggregationType() !=
AggregationFunctionType.COUNT) {
String column = functionColumnPair.getColumn();
- _metricReaders[index] = new PinotSegmentColumnReader(segment, column);
+ PinotSegmentColumnReader metricReader = new
PinotSegmentColumnReader(segment, column);
+ // The distinct arrayAgg star-tree aggregator only supports
single-value, dictionary-encoded source columns.
+ // See ArrayAggDistinctValueAggregator.
+ if (_valueAggregators[index].getAggregationType() ==
AggregationFunctionType.ARRAYAGG) {
+ Preconditions.checkState(metricReader.isSingleValue(),
+ "Star-tree arrayAgg does not support multi-value column: %s",
column);
+ Preconditions.checkState(metricReader.hasDictionary(),
Review Comment:
Non-blocking: the dictionary requirement appears unnecessary for the
aggregation metric. The aggregator consumes scalar values, and both builder
read paths use `PinotSegmentColumnReader.getValue()`, which already supports
raw forward indexes for the supported types. Consider removing this requirement
while retaining the single-value restriction, with a raw-metric build/query
test to cover it.
##########
pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AggregationFunctionUtils.java:
##########
@@ -176,11 +177,15 @@ public static Map<ExpressionContext, BlockValSet>
getBlockValSetMap(AggregationF
///
/// NOTE: We construct the map with original column name as the key but
fetch BlockValSet with the aggregation
/// function pair so that the aggregation result column name is
consistent with or without star-tree.
+ ///
+ /// The BlockValSet is wrapped in [StarTreePreAggregatedBlockValSet] so
aggregation functions can recognize a
+ /// pre-aggregated column by provenance rather than by sniffing the block's
value type, which is ambiguous when the
+ /// raw input type is also BYTES (e.g. distinct arrayAgg over a BYTES
column).
public static Map<ExpressionContext, BlockValSet> getBlockValSetMap(
AggregationFunctionColumnPair aggregationFunctionColumnPair, ValueBlock
valueBlock) {
ExpressionContext expression =
ExpressionContext.forIdentifier(aggregationFunctionColumnPair.getColumn());
BlockValSet blockValSet =
valueBlock.getBlockValueSet(aggregationFunctionColumnPair.toColumnName());
- return Map.of(expression, blockValSet);
+ return Map.of(expression, new
StarTreePreAggregatedBlockValSet(blockValSet));
Review Comment:
Non-blocking: consider applying `StarTreePreAggregatedBlockValSet` only when
the function type is ARRAYAGG. Currently every star-tree aggregation gets a
wrapper, although only the ARRAYAGG implementations inspect the marker.
Narrowing this preserves the provenance fix while avoiding unnecessary
allocations and changes to other aggregations' input objects.
##########
pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/array/BaseArrayAggFunction.java:
##########
@@ -26,24 +28,56 @@
import
org.apache.pinot.core.query.aggregation.function.BaseSingleInputAggregationFunction;
import org.apache.pinot.core.query.aggregation.groupby.GroupByResultHolder;
import
org.apache.pinot.core.query.aggregation.groupby.ObjectGroupByResultHolder;
+import
org.apache.pinot.segment.local.aggregator.ArrayAggDistinctValueAggregator.ElementType;
import org.apache.pinot.segment.spi.AggregationFunctionType;
import org.apache.pinot.spi.data.FieldSpec.DataType;
public abstract class BaseArrayAggFunction<I, F extends Comparable> extends
BaseSingleInputAggregationFunction<I, F> {
private final DataSchema.ColumnDataType _resultColumnType;
+ private final DataType _elementDataType;
Review Comment:
Non-blocking: `_elementDataType` duplicates information already available
from the result type. The compatibility check can use
`getFinalResultColumnType().toDataType().getStoredType()`, since
`ColumnDataType.toDataType()` handles the array result types. This would let us
remove the extra field and getter while retaining the stored-type compatibility
check.
##########
pinot-segment-local/src/main/java/org/apache/pinot/segment/local/aggregator/ArrayAggDistinctValueAggregator.java:
##########
@@ -0,0 +1,419 @@
+/**
+ * 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.segment.local.aggregator;
+
+import it.unimi.dsi.fastutil.objects.ObjectOpenHashSet;
+import it.unimi.dsi.fastutil.objects.ObjectSet;
+import java.math.BigDecimal;
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import javax.annotation.Nullable;
+import org.apache.pinot.segment.spi.AggregationFunctionType;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
+import org.apache.pinot.spi.utils.BigDecimalUtils;
+import org.apache.pinot.spi.utils.ByteArray;
+
+
+/// Value aggregator for distinct `ARRAYAGG` (i.e. `arrayAgg(column, 'type',
true)`) used by the star-tree index.
+///
+/// The aggregated value is the set of distinct (non-null) raw values seen
under a star-tree node. Distinct semantics
+/// make the aggregation associative and idempotent (merge = set-union), which
is what the star-tree requires; the
+/// non-distinct variant would need an unbounded per-node multiset and is
intentionally not supported here.
+///
+/// The set is stored as a variable-length `BYTES` cell in the star-tree
forward index. The serialized layout is
+/// intentionally byte-for-byte identical to `ObjectSerDeUtils.*_SET_SER_DE`
(in pinot-core) so the query-time
+/// `arrayAgg` function can deserialize the stored cell directly. The wire
layout cannot be shared as code because
+/// pinot-segment-local must not depend on pinot-core; any change here must be
mirrored there.
+///
+/// NOTE: Only single-value source columns are supported (enforced at
star-tree build config validation time), so the
+/// raw value handed in by the builder is always a single boxed scalar, never
an array.
+public class ArrayAggDistinctValueAggregator implements
ValueAggregator<Object, ObjectSet<Object>> {
+ public static final DataType AGGREGATED_VALUE_TYPE = DataType.BYTES;
+
+ // Runtime element kind, inferred lazily from the first raw value (mirrors
DistinctCountThetaSketchValueAggregator's
+ // runtime dispatch). All values in a given aggregator instance come from
the same column, so the kind is stable.
+ // Each value carries an explicit stable on-disk `tag` that is persisted as
a leading byte in the serialized cell,
+ // making deserialization self-describing (independent of aggregator
instance state, needed for build re-ingestion).
+ // The query-time arrayAgg function reads and validates the same tag before
decoding the payload.
+ //
+ // IMPORTANT: tags are an on-disk format. Never change or reuse an existing
value's tag; only append new values with
+ // new tags. Reordering the enum is safe because the tag, not the ordinal,
is persisted.
+ public enum ElementType {
+ INT(0), LONG(1), FLOAT(2), DOUBLE(3), BIG_DECIMAL(4), STRING(5), BYTES(6);
+
+ private final byte _tag;
+
+ ElementType(int tag) {
+ _tag = (byte) tag;
+ }
+
+ public byte getTag() {
+ return _tag;
+ }
+
+ private static final ElementType[] BY_TAG;
+
+ static {
+ ElementType[] values = values();
+ int maxTag = 0;
+ for (ElementType value : values) {
+ maxTag = Math.max(maxTag, value._tag);
+ }
+ BY_TAG = new ElementType[maxTag + 1];
+ for (ElementType value : values) {
+ BY_TAG[value._tag] = value;
+ }
+ }
+
+ public static ElementType fromTag(byte tag) {
+ if (tag < 0 || tag >= BY_TAG.length || BY_TAG[tag] == null) {
+ throw new IllegalStateException("Invalid arrayAgg element type tag: "
+ tag);
+ }
+ return BY_TAG[tag];
+ }
+ }
+
+ /// Distinct value set carrying a running total of its variable-width
serialized payload bytes, so cell-size
+ /// tracking is O(1) per insertion instead of re-serializing the growing
set. Instances are only created by this
+ /// aggregator; callers see a plain ObjectSet with content-based equality,
so the extra state is invisible to them.
+ private static class TrackedSet extends ObjectOpenHashSet<Object> {
+ // Sum over elements of (4-byte length prefix + encoded length); only
maintained for variable-width element types.
+ private long _variableWidthPayloadBytes;
+
+ TrackedSet() {
+ }
+
+ TrackedSet(int expectedSize) {
+ super(expectedSize);
+ }
+ }
+
+ @Nullable
+ private ElementType _elementType;
+ private int _maxByteSize;
+
+ @Override
+ public AggregationFunctionType getAggregationType() {
+ return AggregationFunctionType.ARRAYAGG;
+ }
+
+ @Override
+ public DataType getAggregatedValueType() {
+ return AGGREGATED_VALUE_TYPE;
+ }
+
+ @Override
+ public ObjectSet<Object> getInitialAggregatedValue(@Nullable Object
rawValue) {
+ TrackedSet set = new TrackedSet();
+ if (rawValue != null) {
+ addRawValue(set, rawValue);
+ }
+ updateMaxByteSize(set);
+ return set;
+ }
+
+ @Override
+ public ObjectSet<Object> applyRawValue(ObjectSet<Object> value, Object
rawValue) {
+ TrackedSet set = asTrackedSet(value);
+ if (rawValue != null) {
+ addRawValue(set, rawValue);
+ updateMaxByteSize(set);
+ }
+ return set;
+ }
+
+ @Override
+ public ObjectSet<Object> applyAggregatedValue(ObjectSet<Object> value,
ObjectSet<Object> aggregatedValue) {
+ TrackedSet set = asTrackedSet(value);
+ for (Object element : aggregatedValue) {
+ addElement(set, element);
+ }
+ updateMaxByteSize(set);
+ return set;
+ }
+
+ @Override
+ public ObjectSet<Object> cloneAggregatedValue(ObjectSet<Object> value) {
+ TrackedSet clone = new TrackedSet(value.size());
+ for (Object element : value) {
+ addElement(clone, element);
+ }
+ return clone;
+ }
+
+ @Override
+ public boolean isAggregatedValueFixedSize() {
+ return false;
+ }
+
+ @Override
+ public int getMaxAggregatedValueByteSize() {
+ return _maxByteSize;
+ }
+
+ @Override
+ public byte[] serializeAggregatedValue(ObjectSet<Object> value) {
+ // Layout: [1-byte element type tag][int size][payload]. The tag makes the
cell self-describing so it can be
+ // deserialized (during build re-ingestion) without aggregator instance
state. An all-empty aggregator that never
+ // saw a value defaults to LONG; the payload is empty either way so the
tag choice is immaterial for empty cells.
+ ElementType elementType = _elementType != null ? _elementType :
ElementType.LONG;
+ switch (elementType) {
+ case INT:
+ case LONG:
+ case FLOAT:
+ case DOUBLE:
+ return serializeFixedWidth(value, elementType);
+ case BIG_DECIMAL:
+ case STRING:
+ case BYTES:
+ return serializeVariableWidth(value, elementType);
+ default:
+ throw new IllegalStateException("Unsupported element type: " +
elementType);
+ }
+ }
+
+ @Override
+ public ObjectSet<Object> deserializeAggregatedValue(byte[] bytes) {
+ ByteBuffer byteBuffer = ByteBuffer.wrap(bytes);
+ ElementType elementType = ElementType.fromTag(byteBuffer.get());
+ if (_elementType == null) {
+ // A merge-only flow (e.g. re-merging stored cells) may deserialize
before any raw value has been seen; the tag
+ // pins the element type so subsequent size tracking and serialization
interpret elements correctly.
+ _elementType = elementType;
+ }
+ int size = byteBuffer.getInt();
+ TrackedSet set = new TrackedSet(size);
+ for (int i = 0; i < size; i++) {
+ switch (elementType) {
+ case INT:
+ addElement(set, byteBuffer.getInt());
+ break;
+ case LONG:
+ addElement(set, byteBuffer.getLong());
+ break;
+ case FLOAT:
+ addElement(set, byteBuffer.getFloat());
+ break;
+ case DOUBLE:
+ addElement(set, byteBuffer.getDouble());
+ break;
+ case BIG_DECIMAL: {
+ byte[] valueBytes = readLengthPrefixed(byteBuffer);
+ addElement(set, BigDecimalUtils.deserialize(valueBytes));
+ break;
+ }
+ case STRING: {
+ byte[] valueBytes = readLengthPrefixed(byteBuffer);
+ addElement(set, new String(valueBytes, StandardCharsets.UTF_8));
+ break;
+ }
+ case BYTES: {
+ byte[] valueBytes = readLengthPrefixed(byteBuffer);
+ addElement(set, new ByteArray(valueBytes));
+ break;
+ }
+ default:
+ throw new IllegalStateException("Unsupported element type: " +
elementType);
+ }
+ }
+ return set;
+ }
+
+ private void addRawValue(TrackedSet set, Object rawValue) {
+ // Normalize BYTES to ByteArray so hashCode/equals work as a set key;
other types are already value types.
+ addElement(set, rawValue instanceof byte[] ? new ByteArray((byte[])
rawValue) : rawValue);
+ }
+
+ /// Adds an element to the set, maintaining the running variable-width
payload total when the element is new.
+ private void addElement(TrackedSet set, Object element) {
+ if (_elementType == null) {
+ _elementType = inferElementType(element);
+ }
+ if (set.add(element) && fixedElementBytes(_elementType) == 0) {
+ set._variableWidthPayloadBytes += Integer.BYTES +
variableWidthEncodedLength(element);
+ }
+ }
+
+ /// Returns a [TrackedSet] view of the aggregated value. Values produced by
this aggregator already are one; the
+ /// fallback rebuild keeps the method total for defensively handling a
foreign set.
+ private TrackedSet asTrackedSet(ObjectSet<Object> value) {
Review Comment:
Non-blocking: the foreign-set conversion fallback does not appear to have a
current caller. The builder obtains aggregate values through this aggregator's
initialization, cloning, and deserialization methods, all of which return
`TrackedSet`. Consider relying on that invariant instead of rebuilding
arbitrary `ObjectSet` inputs. Incremental size tracking itself should stay.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]