Skip to content
Merged
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
@@ -0,0 +1,98 @@
/*
* 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.hadoop.hive.ql.exec.vector.ptf;

import java.util.List;

import org.apache.hadoop.hive.ql.exec.vector.ColumnVector.Type;
import org.apache.hadoop.hive.ql.exec.vector.DoubleColumnVector;
import org.apache.hadoop.hive.ql.exec.vector.VectorizedRowBatch;
import org.apache.hadoop.hive.ql.metadata.HiveException;
import org.apache.hadoop.hive.ql.plan.ptf.WindowFrameDef;

/**
* Evaluates {@code percent_rank()} as a <b>group-aggregated streaming</b>
* evaluator.
*
* <p>
* The partition is buffered so {@link #setPartitionSize(int)} is known before
* output is written.
* Peer-group rank is tracked incrementally (like
* {@link VectorPTFEvaluatorRank}) during batch
* forward; {@link #addStreamingGroupResults} is a no-op because no pre-pass is
* required.
*/
public class VectorPTFEvaluatorPercentRank extends VectorPTFEvaluatorBase {

private int rank;
private int groupCount;

public VectorPTFEvaluatorPercentRank(WindowFrameDef windowFrameDef, int outputColumnNum) {
super(windowFrameDef, outputColumnNum);
resetEvaluator();
}

@Override
public boolean isGroupAggregatedStreamingEvaluator() {
return true;
}

@Override
public void addStreamingGroupResults(List<Integer> groupRowCounts) {
// No-op: percent_rank only needs the partition size and the current rank.
}

@Override
public void evaluateGroupBatch(VectorizedRowBatch batch) throws HiveException {
if (partitionSize <= 0) {
throw new HiveException("Partition size must be set before computing percent_rank");
}
final double divisor = partitionSize > 1 ? partitionSize - 1 : 1;
DoubleColumnVector outputColVector = (DoubleColumnVector) batch.cols[outputColumnNum];
outputColVector.isRepeating = true;
outputColVector.noNulls = true;
outputColVector.isNull[0] = false;
outputColVector.vector[0] = (rank - 1) / divisor;
groupCount += batch.size;
}

@Override
public void doLastBatchWork() {
rank += groupCount;
groupCount = 0;
}

@Override
public boolean streamsResult() {
return true;
}

@Override
public Type getResultColumnVectorType() {
return Type.DOUBLE;
}

@Override
public void resetEvaluator() {
rank = 1;
partitionSize = -1;
groupCount = 0;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2931,6 +2931,12 @@ private boolean validatePTFOperator(PTFOperator op, VectorizationContext vContex
return false;
}

if (hasUnbufferedPartitionColumnInEvaluatorArgs(vectorPTFDesc)) {
setOperatorIssue(
"Window function argument references partition-only column not buffered in vector PTF");
return false;
}

// Output columns ok?
String[] outputColumnNames = vectorPTFDesc.getOutputColumnNames();
TypeInfo[] outputTypeInfos = vectorPTFDesc.getOutputTypeInfos();
Expand Down Expand Up @@ -3005,10 +3011,7 @@ private boolean validatePTFOperator(PTFOperator op, VectorizationContext vContex
throw new RuntimeException("Unexpected window type " + windowFrameDef.getWindowType());
}

// RANK/DENSE_RANK/CUME_DIST don't care about columns.
if (supportedFunctionType != SupportedFunctionType.RANK &&
supportedFunctionType != SupportedFunctionType.DENSE_RANK &&
supportedFunctionType != SupportedFunctionType.CUME_DIST) {
if (!VectorPTFDesc.COLUMN_AGNOSTIC_FUNCTIONS.contains(supportedFunctionType)) {

if (exprNodeDescList != null) {
// LEAD and LAG now supports multiple arguments in vectorized mode
Expand Down Expand Up @@ -5039,6 +5042,70 @@ private static ExprNodeDesc[] getOrderExprNodeDescs(List<OrderExpressionDef> ord
return exprNodeDescs;
}

// TODO: An evaluator that wants to handle an unbuffered partition-only column in its calculation could
// opt in to vectorization here.
private static boolean hasUnbufferedPartitionColumnInEvaluatorArgs(
VectorPTFDesc vectorPTFDesc) {

// PARTITION BY matches ORDER BY, so partition cols are buffered as order cols.
if (!vectorPTFDesc.getIsPartitionOrderBy()) {
return false;
}

List<ExprNodeDesc> partitionOnlyExprs = getPartitionOnlyExprs(vectorPTFDesc.getPartitionExprNodeDescs(),
vectorPTFDesc.getOrderExprNodeDescs());
if (partitionOnlyExprs.isEmpty()) {
return false;
}

return evaluatorArgsReferencePartitionOnlyExprs(
vectorPTFDesc.getEvaluatorFunctionNames(), vectorPTFDesc.getEvaluatorInputExprNodeDescLists(),
partitionOnlyExprs);
}

private static boolean evaluatorArgsReferencePartitionOnlyExprs(
String[] evaluatorFunctionNames,
List<ExprNodeDesc>[] evaluatorInputExprNodeDescLists,
List<ExprNodeDesc> partitionOnlyExprs) {
for (int i = 0; i < evaluatorFunctionNames.length; i++) {
SupportedFunctionType supportedFunctionType =
VectorPTFDesc.supportedFunctionsMap.get(evaluatorFunctionNames[i].toLowerCase());
List<ExprNodeDesc> exprNodeDescList = evaluatorInputExprNodeDescLists[i];
if (supportedFunctionType == null ||
VectorPTFDesc.COLUMN_AGNOSTIC_FUNCTIONS.contains(supportedFunctionType) ||
Comment thread
deniskuzZ marked this conversation as resolved.
exprNodeDescList == null) {
continue;
}

// Check whether an evaluator argument references a partition-only column.
if (exprNodeDescList.stream()
.anyMatch(expr -> hasPartitionOnlyColumnArg(expr, partitionOnlyExprs))) {
return true;
}
}
return false;
}

private static boolean hasPartitionOnlyColumnArg(
ExprNodeDesc expr, List<ExprNodeDesc> partitionOnlyExprs) {
if (!(expr instanceof ExprNodeColumnDesc)) {
return false;
}
return partitionOnlyExprs.stream().anyMatch(expr::isSame);
}

private static List<ExprNodeDesc> getPartitionOnlyExprs(
Comment thread
deniskuzZ marked this conversation as resolved.
ExprNodeDesc[] partitionExprNodeDescs, ExprNodeDesc[] orderExprNodeDescs) {
List<ExprNodeDesc> partitionOnlyExprs = new ArrayList<ExprNodeDesc>();
for (ExprNodeDesc partitionExpr : partitionExprNodeDescs) {
// Collect partition expressions that are not also ORDER BY expressions.
if (Arrays.stream(orderExprNodeDescs).noneMatch(partitionExpr::isSame)) {
partitionOnlyExprs.add(partitionExpr);
}
}
return partitionOnlyExprs;
}

/*
* Update the VectorPTFDesc with data that is used during validation and that doesn't rely on
* VectorizationContext to lookup column names, etc.
Expand Down
15 changes: 15 additions & 0 deletions ql/src/java/org/apache/hadoop/hive/ql/plan/VectorPTFDesc.java
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,10 @@
package org.apache.hadoop.hive.ql.plan;

import java.util.ArrayList;
import java.util.EnumSet;
import java.util.HashMap;
import java.util.List;
import java.util.Set;
import java.util.TreeSet;

import org.apache.commons.lang3.ArrayUtils;
Expand Down Expand Up @@ -58,6 +60,7 @@
import org.apache.hadoop.hive.ql.exec.vector.ptf.VectorPTFEvaluatorLongMax;
import org.apache.hadoop.hive.ql.exec.vector.ptf.VectorPTFEvaluatorLongMin;
import org.apache.hadoop.hive.ql.exec.vector.ptf.VectorPTFEvaluatorLongSum;
import org.apache.hadoop.hive.ql.exec.vector.ptf.VectorPTFEvaluatorPercentRank;
import org.apache.hadoop.hive.ql.exec.vector.ptf.VectorPTFEvaluatorRank;
import org.apache.hadoop.hive.ql.exec.vector.ptf.VectorPTFEvaluatorRowNumber;
import org.apache.hadoop.hive.ql.exec.vector.ptf.VectorPTFEvaluatorStreamingDecimalAvg;
Expand Down Expand Up @@ -93,6 +96,7 @@ public enum SupportedFunctionType {
ROW_NUMBER,
RANK,
DENSE_RANK,
PERCENT_RANK,
CUME_DIST,
MIN,
MAX,
Expand Down Expand Up @@ -133,6 +137,14 @@ public boolean isSupportDistinct() {
supportedFunctionNames.addAll(treeSet);
}

// Functions that do not depend on input columns.
public static final Set<SupportedFunctionType> COLUMN_AGNOSTIC_FUNCTIONS =
EnumSet.of(
SupportedFunctionType.RANK,
SupportedFunctionType.DENSE_RANK,
SupportedFunctionType.PERCENT_RANK,
SupportedFunctionType.CUME_DIST);

private TypeInfo[] reducerBatchTypeInfos;
private DataTypePhysicalVariation[] reducerBatchDataTypePhysicalVariations;

Expand Down Expand Up @@ -204,6 +216,9 @@ public static VectorPTFEvaluatorBase getEvaluator(SupportedFunctionType function
case DENSE_RANK:
evaluator = new VectorPTFEvaluatorDenseRank(windowFrameDef, outputColumnNum);
break;
case PERCENT_RANK:
evaluator = new VectorPTFEvaluatorPercentRank(windowFrameDef, outputColumnNum);
break;
case CUME_DIST:
evaluator = new VectorPTFEvaluatorCumeDist(windowFrameDef, outputColumnNum);
break;
Expand Down
113 changes: 113 additions & 0 deletions ql/src/test/queries/clientpositive/vector_ptf_percent_rank.q
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
set hive.vectorized.testing.reducer.batch.size=2;

DROP TABLE IF EXISTS vector_ptf_percent_rank_int;

CREATE TABLE vector_ptf_percent_rank_int(name string, rowindex int, mynumber int) stored as orc;

INSERT INTO vector_ptf_percent_rank_int values
('five', 1, 10),
('five', 2, 20),
('five', 3, 30),
('five', 4, 40),
('five', 5, 50),
('six', 1, 10),
('six', 2, 20),
('six', 3, 30),
('six', 4, 40),
('six', 5, 50),
('six', 6, 60),
-- single-row partition: percent_rank 0.0
('lonely', 99, 42),
-- two-row null partition
(NULL, 1, 100),
(NULL, 2, 100);

-- NON-VECTORIZED
set hive.vectorized.execution.ptf.enabled=false;

select name, rowindex, mynumber,
percent_rank() over (partition by name order by mynumber) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over (partition by name order by mynumber) as r,
dense_rank() over (partition by name order by mynumber) as dr,
percent_rank() over (partition by name order by mynumber) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over (order by mynumber) as r,
dense_rank() over (order by mynumber) as dr,
percent_rank() over (order by mynumber) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over (partition by name) as r,
dense_rank() over (partition by name) as dr,
percent_rank() over (partition by name) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over () as r,
dense_rank() over () as dr,
percent_rank() over () as pr
from vector_ptf_percent_rank_int;

-- VECTORIZED
set hive.vectorized.execution.ptf.enabled=true;

explain vectorization detail select name, rowindex, mynumber,
percent_rank() over (partition by name order by mynumber) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
percent_rank() over (partition by name order by mynumber) as pr
from vector_ptf_percent_rank_int;

explain vectorization detail select name, rowindex, mynumber,
rank() over (partition by name order by mynumber) as r,
dense_rank() over (partition by name order by mynumber) as dr,
percent_rank() over (partition by name order by mynumber) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over (partition by name order by mynumber) as r,
dense_rank() over (partition by name order by mynumber) as dr,
percent_rank() over (partition by name order by mynumber) as pr
from vector_ptf_percent_rank_int;

explain vectorization detail select name, rowindex, mynumber,
rank() over (order by mynumber) as r,
dense_rank() over (order by mynumber) as dr,
percent_rank() over (order by mynumber) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over (order by mynumber) as r,
dense_rank() over (order by mynumber) as dr,
percent_rank() over (order by mynumber) as pr
from vector_ptf_percent_rank_int;

explain vectorization detail select name, rowindex, mynumber,
rank() over (partition by name) as r,
dense_rank() over (partition by name) as dr,
percent_rank() over (partition by name) as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over (partition by name) as r,
dense_rank() over (partition by name) as dr,
percent_rank() over (partition by name) as pr
from vector_ptf_percent_rank_int;

explain vectorization detail select name, rowindex, mynumber,
rank() over () as r,
dense_rank() over () as dr,
percent_rank() over () as pr
from vector_ptf_percent_rank_int;

select name, rowindex, mynumber,
rank() over () as r,
dense_rank() over () as dr,
percent_rank() over () as pr
from vector_ptf_percent_rank_int;
Loading
Loading