deniskuzZ commented on code in PR #6716:
URL: https://github.com/apache/hive/pull/6716#discussion_r4020092620


##########
ql/src/java/org/apache/hadoop/hive/ql/optimizer/StatsOptimizer.java:
##########
@@ -899,39 +590,198 @@ else if (udaf instanceof GenericUDAFCount) {
       }
     }
 
-    private ColumnStatisticsData 
validateSingleColStat(List<ColumnStatisticsObj> statObj) {
-      if (statObj.size() > 1) {
-        Logger.error("More than one stat for a single column!");
-        return null;
-      } else if (statObj.isEmpty()) {
-        Logger.debug("No stats for some partition and column");
-        return null;
-      }
-      return statObj.get(0).getStatsData();
+    /** The columns the aggregates read, which are the ones statistics have to 
be fetched for. */
+    private static List<String> aggregateColumns(GroupByOperator pgbyOp, 
Map<String, ExprNodeDesc> exprMap) {
+      return pgbyOp.getConf().getAggregators().stream()
+          .filter(aggr -> !aggr.getParameters().isEmpty())
+          .map(aggr -> aggr.getParameters().get(0))
+          .filter(ExprNodeColumnDesc.class::isInstance)
+          .map(desc -> exprMap.get(((ExprNodeColumnDesc) desc).getColumn()))
+          .filter(ExprNodeColumnDesc.class::isInstance)
+          .map(desc -> ((ExprNodeColumnDesc) desc).getColumn())
+          .distinct()
+          .collect(Collectors.toList());
     }
 
-    private Collection<List<ColumnStatisticsObj>> verifyAndGetPartColumnStats(
-        Hive hive, Table tbl, String colName, Set<Partition> parts) throws 
TException, LockException {
-      List<String> partNames = new ArrayList<String>(parts.size());
-      for (Partition part : parts) {
-        if 
(!StatsUtils.areColumnStatsUptoDateForQueryAnswering(part.getTable(), 
part.getParameters(), colName)) {
-          Logger.debug("Stats for part : " + part.getSpec() + " column " + 
colName
+    /**
+     * The statistics of the columns a scan's aggregates read, fetched once 
when the first
+     * aggregate needs them and shared by the rest. An aggregate this rewrite 
cannot answer leaves
+     * the query for execution, whole or not at all. Answers for a scan of a 
partitioned table.
+     */
+    private static final class ScanColStats {
+      private final Hive hive;
+      private final Table tbl;
+      private final List<String> colNames;
+      private final PrunedPartitionList prunedList;
+      private Map<String, ColumnStatisticsObj> colStatsByName;
+      private boolean fetched;
+
+      ScanColStats(Hive hive, Table tbl, List<String> colNames, 
PrunedPartitionList prunedList) {
+        this.hive = hive;
+        this.tbl = tbl;
+        this.colNames = colNames;
+        this.prunedList = prunedList;
+      }
+
+      /**
+       * One column's statistics. A scan pruned to no partitions reads no 
rows, and the statistics
+       * of no rows are the empty ones: nothing counted, and no least or 
greatest to name.
+       */
+      ColumnStatisticsData statsFor(String colName, StatType type) throws 
HiveException {
+        if (prunedList != null && prunedList.getPartitions().isEmpty()) {
+          return emptyColStats(type);
+        }
+        if (!fetched) {
+          fetched = true;
+          colStatsByName = prunedList == null ? tableColStats() : 
partitionColStats();
+        }
+        ColumnStatisticsObj stat = colStatsByName == null ? null : 
colStatsByName.get(colName);
+        if (stat == null) {
+          Logger.debug("No stats for " + tbl.getTableName() + " column " + 
colName);
+          return null;
+        }
+        return stat.getStatsData();
+      }
+
+      /**
+       * Whether the table's own statistics answer for this scan: it keeps 
them for the table as
+       * a whole, and the scan reads every partition. They then describe 
exactly the rows read.
+       */
+      private boolean answeredByTableStats() {
+        return !StatsUtils.isPartitionStats(tbl, hive.getConf()) &&
+            prunedList.getReferredPartCols().isEmpty() && 
!prunedList.hasUnknownPartitions();
+      }
+
+      /** The table's own statistics, taken only while they still describe it. 
*/
+      private Map<String, ColumnStatisticsObj> tableColStats() throws 
HiveException {
+        if (!StatsUtils.areColumnStatsUptoDateForQueryAnswering(tbl, 
tbl.getParameters(), colNames)) {
+          Logger.debug("Stats for table : " + tbl.getTableName() + " columns " 
+ colNames
               + " are not up to date.");
           return null;
         }
-        partNames.add(part.getName());
+        return indexByColumnName(hive.getTableColumnStatistics(tbl, colNames, 
true));
+      }
+
+      /** What the scan's partitions hold for every column asked about, or 
null to decline. */
+      private Map<String, ColumnStatisticsObj> partitionColStats() throws 
HiveException {
+        Set<Partition> parts = prunedList.getPartitions();
+        List<String> partNames = new ArrayList<>(parts.size());
+        // a storage handler holds no partition parameters, and one kept per 
partition describes no
+        // partition in particular: whether each still describes itself is 
answered by the aggregate
+        // below, which is told the partitions this query pruned to
+        if (tbl.isNonNative()) {
+          if (!StatsUtils.checkCanProvideColumnStats(tbl)) {
+            Logger.debug("Table : " + tbl.getTableName() + " provides no 
column statistics.");
+            return null;
+          }
+          if (answeredByTableStats()) {
+            return tableColStats();
+          }
+          parts.forEach(part -> partNames.add(part.getName()));
+        } else {
+          for (Partition part : parts) {
+            if (!StatsUtils.areColumnStatsUptoDateForQueryAnswering(
+                part.getTable(), part.getParameters(), colNames)) {
+              Logger.debug("Stats for part : " + part.getSpec() + " columns " 
+ colNames
+                  + " are not up to date.");
+              return null;
+            }
+            partNames.add(part.getName());
+          }
+        }
+        // Aggregated rather than per partition: the callers fold these with 
min, max or a sum, so
+        // merging first gives the same answer. A handler aggregates its own 
statistics, which
+        // the metastore cannot hold: PART_COL_STATS rows need a partition 
Iceberg never creates.
+        AggrStats aggrStats;
+        try {
+          aggrStats = tbl.isNonNative()
+              ? tbl.getStorageHandler().getAggrColStatsFor(tbl, colNames, 
partNames)
+              : exactAggrColStats(partNames);
+        } catch (MetaException e) {
+          throw new HiveException(e);
+        }
+        if (aggrStats == null || aggrStats.getColStats() == null) {
+          Logger.debug("No stats for " + tbl.getTableName() + " columns " + 
colNames);
+          return null;
+        }
+        if (aggrStats.getPartsFound() != parts.size()) {
+          // a partition whose statistics are missing would leave the answer 
describing a subset
+          Logger.debug("Received " + aggrStats.getPartsFound() + " stats for " 
+ parts.size() + " partitions");
+          return null;
+        }
+        return indexByColumnName(aggrStats.getColStats());
+      }
+
+      /**
+       * Each partition fetched and folded the way a storage handler folds its 
own: the
+       * metastore's aggregate endpoint may serve a cached aggregate of a 
different partition
+       * set within its variance, which estimates a plan fine but must not 
answer a query.
+       */
+      private AggrStats exactAggrColStats(List<String> partNames) throws 
HiveException, MetaException {
+        Map<String, List<ColumnStatisticsObj>> statsByPart = 
hive.getPartitionColumnStatistics(
+            tbl.getDbName(), tbl.getTableName(), partNames, colNames, true);
+        List<ColumnStatistics> partStats = new ArrayList<>();
+        statsByPart.forEach((partitionName, statsObjs) -> {
+          // a partition counts as found only when it holds every column asked 
about
+          if (statsObjs.size() == colNames.size()) {
+            ColumnStatisticsDesc statsDesc = new ColumnStatisticsDesc(false, 
tbl.getDbName(), tbl.getTableName());
+            statsDesc.setPartName(partitionName);
+            partStats.add(new ColumnStatistics(statsDesc, statsObjs));
+          }
+        });
+        HiveConf conf = hive.getConf();
+        List<ColumnStatisticsObj> aggregated = 
MetaStoreServerUtils.aggrPartitionStats(partStats,
+            MetaStoreUtils.getDefaultCatalog(conf), tbl.getDbName(), 
tbl.getTableName(),
+            partNames, colNames,
+            partStats.size() == partNames.size(),
+            MetastoreConf.getBoolVar(conf, 
MetastoreConf.ConfVars.STATS_NDV_DENSITY_FUNCTION),
+            MetastoreConf.getDoubleVar(conf, 
MetastoreConf.ConfVars.STATS_NDV_TUNER));
+        return new AggrStats(aggregated, partStats.size());
+      }
+
+      /**
+       * The statistics by the column they describe. A source naming one 
column twice disagrees with
+       * itself: collecting without a merge function throws, and the query 
leaves for execution
+       * rather than an arbitrary one of them standing as an exact answer.
+       */
+      private static Map<String, ColumnStatisticsObj> indexByColumnName(
+          List<ColumnStatisticsObj> colStats) {
+        return colStats.stream().collect(
+            Collectors.toMap(ColumnStatisticsObj::getColName, 
Function.identity()));
+      }
+    }
+
+    /** The rows a COUNT reads, or null to decline - logged. */
+    private Long countFor(AggregationDesc aggr, Map<String, ExprNodeDesc> 
exprMap, long rowCnt,
+        ScanColStats scanColStats) throws HiveException {
+      if (aggr.getParameters().isEmpty()) {
+        // count(*) or count()
+        return rowCnt;
+      }
+      ExprNodeDesc param = aggr.getParameters().get(0);
+      if (param instanceof ExprNodeColumnDesc column) {
+        param = exprMap.get(column.getColumn());
+      }
+      if (param instanceof ExprNodeConstantDesc constant) {
+        // count(1) reads every row, count(null) none
+        return constant.getValue() == null ? 0L : rowCnt;
+      }
+      // count(col): the rows where it is set
+      ExprNodeColumnDesc desc = (ExprNodeColumnDesc) param;
+      String colName = desc.getColumn();
+      StatType type = getType(desc.getTypeString());
+
+      ColumnStatisticsData statData = scanColStats.statsFor(colName, type);
+      if (statData == null) {
+        return null; // logging inside
       }
-      AcidUtils.TableSnapshot tableSnapshot =
-          AcidUtils.getTableSnapshot(hive.getConf(), tbl);
-
-      Map<String, List<ColumnStatisticsObj>> result = 
hive.getMSC().getPartitionColumnStatistics(
-          tbl.getDbName(), tbl.getTableName(), partNames, 
Lists.newArrayList(colName),
-          Constants.HIVE_ENGINE, tableSnapshot != null ? 
tableSnapshot.getValidWriteIdList() : null);
-      if (result.size() != parts.size()) {
-        Logger.debug("Received " + result.size() + " stats for " + 
parts.size() + " partitions");
+      Long nullCnt = getNullCountFor(type, statData);
+      if (nullCnt == null) {

Review Comment:
   master does the identical rowCnt -= nullCnt with no guard



-- 
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]

Reply via email to