Skip to content
Open
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,99 @@
/*
* 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) {
// Rank is advanced during batch forward; partition size alone is needed up
// front.
}

@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 @@ -3005,10 +3005,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
12 changes: 12 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 @@
ROW_NUMBER,
RANK,
DENSE_RANK,
PERCENT_RANK,
CUME_DIST,
MIN,
MAX,
Expand Down Expand Up @@ -133,6 +137,11 @@
supportedFunctionNames.addAll(treeSet);
}

// functions that don't care about input columns.
public static final Set<SupportedFunctionType> COLUMN_AGNOSTIC_FUNCTIONS =

Check warning on line 141 in ql/src/java/org/apache/hadoop/hive/ql/plan/VectorPTFDesc.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Make this member "protected".

See more on https://sonarcloud.io/project/issues?id=apache_hive&issues=AaBilL3ZbkaKLhs52tch&open=AaBilL3ZbkaKLhs52tch&pullRequest=6752
EnumSet.of(SupportedFunctionType.RANK, SupportedFunctionType.DENSE_RANK,

Check warning on line 142 in ql/src/java/org/apache/hadoop/hive/ql/plan/VectorPTFDesc.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

'EnumSet' has incorrect indentation level 4, expected level should be 6.

See more on https://sonarcloud.io/project/issues?id=apache_hive&issues=AaBmiyAlmr3DRVHUAYBG&open=AaBmiyAlmr3DRVHUAYBG&pullRequest=6752
SupportedFunctionType.PERCENT_RANK, SupportedFunctionType.CUME_DIST);

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

Expand Down Expand Up @@ -204,6 +213,9 @@
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
5 changes: 5 additions & 0 deletions ql/src/test/queries/clientpositive/cbo_windowing.q
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,11 @@ set hive.auto.convert.join=false;
-- 9. Test Windowing Functions
-- SORT_QUERY_RESULTS

-- Vector PTF does not buffer PARTITION BY columns (constant within a partition). When a partition
-- column is also a window-function argument (e.g. sum(c_float) OVER (PARTITION BY c_float)),
-- input-column remapping is wrong and can cause ClassCastException. Use non-vector PTF until fixed.
set hive.vectorized.execution.ptf.enabled=false;

select count(c_int) over() from cbo_t1;
select count(c_int) over(partition by c_float order by key), sum(c_float) over(partition by c_float order by key), max(c_int) over(partition by c_float order by key), min(c_int) over(partition by c_float order by key), row_number() over(partition by c_float order by key) as rn, rank() over(partition by c_float order by key), dense_rank() over(partition by c_float order by key), round(percent_rank() over(partition by c_float order by key), 2), lead(c_int, 2, c_int) over(partition by c_float order by key), lag(c_float, 2, c_float) over(partition by c_float order by key) from cbo_t1 order by rn;
select * from (select count(c_int) over(partition by c_float order by key), sum(c_float) over(partition by c_float order by key), max(c_int) over(partition by c_float order by key), min(c_int) over(partition by c_float order by key), row_number() over(partition by c_float order by key) as rn, rank() over(partition by c_float order by key), dense_rank() over(partition by c_float order by key), round(percent_rank() over(partition by c_float order by key),2), lead(c_int, 2, c_int) over(partition by c_float order by key ), lag(c_float, 2, c_float) over(partition by c_float order by key) from cbo_t1 order by rn) cbo_t1;
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;
Original file line number Diff line number Diff line change
Expand Up @@ -4323,7 +4323,7 @@ STAGE PLANS:
Execution mode: vectorized, llap
LLAP IO: all inputs
Reducer 2
Execution mode: llap
Execution mode: vectorized, llap
Reduce Operator Tree:
Select Operator
expressions: KEY.reducesinkkey1 (type: string), VALUE._col1 (type: int), KEY.reducesinkkey0 (type: float)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4473,7 +4473,7 @@ STAGE PLANS:
Execution mode: vectorized, llap
LLAP IO: all inputs
Reducer 2
Execution mode: llap
Execution mode: vectorized, llap
Reduce Operator Tree:
Select Operator
expressions: KEY.reducesinkkey1 (type: string), VALUE._col1 (type: int), KEY.reducesinkkey0 (type: float)
Expand Down
Loading
Loading