Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -668,6 +668,12 @@ public static Integer getNumGroupsWarningLimit(Map<String, String> queryOptions)
return checkedParseIntPositive(QueryOptionKey.NUM_GROUPS_WARNING_LIMIT, numGroupsWarningLimit);
}

@Nullable
public static Boolean isGroupByOffHeap(Map<String, String> queryOptions) {
String groupByOffHeap = queryOptions.get(QueryOptionKey.GROUP_BY_OFF_HEAP);
return groupByOffHeap != null ? Boolean.parseBoolean(groupByOffHeap) : null;
}

@Nullable
public static Integer getMaxInitialResultHolderCapacity(Map<String, String> queryOptions) {
String maxInitialResultHolderCapacity = queryOptions.get(QueryOptionKey.MAX_INITIAL_RESULT_HOLDER_CAPACITY);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -329,7 +329,7 @@ private Collection<Record> getUnsortedTopRecords(Map<Key, Record> recordsMap, in
/// This method is to be called from individual segment if the intermediate results need to be trimmed.
public List<IntermediateRecord> sortInSegmentResults(GroupKeyGenerator groupKeyGenerator,
GroupByResultHolder[] groupByResultHolders, int size) {
// getNumKeys() does not count nulls
// NOTE: getNumKeys() counts every group, including the null group when null handling is enabled
assert groupKeyGenerator.getNumKeys() <= size;
Iterator<GroupKeyGenerator.GroupKey> groupKeyIterator = groupKeyGenerator.getGroupKeys();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -106,38 +106,39 @@ protected void processSegments() {
((AcquireReleaseColumnsSegmentOperator) operator).acquire();
}
GroupByResultsBlock resultsBlock = (GroupByResultsBlock) operator.nextBlock();
if (_indexedTable == null) {
synchronized (this) {
if (_indexedTable == null) {
_indexedTable = GroupByUtils.createIndexedTableForCombineOperator(resultsBlock, _queryContext, _numTasks,
_executorService);
// Hold the group-by result so its group key generator (which may own off-heap resources) is always
// released, even when indexed-table creation or the merge below throws
AggregationGroupByResult aggregationGroupByResult = resultsBlock.getAggregationGroupByResult();
try {
if (_indexedTable == null) {
synchronized (this) {
if (_indexedTable == null) {
_indexedTable = GroupByUtils.createIndexedTableForCombineOperator(resultsBlock, _queryContext,
_numTasks, _executorService);
}
}
}
}

if (resultsBlock.isGroupsTrimmed()) {
_groupsTrimmed = true;
}
// Set groups limit reached flag.
if (resultsBlock.isNumGroupsLimitReached()) {
_numGroupsLimitReached = true;
}
if (resultsBlock.isNumGroupsWarningLimitReached()) {
_numGroupsWarningLimitReached = true;
}
if (resultsBlock.isGroupsTrimmed()) {
_groupsTrimmed = true;
}
// Set groups limit reached flag.
if (resultsBlock.isNumGroupsLimitReached()) {
_numGroupsLimitReached = true;
}
if (resultsBlock.isNumGroupsWarningLimitReached()) {
_numGroupsWarningLimitReached = true;
}

// Merge aggregation group-by result.
// Iterate over the group-by keys, for each key, update the group-by result in the indexedTable
Collection<IntermediateRecord> intermediateRecords = resultsBlock.getIntermediateRecords();
// Count the number of merged keys
int mergedKeys = 0;
// For now, only GroupBy OrderBy query has pre-constructed intermediate records
if (intermediateRecords == null) {
// Merge aggregation group-by result.
AggregationGroupByResult aggregationGroupByResult = resultsBlock.getAggregationGroupByResult();
if (aggregationGroupByResult != null) {
// Iterate over the group-by keys, for each key, update the group-by result in the indexedTable
try {
// Iterate over the group-by keys, for each key, update the group-by result in the indexedTable
Collection<IntermediateRecord> intermediateRecords = resultsBlock.getIntermediateRecords();
// Count the number of merged keys
int mergedKeys = 0;
// For now, only GroupBy OrderBy query has pre-constructed intermediate records
if (intermediateRecords == null) {
if (aggregationGroupByResult != null) {
// Iterate over the group-by keys, for each key, update the group-by result in the indexedTable
Iterator<GroupKeyGenerator.GroupKey> dicGroupKeyIterator = aggregationGroupByResult.getGroupKeyIterator();
while (dicGroupKeyIterator.hasNext()) {
QueryThreadContext.checkTerminationAndSampleUsagePeriodically(mergedKeys++, EXPLAIN_NAME);
Expand All @@ -150,16 +151,18 @@ protected void processSegments() {
}
_indexedTable.upsert(new Key(keys), new Record(values));
}
} finally {
// Release the resources used by the group key generator
aggregationGroupByResult.closeGroupKeyGenerator();
}
} else {
for (IntermediateRecord intermediateResult : intermediateRecords) {
QueryThreadContext.checkTerminationAndSampleUsagePeriodically(mergedKeys++, EXPLAIN_NAME);
//TODO: change upsert api so that it accepts intermediateRecord directly
_indexedTable.upsert(intermediateResult._key, intermediateResult._record);
}
}
} else {
for (IntermediateRecord intermediateResult : intermediateRecords) {
QueryThreadContext.checkTerminationAndSampleUsagePeriodically(mergedKeys++, EXPLAIN_NAME);
//TODO: change upsert api so that it accepts intermediateRecord directly
_indexedTable.upsert(intermediateResult._key, intermediateResult._record);
} finally {
if (aggregationGroupByResult != null) {
// Release the resources used by the group key generator
aggregationGroupByResult.closeGroupKeyGenerator();
}
}
} catch (RuntimeException e) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,23 @@ protected GroupByResultsBlock getNextBlock() {
resultHolderIndexMap.put(_aggregationFunctions[i], i);
}

GroupKeyGenerator[] createdGroupKeyGenerator = new GroupKeyGenerator[1];
try {
return processAndBuildResultsBlock(groupByResultHolders, resultHolderIndexMap, createdGroupKeyGenerator);
} catch (Throwable t) {
// Release group-by resources (including off-heap key tables and result holders) that would otherwise leak.
// Close is idempotent on all generators; on the success path the generator is either closed on the trim/sort
// paths below or handed to the combine operator, which closes it after the merge.
if (createdGroupKeyGenerator[0] != null) {
createdGroupKeyGenerator[0].close();
}
throw t;
}
}

private GroupByResultsBlock processAndBuildResultsBlock(GroupByResultHolder[] groupByResultHolders,
IdentityHashMap<AggregationFunction, Integer> resultHolderIndexMap,
GroupKeyGenerator[] createdGroupKeyGenerator) {
GroupKeyGenerator groupKeyGenerator = null;
for (AggregationInfo aggregationInfo : _aggregationInfos) {
AggregationFunction[] aggregationFunctions = aggregationInfo.getFunctions();
Expand All @@ -162,6 +179,7 @@ protected GroupByResultsBlock getNextBlock() {
// GroupByExecutor with a pre-existing GroupKeyGenerator so that the GroupKeyGenerator can be shared across
// loop iterations i.e. across all aggs.
groupKeyGenerator = groupByExecutor.getGroupKeyGenerator();
createdGroupKeyGenerator[0] = groupKeyGenerator;

int numDocsScanned = 0;
ValueBlock valueBlock;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,19 @@ protected GroupByResultsBlock getNextBlock() {
} else {
groupByExecutor = new DefaultGroupByExecutor(_queryContext, _groupByExpressions, _projectOperator);
}
try {
return processAndBuildResultsBlock(groupByExecutor);
} catch (Throwable t) {
// Release group-by resources (including off-heap key tables and result holders) that would otherwise leak.
// On the success path, ownership either ends inside processAndBuildResultsBlock (trim/sort paths close the
// generator there) or moves to the results block consumer (the combine operator closes the generator after
// merging the AggregationGroupByResult). Close is idempotent on all generators.
groupByExecutor.getGroupKeyGenerator().close();
throw t;
}
}

private GroupByResultsBlock processAndBuildResultsBlock(GroupByExecutor groupByExecutor) {
ValueBlock valueBlock;

while ((valueBlock = _projectOperator.nextBlock()) != null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@
import org.apache.pinot.core.plan.StreamingInstanceResponsePlanNode;
import org.apache.pinot.core.plan.StreamingSelectionPlanNode;
import org.apache.pinot.core.query.aggregation.function.AggregationFunction;
import org.apache.pinot.core.query.aggregation.groupby.offheap.OffHeapGroupByBufferPool;
import org.apache.pinot.core.query.executor.ResultsBlockStreamer;
import org.apache.pinot.core.query.prefetch.FetchPlanner;
import org.apache.pinot.core.query.prefetch.FetchPlannerRegistry;
Expand Down Expand Up @@ -117,6 +118,8 @@ public class InstancePlanMakerImplV2 implements PlanMaker {
private int _groupByTrimThreshold = Server.DEFAULT_QUERY_EXECUTOR_GROUPBY_TRIM_THRESHOLD;
private AndRestrictionPushdownMode _andRestrictionPushdownMode =
Server.DEFAULT_QUERY_EXECUTOR_AND_RESTRICTION_PUSHDOWN_MODE;
// Whether to store group-by key tables and fixed-width result holders in off-heap (direct) memory
private boolean _groupByOffHeap = Server.DEFAULT_QUERY_EXECUTOR_GROUPBY_OFF_HEAP;

@Override
public void init(PinotConfiguration queryExecutorConfig) {
Expand Down Expand Up @@ -150,11 +153,16 @@ public void init(PinotConfiguration queryExecutorConfig) {
Server.DEFAULT_QUERY_EXECUTOR_AND_RESTRICTION_PUSHDOWN_MODE.name()));
Preconditions.checkState(_groupByTrimThreshold > 0,
"Invalid configurable: groupByTrimThreshold: %d must be positive", _groupByTrimThreshold);
_groupByOffHeap =
queryExecutorConfig.getProperty(Server.GROUPBY_OFF_HEAP, Server.DEFAULT_QUERY_EXECUTOR_GROUPBY_OFF_HEAP);
OffHeapGroupByBufferPool.setMaxBytesPerThread(
queryExecutorConfig.getProperty(Server.GROUPBY_OFF_HEAP_POOL_MAX_BYTES_PER_THREAD,
Server.DEFAULT_QUERY_EXECUTOR_GROUPBY_OFF_HEAP_POOL_MAX_BYTES_PER_THREAD));
LOGGER.info("Initialized plan maker with maxExecutionThreads: {}, defaultExecutionThreads: {}, "
+ "maxInitialResultHolderCapacity: {}, numGroupsLimit: {}, minSegmentGroupTrimSize: {}, "
+ "minServerGroupTrimSize: {}, groupByTrimThreshold: {}",
+ "minServerGroupTrimSize: {}, groupByTrimThreshold: {}, groupByOffHeap: {}",
_maxExecutionThreads, _defaultExecutionThreads, _maxInitialResultHolderCapacity, _numGroupsLimit,
_minSegmentGroupTrimSize, _minServerGroupTrimSize, _groupByTrimThreshold);
_minSegmentGroupTrimSize, _minServerGroupTrimSize, _groupByTrimThreshold, _groupByOffHeap);
}

@VisibleForTesting
Expand Down Expand Up @@ -207,6 +215,11 @@ public void setGroupByTrimThreshold(int groupByTrimThreshold) {
_groupByTrimThreshold = groupByTrimThreshold;
}

@VisibleForTesting
public void setGroupByOffHeap(boolean groupByOffHeap) {
_groupByOffHeap = groupByOffHeap;
}

@Override
public Plan makeInstancePlan(List<SegmentContext> segmentContexts, QueryContext queryContext,
ExecutorService executorService) {
Expand Down Expand Up @@ -334,6 +347,9 @@ void applyQueryOptions(QueryContext queryContext) {
} else {
queryContext.setNumGroupsLimit(_numGroupsLimit);
}
// Set groupByOffHeap
Boolean groupByOffHeap = QueryOptionsUtils.isGroupByOffHeap(queryOptions);
queryContext.setGroupByOffHeap(groupByOffHeap != null ? groupByOffHeap : _groupByOffHeap);
// Set numGroupsWarningThreshold
queryContext.setNumGroupsWarningLimit(_numGroupsWarningLimit);
// Set minSegmentGroupTrimSize
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,18 @@ public interface AggregationFunction<IntermediateResult, FinalResult extends Com
/// See the null contract on [AggregationFunction].
GroupByResultHolder createGroupByResultHolder(int initialCapacity, int maxCapacity);

/// Returns a group-by result holder that keeps this function's per-group state in off-heap memory, or `null`
/// (the default) when the function has no off-heap state implementation and should keep its on-heap holder.
///
/// Only consulted when off-heap group-by is enabled for the query. A non-null holder must implement
/// [AutoCloseable] — the executor registers it on the query's resource tracker so the existing group key
/// generator close sites release the direct memory — and must be indistinguishable from the on-heap holder
/// through [#extractGroupByResult], including returning `null` for untouched groups.
@Nullable
default GroupByResultHolder createOffHeapGroupByResultHolder(int initialCapacity, int maxCapacity) {
return null;
}

/// Performs aggregation on the given block value sets (aggregation only).
///
/// With null handling enabled, null rows must be skipped rather than folded in as the column's default. See the null
Expand Down
Loading
Loading