This is an automated email from the ASF dual-hosted git repository.

soumyakantidas pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hive.git


The following commit(s) were added to refs/heads/master by this push:
     new 8086a9e4055 HIVE-28655: Implement HMS Related Drop Stats Changes 
(Part2) COL_STATS_ACCURATE related changes (#6198)
8086a9e4055 is described below

commit 8086a9e40554e2c834f51178d4cd1ce35a524438
Author: Daniel (Hongdan) Zhu <[email protected]>
AuthorDate: Thu Mar 26 13:21:33 2026 -0700

    HIVE-28655: Implement HMS Related Drop Stats Changes (Part2) 
COL_STATS_ACCURATE related changes (#6198)
    
    Co-authored-by: zdeng <[email protected]>
---
 .../apache/hadoop/hive/common/StatsSetupConst.java |   8 +-
 .../hadoop/hive/metastore/utils/StringUtils.java   |  17 ++
 .../hadoop/hive/metastore/DatabaseProduct.java     |   8 +-
 .../apache/hadoop/hive/metastore/HMSHandler.java   |  24 +-
 .../apache/hadoop/hive/metastore/ObjectStore.java  |  91 +++++---
 .../metastore/PartitionProjectionEvaluator.java    |   5 +-
 .../hadoop/hive/metastore/StatObjectConverter.java |   1 +
 .../hadoop/hive/metastore/cache/CachedStore.java   |  23 +-
 .../{ => directsql}/DirectSqlAggrStats.java        |  25 ++-
 .../metastore/{ => directsql}/DirectSqlBase.java   |   3 +-
 .../metastore/directsql/DirectSqlDeleteStats.java  | 243 +++++++++++++++++++++
 .../{ => directsql}/DirectSqlInsertPart.java       |   8 +-
 .../{ => directsql}/DirectSqlUpdateParams.java     |  13 +-
 .../{ => directsql}/DirectSqlUpdatePart.java       |  30 ++-
 .../{ => directsql}/MetaStoreDirectSql.java        |  91 +++-----
 .../{ => directsql}/MetastoreDirectSqlUtils.java   |  58 +++--
 .../hadoop/hive/metastore/DummyCustomRDBMS.java    |   2 +-
 .../hadoop/hive/metastore/TestHiveMetaStore.java   |  16 ++
 .../hadoop/hive/metastore/TestObjectStore.java     | 110 ++++++++++
 .../hive/metastore/VerifyingObjectStore.java       |   8 +-
 20 files changed, 601 insertions(+), 183 deletions(-)

diff --git 
a/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/common/StatsSetupConst.java
 
b/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/common/StatsSetupConst.java
index f2c6ef27f57..5a7d85107f5 100644
--- 
a/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/common/StatsSetupConst.java
+++ 
b/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/common/StatsSetupConst.java
@@ -374,8 +374,12 @@ public static void removeColumnStatsState(Map<String, 
String> params, List<Strin
       return;
     }
     ColumnStatsAccurate stats = 
parseStatsAcc(params.get(COLUMN_STATS_ACCURATE));
-    colNames.forEach(colName ->
-        stats.columnStats.remove(colName.toLowerCase()));
+    colNames.forEach(colName -> {
+      // it is possible that colName could be null
+      if (colName != null) {
+        stats.columnStats.remove(colName.toLowerCase());
+      }
+    });
 
     try {      
       params.put(COLUMN_STATS_ACCURATE, 
ColumnStatsAccurate.objectWriter.writeValueAsString(stats));
diff --git 
a/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/metastore/utils/StringUtils.java
 
b/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/metastore/utils/StringUtils.java
index 61258acbef7..f3af724cf26 100644
--- 
a/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/metastore/utils/StringUtils.java
+++ 
b/standalone-metastore/metastore-common/src/main/java/org/apache/hadoop/hive/metastore/utils/StringUtils.java
@@ -27,6 +27,8 @@
 import java.util.Map;
 import java.util.Set;
 
+import org.apache.hadoop.hive.metastore.api.MetaException;
+
 public class StringUtils {
 
   /**
@@ -94,6 +96,21 @@ public static String normalizeIdentifier(String identifier) {
     return identifier.trim().toLowerCase();
   }
 
+  public static List<String> normalizeIdentifiers(List<String> identifiers)
+      throws MetaException {
+    if (identifiers != null && !identifiers.isEmpty()) {
+      List<String> normalizedIdents = new ArrayList<>();
+      for (String ident : identifiers) {
+        if (StringUtils.isEmpty(ident)) {
+          throw new MetaException("An unexpected empty identifier found in the 
list");
+        }
+        normalizedIdents.add(normalizeIdentifier(ident));
+      }
+      return normalizedIdents;
+    }
+    return Collections.emptyList();
+  }
+
   /**
    * Make a string representation of the exception.
    * @param e The exception to stringify
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DatabaseProduct.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DatabaseProduct.java
index 686c1e9c371..cd78a18bf53 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DatabaseProduct.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DatabaseProduct.java
@@ -282,14 +282,14 @@ public String toVarChar(String column) {
   /**
    * Whether the RDBMS has restrictions on IN list size (explicit, or poor 
perf-based).
    */
-  protected boolean needsInBatching() {
+  public boolean needsInBatching() {
     return isORACLE() || isSQLSERVER();
   }
 
   /**
    * Whether the RDBMS has a bug in join and filter operation order described 
in DERBY-6358.
    */
-  protected boolean hasJoinOperationOrderBug() {
+  public boolean hasJoinOperationOrderBug() {
     return isDERBY() || isORACLE() || isPOSTGRES();
   }
 
@@ -314,7 +314,7 @@ public static void reset() {
     theDatabaseProduct = null;
   }
 
-  protected String toDate(String tableValue) {
+  public String toDate(String tableValue) {
     if (isORACLE()) {
       return "TO_DATE(" + tableValue + ", 'YYYY-MM-DD')";
     } else {
@@ -322,7 +322,7 @@ protected String toDate(String tableValue) {
     }
   }
 
-  protected String toTimestamp(String tableValue) {
+  public String toTimestamp(String tableValue) {
     if (isORACLE()) {
       return "TO_TIMESTAMP(" + tableValue + ", 'YYYY-MM-DD HH24:mi:ss')";
     } else if (isSQLSERVER()) {
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/HMSHandler.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/HMSHandler.java
index b7801e08e77..9838d9d2573 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/HMSHandler.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/HMSHandler.java
@@ -4572,14 +4572,17 @@ public boolean 
delete_partition_column_statistics(String dbName, String tableNam
     dbName = dbName.toLowerCase();
     String[] parsedDbName = parseDbName(dbName, conf);
     tableName = tableName.toLowerCase();
+    DeleteColumnStatisticsRequest request = new 
DeleteColumnStatisticsRequest(parsedDbName[DB_NAME], tableName);
     if (colName != null) {
       colName = colName.toLowerCase();
+      request.addToCol_names(colName);
     }
-    DeleteColumnStatisticsRequest request = new 
DeleteColumnStatisticsRequest(parsedDbName[DB_NAME], tableName);
+
     request.setEngine(engine);
     request.setCat_name(parsedDbName[CAT_NAME]);
-    request.addToCol_names(colName);
-    request.addToPart_names(partName);
+    if (partName != null) {
+      request.addToPart_names(partName);
+    }
     return delete_column_statistics_req(request);
   }
 
@@ -4611,8 +4614,8 @@ public boolean 
delete_column_statistics_req(DeleteColumnStatisticsRequest req) t
         ret = rawStore.deleteTableColumnStatistics(parsedDbName[CAT_NAME], 
parsedDbName[DB_NAME], tableName, colNames, engine);
         if (ret) {
           eventType = EventType.DELETE_TABLE_COLUMN_STAT;
-          for (String colName :
-              colNames == null ? 
table.getSd().getCols().stream().map(FieldSchema::getName).collect(Collectors.toList())
 : colNames) {
+          for (String colName : colNames == null || colNames.isEmpty() ?
+              
table.getSd().getCols().stream().map(FieldSchema::getName).toList() : colNames) 
{
             if (transactionalListeners != null && 
!transactionalListeners.isEmpty()) {
               MetaStoreListenerNotifier.notifyEvent(transactionalListeners, 
eventType,
                   new DeleteTableColumnStatEvent(parsedDbName[CAT_NAME], 
parsedDbName[DB_NAME], tableName, colName, engine, this));
@@ -4635,8 +4638,8 @@ public boolean 
delete_column_statistics_req(DeleteColumnStatisticsRequest req) t
                 partNames, colNames, engine);
         if (ret) {
           eventType = EventType.DELETE_PARTITION_COLUMN_STAT;
-          for (String colName : colNames == null ? 
table.getSd().getCols().stream().map(FieldSchema::getName)
-              .collect(Collectors.toList()) : colNames) {
+          for (String colName : colNames == null || colNames.isEmpty() ?
+              
table.getSd().getCols().stream().map(FieldSchema::getName).toList() : colNames) 
{
             for (String partName : partNames) {
               List<String> partVals = getPartValsFromName(table, partName);
               if (transactionalListeners != null && 
!transactionalListeners.isEmpty()) {
@@ -4671,17 +4674,14 @@ public boolean delete_table_column_statistics(String 
dbName, String tableName, S
       throws TException {
     dbName = dbName.toLowerCase();
     tableName = tableName.toLowerCase();
-
     String[] parsedDbName = parseDbName(dbName, conf);
-
+    DeleteColumnStatisticsRequest request = new 
DeleteColumnStatisticsRequest(parsedDbName[DB_NAME], tableName);
     if (colName != null) {
       colName = colName.toLowerCase();
+      request.addToCol_names(colName);
     }
-
-    DeleteColumnStatisticsRequest request = new 
DeleteColumnStatisticsRequest(parsedDbName[DB_NAME], tableName);
     request.setEngine(engine);
     request.setCat_name(parsedDbName[CAT_NAME]);
-    request.addToCol_names(colName);
     request.setTableLevel(true);
     return delete_column_statistics_req(request);
   }
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/ObjectStore.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/ObjectStore.java
index 340c5c35b13..447d1a6a2c2 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/ObjectStore.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/ObjectStore.java
@@ -24,6 +24,7 @@
 import static 
org.apache.hadoop.hive.metastore.utils.MetaStoreUtils.getDefaultCatalog;
 import static 
org.apache.hadoop.hive.metastore.utils.MetaStoreUtils.newMetaException;
 import static 
org.apache.hadoop.hive.metastore.utils.StringUtils.normalizeIdentifier;
+import static 
org.apache.hadoop.hive.metastore.utils.StringUtils.normalizeIdentifiers;
 
 import java.io.IOException;
 import java.net.InetAddress;
@@ -75,8 +76,8 @@
 
 import org.apache.commons.collections4.CollectionUtils;
 import org.apache.commons.lang3.ArrayUtils;
-import org.apache.commons.lang3.ObjectUtils;
 import org.apache.commons.lang3.StringUtils;
+import org.apache.commons.lang3.ObjectUtils;
 import org.apache.commons.lang3.exception.ExceptionUtils;
 import org.apache.commons.lang3.tuple.Pair;
 import org.apache.hadoop.classification.InterfaceAudience;
@@ -88,7 +89,9 @@
 import org.apache.hadoop.hive.common.TableName;
 import org.apache.hadoop.hive.common.ValidReaderWriteIdList;
 import org.apache.hadoop.hive.common.ValidWriteIdList;
-import 
org.apache.hadoop.hive.metastore.MetaStoreDirectSql.SqlFilterForPushdown;
+import org.apache.hadoop.hive.metastore.directsql.DirectSqlDeleteStats;
+import org.apache.hadoop.hive.metastore.directsql.MetaStoreDirectSql;
+import 
org.apache.hadoop.hive.metastore.directsql.MetaStoreDirectSql.SqlFilterForPushdown;
 import org.apache.hadoop.hive.metastore.api.AggrStats;
 import org.apache.hadoop.hive.metastore.api.AllTableConstraintsRequest;
 import org.apache.hadoop.hive.metastore.api.AlreadyExistsException;
@@ -199,6 +202,7 @@
 import org.apache.hadoop.hive.metastore.client.builder.GetPartitionsArgs;
 import org.apache.hadoop.hive.metastore.conf.MetastoreConf;
 import org.apache.hadoop.hive.metastore.conf.MetastoreConf.ConfVars;
+import org.apache.hadoop.hive.metastore.directsql.DirectSqlAggrStats;
 import org.apache.hadoop.hive.metastore.metrics.Metrics;
 import org.apache.hadoop.hive.metastore.metrics.MetricsConstants;
 import org.apache.hadoop.hive.metastore.model.FetchGroups;
@@ -2134,14 +2138,7 @@ public List<Table> getTableObjectsByName(String catName, 
String db, List<String>
       openTransaction();
       catName = normalizeIdentifier(catName);
 
-      List<String> lowered_tbl_names = new ArrayList<>();
-      if(tbl_names != null) {
-        lowered_tbl_names = new ArrayList<>(tbl_names.size());
-        for (String t : tbl_names) {
-          lowered_tbl_names.add(normalizeIdentifier(t));
-        }
-      }
-
+      List<String> lowered_tbl_names = normalizeIdentifiers(tbl_names);
       StringBuilder filterBuilder = new StringBuilder();
       List<String> parameterVals = new ArrayList<>();
       appendPatternCondition(filterBuilder, "database.name", db, 
parameterVals);
@@ -3026,7 +3023,7 @@ private void dropPartitionsViaJdo(String catName, String 
dbName, String tblName,
       return;
     }
     openTransaction();
-    
+
     int batch = batchSize == NO_BATCHING ? 1 : (partNames.size() + batchSize) 
/ batchSize;
     AtomicLong batchIdx = new AtomicLong(1);
     AtomicLong timeSpent = new AtomicLong(0);
@@ -10016,6 +10013,7 @@ public boolean deletePartitionColumnStatistics(String 
catName, String dbName, St
     }
     dbName = org.apache.commons.lang3.StringUtils.defaultString(dbName, 
Warehouse.DEFAULT_DATABASE_NAME);
     catName = normalizeIdentifier(catName);
+    List<String> cols = normalizeIdentifiers(colNames);
     return new GetHelper<Boolean>(catName, dbName, tableName, true, true) {
       @Override
       protected String describeResult() {
@@ -10023,12 +10021,13 @@ protected String describeResult() {
       }
       @Override
       protected Boolean getSqlResult(GetHelper<Boolean> ctx) throws 
MetaException {
-        return directSql.deletePartitionColumnStats(catName, dbName, 
tableName, partNames, colNames, engine);
+        DirectSqlDeleteStats deleteStats = new DirectSqlDeleteStats(directSql, 
pm);
+        return deleteStats.deletePartitionColumnStats(catName, dbName, 
tableName, partNames, cols, engine);
       }
       @Override
       protected Boolean getJdoResult(GetHelper<Boolean> ctx)
               throws MetaException, NoSuchObjectException, 
InvalidObjectException, InvalidInputException {
-        return deletePartitionColumnStatisticsViaJdo(catName, dbName, 
tableName, partNames, colNames, engine);
+        return deletePartitionColumnStatisticsViaJdo(catName, dbName, 
tableName, partNames, cols, engine);
       }
     }.run(false);
   }
@@ -10066,12 +10065,7 @@ public List<Void> run(List<String> input) throws 
Exception {
           params.add(normalizeIdentifier(database));
           params.add(normalizeIdentifier(tableName));
           if (colNames != null && !colNames.isEmpty()) {
-            List<String> normalizedColNames = new ArrayList<>();
-            for (String colName : colNames){
-              // trim the extra spaces, and change to lowercase
-              normalizedColNames.add(normalizeIdentifier(colName));
-            }
-            params.add(normalizedColNames);
+            params.add(colNames);
           }
           params.add(catalog);
           if (engine != null) {
@@ -10082,9 +10076,6 @@ public List<Void> run(List<String> input) throws 
Exception {
           pm.retrieveAll(mStatsObjColl);
           if (mStatsObjColl != null) {
             pm.deletePersistentAll(mStatsObjColl);
-          } else {
-            throw new NoSuchObjectException("partition stats doesn't exist for 
db=" + dbName + " table="
-                + tableName + " col=" + String.join(", ", colNames) + " 
partNames=" + String.join(", ", input));
           }
           return null;
         }
@@ -10094,6 +10085,31 @@ public List<Void> run(List<String> input) throws 
Exception {
       } finally {
         b.closeAllQueries();
       }
+
+      Batchable.runBatched(batchSize, partNames, new Batchable<String, Void>() 
{
+        @Override
+        public List<Void> run(List<String> input) throws MetaException {
+          Pair<Query, Map<String, String>> queryWithParams = 
getPartQueryWithParams(catalog, database, tableName,
+              input);
+          try (QueryWrapper qw = new QueryWrapper(queryWithParams.getLeft())) {
+            qw.setResultClass(MPartition.class);
+            qw.setClass(MPartition.class);
+            List<MPartition> mparts = (List<MPartition>) 
qw.executeWithMap(queryWithParams.getRight());
+            for (MPartition mPart : mparts) {
+              Map<String, String> params = mPart.getParameters();
+              if (params != null && 
params.containsKey(StatsSetupConst.COLUMN_STATS_ACCURATE)) {
+                if (colNames == null || colNames.isEmpty()) {
+                  StatsSetupConst.clearColumnStatsState(params);
+                } else {
+                  StatsSetupConst.removeColumnStatsState(params, colNames);
+                }
+                mPart.setParameters(params);
+              }
+            }
+          }
+          return Collections.emptyList();
+        }
+      });
       ret = commitTransaction();
     } finally {
       rollbackAndCleanup(ret, null);
@@ -10109,6 +10125,7 @@ public boolean deleteTableColumnStatistics(String 
catName, String dbName, String
     if (tableName == null) {
       throw new InvalidInputException("Table name is null.");
     }
+    List<String> cols = normalizeIdentifiers(colNames);
     return new GetHelper<Boolean>(catName, dbName, tableName, true, true) {
       @Override
       protected String describeResult() {
@@ -10116,12 +10133,13 @@ protected String describeResult() {
       }
       @Override
       protected Boolean getSqlResult(GetHelper<Boolean> ctx) throws 
MetaException {
-        return directSql.deleteTableColumnStatistics(getTable().getId(), 
colNames, engine);
+        DirectSqlDeleteStats deleteStats = new DirectSqlDeleteStats(directSql, 
pm);
+        return deleteStats.deleteTableColumnStatistics(getTable(), cols, 
engine);
       }
       @Override
       protected Boolean getJdoResult(GetHelper<Boolean> ctx)
               throws MetaException, NoSuchObjectException, 
InvalidObjectException, InvalidInputException {
-        return deleteTableColumnStatisticsViaJdo(catName, dbName, tableName, 
colNames, engine);
+        return deleteTableColumnStatisticsViaJdo(catName, dbName, tableName, 
cols, engine);
       }
     }.run(true);
   }
@@ -10151,14 +10169,9 @@ private boolean 
deleteTableColumnStatisticsViaJdo(String catName, String dbName,
       List<Object> params = new ArrayList<>();
       params.add(normalizeIdentifier(tableName));
       params.add(normalizeIdentifier(dbName));
-      params.add(normalizeIdentifier(catName));
+      params.add(catName == null ? null : normalizeIdentifier(catName));
       if (colNames != null && !colNames.isEmpty()) {
-        List<String> normalizedColNames = new ArrayList<>();
-        for (String colName : colNames){
-          // trim the extra spaces, and change to lowercase
-          normalizedColNames.add(normalizeIdentifier(colName));
-        }
-        params.add(normalizedColNames);
+        params.add(colNames);
       }
       if (engine != null) {
         params.add(engine);
@@ -10167,9 +10180,19 @@ private boolean 
deleteTableColumnStatisticsViaJdo(String catName, String dbName,
       pm.retrieveAll(mStatsObjColl);
       if (mStatsObjColl != null) {
         pm.deletePersistentAll(mStatsObjColl);
-      } else {
-        throw new NoSuchObjectException("Column stats doesn't exist for db=" + 
dbName + " table="
-            + tableName + " col=" + String.join(", ", colNames));
+      }
+
+      MTable mTable = getMTable(catName, dbName, tableName);
+      if (mTable != null) {
+        Map<String, String> tableParams = mTable.getParameters();
+        if (tableParams != null && 
tableParams.containsKey(StatsSetupConst.COLUMN_STATS_ACCURATE)) {
+          if (colNames == null || colNames.isEmpty()) {
+            StatsSetupConst.clearColumnStatsState(tableParams);
+          } else {
+            StatsSetupConst.removeColumnStatsState(tableParams, colNames);
+          }
+          mTable.setParameters(tableParams);
+        }
       }
       ret = commitTransaction();
     } finally {
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/PartitionProjectionEvaluator.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/PartitionProjectionEvaluator.java
index b1a1328c18b..72f69ce019e 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/PartitionProjectionEvaluator.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/PartitionProjectionEvaluator.java
@@ -23,7 +23,6 @@
 import com.fasterxml.jackson.databind.node.ArrayNode;
 import com.google.common.annotations.VisibleForTesting;
 import com.google.common.base.Joiner;
-import com.google.common.collect.ImmutableBiMap;
 import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.ImmutableSet;
 
@@ -32,7 +31,7 @@
 import org.apache.hadoop.hive.metastore.api.Partition;
 import org.apache.hadoop.hive.metastore.api.SerDeInfo;
 import org.apache.hadoop.hive.metastore.api.StorageDescriptor;
-import org.apache.hadoop.hive.metastore.api.Table;
+import org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils;
 import org.apache.hadoop.hive.metastore.utils.MetaStoreServerUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -54,7 +53,7 @@
 import java.util.TreeMap;
 import java.util.regex.Pattern;
 
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.extractSqlLong;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.extractSqlLong;
 
 /**
  * Evaluator for partition projection filters which specify parts of the 
partition that should be
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/StatObjectConverter.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/StatObjectConverter.java
index 24237d75829..e55d4ebcdf9 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/StatObjectConverter.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/StatObjectConverter.java
@@ -49,6 +49,7 @@
 import 
org.apache.hadoop.hive.metastore.columnstats.cache.LongColumnStatsDataInspector;
 import 
org.apache.hadoop.hive.metastore.columnstats.cache.StringColumnStatsDataInspector;
 import 
org.apache.hadoop.hive.metastore.columnstats.cache.TimestampColumnStatsDataInspector;
+import org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils;
 import org.apache.hadoop.hive.metastore.model.MPartition;
 import org.apache.hadoop.hive.metastore.model.MPartitionColumnStatistics;
 import org.apache.hadoop.hive.metastore.model.MTable;
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/cache/CachedStore.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/cache/CachedStore.java
index ea3523051fa..7df41ec68a3 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/cache/CachedStore.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/cache/CachedStore.java
@@ -2169,15 +2169,24 @@ private static void 
updateTableColumnsStatsInternal(Configuration conf, ColumnSt
       throw new RuntimeException("CachedStore can only be enabled for Hive 
engine");
     }
     boolean succ = rawStore.deleteTableColumnStatistics(catName, dbName, 
tblName, colNames, engine);
+    catName = normalizeIdentifier(catName);
+    dbName = normalizeIdentifier(dbName);
+    tblName = normalizeIdentifier(tblName);
+    if (!shouldCacheTable(catName, dbName, tblName)) {
+      return succ;
+    }
+    Table cachedTable = sharedCache.getTableFromCache(catName, dbName, 
tblName);
+    if (cachedTable != null &&
+        
cachedTable.getParameters().containsKey(StatsSetupConst.COLUMN_STATS_ACCURATE) 
&& succ) {
+      if (colNames == null || colNames.isEmpty()) {
+        StatsSetupConst.clearColumnStatsState(cachedTable.getParameters());
+      } else {
+        StatsSetupConst.removeColumnStatsState(cachedTable.getParameters(), 
colNames);
+      }
+      sharedCache.alterTableInCache(catName, dbName, tblName, cachedTable);
+    }
     // in case of event based cache update, cache is updated during commit txn
     if (succ && !canUseEvents) {
-      catName = normalizeIdentifier(catName);
-      dbName = normalizeIdentifier(dbName);
-      tblName = normalizeIdentifier(tblName);
-      if (!shouldCacheTable(catName, dbName, tblName)) {
-        return succ;
-      }
-
       if (colNames == null || colNames.isEmpty()) {
         colNames = getTable(catName, dbName, tblName)
             .getSd().getCols().stream().map(FieldSchema::getName)
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlAggrStats.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlAggrStats.java
similarity index 96%
rename from 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlAggrStats.java
rename to 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlAggrStats.java
index 9e1d4273060..0dc04ccee0f 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlAggrStats.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlAggrStats.java
@@ -16,10 +16,18 @@
  * limitations under the License.
  */
 
-package org.apache.hadoop.hive.metastore;
+package org.apache.hadoop.hive.metastore.directsql;
 
 import com.google.common.collect.ImmutableMap;
 import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.hive.metastore.Batchable;
+import org.apache.hadoop.hive.metastore.DatabaseProduct;
+import org.apache.hadoop.hive.metastore.Deadline;
+import org.apache.hadoop.hive.metastore.IExtrapolatePartStatus;
+import org.apache.hadoop.hive.metastore.LinearExtrapolatePartStatus;
+import org.apache.hadoop.hive.metastore.PersistenceManagerProvider;
+import org.apache.hadoop.hive.metastore.QueryWrapper;
+import org.apache.hadoop.hive.metastore.StatObjectConverter;
 import org.apache.hadoop.hive.metastore.api.ColumnStatistics;
 import org.apache.hadoop.hive.metastore.api.ColumnStatisticsData;
 import org.apache.hadoop.hive.metastore.api.ColumnStatisticsDesc;
@@ -58,12 +66,11 @@
 import static 
org.apache.hadoop.hive.metastore.IExtrapolatePartStatus.DBStatsAggrIndices.SUM_NDV_DECIMAL;
 import static 
org.apache.hadoop.hive.metastore.IExtrapolatePartStatus.DBStatsAggrIndices.COUNT_ROWS;
 import static 
org.apache.hadoop.hive.metastore.IExtrapolatePartStatus.DBStatsAggrIndices.SUM_NUM_DISTINCTS;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.getFullyQualifiedName;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.makeParams;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.executeWithArray;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.prepareParams;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.getFullyQualifiedName;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.makeParams;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.prepareParams;
 
-class DirectSqlAggrStats {
+public class DirectSqlAggrStats {
   private static final int NO_BATCHING = -1;
   private static final int DETECT_BATCHING = 0;
 
@@ -146,7 +153,7 @@ public List<Object[]> run(List<String> input) throws 
MetaException {
         long start = doTrace ? System.nanoTime() : 0;
         Query query = pm.newQuery("javax.jdo.query.SQL", queryText);
         try {
-          Object qResult = executeWithArray(query, params, queryText);
+          Object qResult = MetastoreDirectSqlUtils.executeWithArray(query, 
params, queryText);
           MetastoreDirectSqlUtils.timingTrace(doTrace, queryText0 + "...)", 
start,
               (doTrace ? System.nanoTime() : 0));
           if (qResult == null) {
@@ -479,7 +486,7 @@ public List<Object[]> run(List<String> inputPartNames)
             long start = doTrace ? System.nanoTime() : 0;
             Query<?> query = pm.newQuery("javax.jdo.query.SQL", queryText);
             try {
-              Object qResult = executeWithArray(query,
+              Object qResult = MetastoreDirectSqlUtils.executeWithArray(query,
                   prepareParams(catName, dbName, tableName, inputPartNames, 
inputColNames, engine), queryText);
               long end = doTrace ? System.nanoTime() : 0;
               MetastoreDirectSqlUtils.timingTrace(doTrace, queryText0, start, 
end);
@@ -518,7 +525,7 @@ public List<MetaStoreServerUtils.ColStatsObjWithSourceInfo> 
getColStatsForAllTab
     start = doTrace ? System.nanoTime() : 0;
     List<MetaStoreServerUtils.ColStatsObjWithSourceInfo> colStatsForDB = new 
ArrayList<>();
     try (QueryWrapper query = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", queryText))) {
-      qResult = executeWithArray(query.getInnerQuery(), new Object[] { dbName, 
catName }, queryText);
+      qResult = 
MetastoreDirectSqlUtils.executeWithArray(query.getInnerQuery(), new Object[] { 
dbName, catName }, queryText);
       if (qResult == null) {
         return colStatsForDB;
       }
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlBase.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlBase.java
similarity index 95%
rename from 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlBase.java
rename to 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlBase.java
index e15083f9502..8d454373d74 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlBase.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlBase.java
@@ -16,8 +16,9 @@
  * limitations under the License.
  */
 
-package org.apache.hadoop.hive.metastore;
+package org.apache.hadoop.hive.metastore.directsql;
 
+import org.apache.hadoop.hive.metastore.DatabaseProduct;
 import org.apache.hadoop.hive.metastore.api.MetaException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlDeleteStats.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlDeleteStats.java
new file mode 100644
index 00000000000..9ff23576605
--- /dev/null
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlDeleteStats.java
@@ -0,0 +1,243 @@
+/*
+ * 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.metastore.directsql;
+
+import javax.jdo.PersistenceManager;
+import javax.jdo.datastore.JDOConnection;
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.SQLException;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import org.apache.commons.lang3.tuple.Pair;
+import org.apache.hadoop.hive.common.StatsSetupConst;
+import org.apache.hadoop.hive.common.TableName;
+import org.apache.hadoop.hive.metastore.Batchable;
+import org.apache.hadoop.hive.metastore.QueryWrapper;
+import org.apache.hadoop.hive.metastore.api.MetaException;
+import org.apache.hadoop.hive.metastore.api.Table;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import static 
org.apache.hadoop.hive.metastore.directsql.MetaStoreDirectSql.getIdListForIn;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.closeDbConn;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.executeWithArray;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.extractSqlClob;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.makeParams;
+
+public class DirectSqlDeleteStats {
+  private static final Logger LOG = 
LoggerFactory.getLogger(DirectSqlDeleteStats.class);
+  private final MetaStoreDirectSql directSql;
+  private final PersistenceManager pm;
+  private final int batchSize;
+
+  public DirectSqlDeleteStats(MetaStoreDirectSql directSql, PersistenceManager 
pm) {
+    this.directSql = directSql;
+    this.pm = pm;
+    this.batchSize = directSql.getDirectSqlBatchSize();
+  }
+
+  public boolean deletePartitionColumnStats(String catName, String dbName, 
String tblName,
+      List<String> partNames, List<String> colNames, String engine) throws 
MetaException {
+    List<Long> partIds = Batchable.runBatched(batchSize, partNames, new 
Batchable<String, Long>() {
+      @Override
+      public List<Long> run(List<String> input) throws Exception {
+        String sqlFilter = "\"PARTITIONS\".\"PART_NAME\" in  (" + 
makeParams(input.size()) + ")";
+        List<Long> partitionIds = 
directSql.getPartitionIdsViaSqlFilter(catName, dbName, tblName, sqlFilter,
+            input, Collections.emptyList(), -1);
+        if (!partitionIds.isEmpty()) {
+          String deleteSql = "delete from \"PART_COL_STATS\" where \"PART_ID\" 
in ( " + getIdListForIn(partitionIds) + ")";
+          List<Object> params = new ArrayList<>(colNames == null ? 1 : 
colNames.size() + 1);
+
+          if (colNames != null && !colNames.isEmpty()) {
+            deleteSql += " and \"COLUMN_NAME\" in (" + 
makeParams(colNames.size()) + ")";
+            params.addAll(colNames);
+          }
+
+          if (engine != null) {
+            deleteSql += " and \"ENGINE\" = ?";
+            params.add(engine);
+          }
+          try (QueryWrapper queryParams = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", deleteSql))) {
+            executeWithArray(queryParams.getInnerQuery(), params.toArray(), 
deleteSql);
+          }
+        }
+        return partitionIds;
+      }
+    });
+    try {
+      return updateColumnStatsAccurateForPartitions(partIds, colNames);
+    } catch (SQLException e) {
+      String errorMsg = String.format(
+          "Failed to update partitions' COLUMN_STATS_ACCURATE for table '%s' "
+              + "(partitionCount=%d, columnCount=%d)",
+          new TableName(catName, dbName, tblName),
+          partIds == null ? 0 : partIds.size(),
+          colNames == null ? 0 : colNames.size());
+      LOG.error(errorMsg, e);
+      MetaException metaException = new MetaException(errorMsg + ": " + 
e.getMessage());
+      metaException.initCause(e);
+      throw metaException;
+    }
+  }
+
+  /**
+   * A helper function which will get the current COLUMN_STATS_ACCURATE 
parameter on table level
+   * and update the COLUMN_STATS_ACCURATE parameter with the new value on 
table level using directSql
+   */
+  private long updateColumnStatsAccurateForTable(Table table, List<String> 
droppedCols) throws MetaException {
+    Map<String, String> params = table.getParameters();
+    // get the current COLUMN_STATS_ACCURATE
+    String currentValue;
+    if (params == null || (currentValue = 
params.get(StatsSetupConst.COLUMN_STATS_ACCURATE)) == null) {
+      return 0;
+    }
+    // if the dropping columns is empty, that means we delete all the columns
+    if (droppedCols == null || droppedCols.isEmpty()) {
+      StatsSetupConst.clearColumnStatsState(params);
+    } else {
+      StatsSetupConst.removeColumnStatsState(params, droppedCols);
+    }
+
+    String updatedValue = params.get(StatsSetupConst.COLUMN_STATS_ACCURATE);
+    // if the COL_STATS_ACCURATE has changed, then update it using directSql
+    if (currentValue.equals(updatedValue)) {
+      return 0;
+    }
+    return directSql.updateTableParam(table, 
StatsSetupConst.COLUMN_STATS_ACCURATE, currentValue, updatedValue);
+  }
+
+  private boolean updateColumnStatsAccurateForPartitions(List<Long> partIds, 
List<String> colNames)
+      throws MetaException, SQLException {
+    // Get the list of params that need to be updated
+    List<Pair<Long, String>> updates = getPartColAccuToUpdate(partIds, 
colNames);
+    if (updates.isEmpty()) {
+      // Nothing to update: treat as successful completion
+      return true;
+    }
+    JDOConnection jdoConn = null;
+    try {
+      jdoConn = pm.getDataStoreConnection();
+      Connection dbConn = (Connection) jdoConn.getNativeConnection();
+      String update = "UPDATE \"PARTITION_PARAMS\" SET " + " \"PARAM_VALUE\" = 
?" +
+          " WHERE \"PART_ID\" = ? AND \"PARAM_KEY\" = ?";
+      try (PreparedStatement pst = dbConn.prepareStatement(update)) {
+        List<Long> updated = new ArrayList<>();
+        for (Pair<Long, String> accurate : updates) {
+          pst.setString(1, accurate.getRight());
+          pst.setLong(2, accurate.getLeft());
+          pst.setString(3, StatsSetupConst.COLUMN_STATS_ACCURATE);
+          pst.addBatch();
+          updated.add(accurate.getLeft());
+          if (updated.size() == batchSize) {
+            LOG.debug("Execute updates on part: {}", updated);
+            verifyUpdates(pst.executeBatch(), updated);
+            updated = new ArrayList<>();
+          }
+        }
+        if (!updated.isEmpty()) {
+          verifyUpdates(pst.executeBatch(), updated);
+        }
+      }
+      return true;
+    } finally {
+      closeDbConn(jdoConn);
+    }
+  }
+
+  private void verifyUpdates(int[] numUpdates, List<Long> partIds) throws 
MetaException {
+    for (int i = 0; i < numUpdates.length; i++) {
+      if (numUpdates[i] != 1) {
+        throw new MetaException("Invalid state of PARTITION_PARAMS ("
+            + StatsSetupConst.COLUMN_STATS_ACCURATE + ") for PART_ID " + 
partIds.get(i));
+      }
+    }
+  }
+
+  private List<Pair<Long, String>> getPartColAccuToUpdate(List<Long> partIds, 
List<String> colNames) throws MetaException {
+    return Batchable.runBatched(batchSize, partIds, new Batchable<>() {
+      @Override
+      public List<Pair<Long, String>> run(List<Long> input) throws Exception {
+        // 3. Get current COLUMN_STATS_ACCURATE values
+        String queryText = "SELECT \"PART_ID\", \"PARAM_VALUE\" FROM 
\"PARTITION_PARAMS\"" +
+            " WHERE \"PARAM_KEY\" = ? AND \"PART_ID\" IN (" + 
makeParams(input.size()) + ")";
+        Object[] params = new Object[1 + input.size()];
+        params[0] = StatsSetupConst.COLUMN_STATS_ACCURATE;
+        for (int i = 0; i < input.size(); i++) {
+          params[i + 1] = input.get(i);
+        }
+
+        List<Pair<Long, String>> result = new ArrayList<>();
+        try (QueryWrapper query = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", queryText))) {
+          @SuppressWarnings("unchecked") List<Object> sqlResult = 
executeWithArray(query.getInnerQuery(), params,
+              queryText);
+          for (Object row : sqlResult) {
+            Object[] fields = (Object[]) row;
+            Long partId = MetastoreDirectSqlUtils.extractSqlLong(fields[0]);
+            if (fields[1] == null) {
+              continue;
+            }
+            Map<String, String> parameters = new HashMap<>();
+            String accurateBefore = extractSqlClob(fields[1]);
+            parameters.put(StatsSetupConst.COLUMN_STATS_ACCURATE, 
accurateBefore);
+            if (colNames == null || colNames.isEmpty()) {
+              StatsSetupConst.clearColumnStatsState(parameters);
+            } else {
+              StatsSetupConst.removeColumnStatsState(parameters, colNames);
+            }
+            String accurateAfter = 
parameters.get(StatsSetupConst.COLUMN_STATS_ACCURATE);
+            if (accurateBefore.equals(accurateAfter)) {
+              continue;
+            }
+            result.add(Pair.of(partId, accurateAfter));
+          }
+        }
+        return result;
+      }
+    });
+  }
+
+  public boolean deleteTableColumnStatistics(Table table, List<String> 
colNames, String engine) throws MetaException {
+    String deleteSql = "delete from \"TAB_COL_STATS\" where \"TBL_ID\" = ?";
+    List<Object> params = new ArrayList<>();
+    params.add(table.getId());
+
+    if (colNames != null && !colNames.isEmpty()) {
+      deleteSql += " and \"COLUMN_NAME\" in (" + makeParams(colNames.size()) + 
")";
+      params.addAll(colNames);
+    }
+    if (engine != null) {
+      deleteSql += " and \"ENGINE\" = ?";
+      params.add(engine);
+    }
+    try (QueryWrapper queryParams = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", deleteSql))) {
+      executeWithArray(queryParams.getInnerQuery(), params.toArray(), 
deleteSql);
+    }
+    long numUpdated = updateColumnStatsAccurateForTable(table, colNames);
+    if (numUpdated == 0 && LOG.isDebugEnabled()) {
+      LOG.debug("No COLUMN_STATS_ACCURATE rows updated for table {}", 
table.getTableName());
+    }
+    // Return true as long as the delete (and any required metadata updates) 
completed without exception.
+    return true;
+  }
+}
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlInsertPart.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlInsertPart.java
similarity index 99%
rename from 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlInsertPart.java
rename to 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlInsertPart.java
index 3494a402f71..1a489ad4522 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlInsertPart.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlInsertPart.java
@@ -16,10 +16,10 @@
  * limitations under the License.
  */
 
-package org.apache.hadoop.hive.metastore;
+package org.apache.hadoop.hive.metastore.directsql;
 
 import static org.apache.commons.lang3.StringUtils.repeat;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.getModelIdentity;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.getModelIdentity;
 
 import java.util.ArrayList;
 import java.util.Collections;
@@ -32,6 +32,8 @@
 import javax.jdo.identity.LongIdentity;
 
 import org.apache.commons.collections4.CollectionUtils;
+import org.apache.hadoop.hive.metastore.DatabaseProduct;
+import org.apache.hadoop.hive.metastore.QueryWrapper;
 import org.apache.hadoop.hive.metastore.api.MetaException;
 import org.apache.hadoop.hive.metastore.model.MColumnDescriptor;
 import org.apache.hadoop.hive.metastore.model.MFieldSchema;
@@ -51,7 +53,7 @@ class DirectSqlInsertPart {
   private final DatabaseProduct dbType;
   private final int batchSize;
 
-  public DirectSqlInsertPart(PersistenceManager pm, DatabaseProduct dbType, 
int batchSize) {
+  DirectSqlInsertPart(PersistenceManager pm, DatabaseProduct dbType, int 
batchSize) {
     this.pm = pm;
     this.dbType = dbType;
     this.batchSize = batchSize;
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlUpdateParams.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlUpdateParams.java
similarity index 91%
rename from 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlUpdateParams.java
rename to 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlUpdateParams.java
index 0d40ef87501..3ad22fe5856 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlUpdateParams.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlUpdateParams.java
@@ -15,7 +15,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-package org.apache.hadoop.hive.metastore;
+package org.apache.hadoop.hive.metastore.directsql;
 
 import java.sql.SQLException;
 import java.util.ArrayList;
@@ -27,12 +27,15 @@
 import javax.jdo.PersistenceManager;
 
 import org.apache.commons.lang3.tuple.Pair;
+import org.apache.hadoop.hive.metastore.Batchable;
+import org.apache.hadoop.hive.metastore.DatabaseProduct;
+import org.apache.hadoop.hive.metastore.QueryWrapper;
 import org.apache.hadoop.hive.metastore.api.MetaException;
 import org.apache.hadoop.hive.metastore.txn.TxnUtils;
 
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.executeWithArray;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.extractSqlClob;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.extractSqlLong;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.executeWithArray;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.extractSqlClob;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.extractSqlLong;
 
 /**
  * Shared helper to diff and apply parameter table changes 
(delete/update/insert) in batches.
@@ -43,7 +46,7 @@ class DirectSqlUpdateParams extends DirectSqlBase {
     super(pm, dbType, batchSize);
   }
 
-  void run(String paramTable, String idColumn, List<Long> ids, Map<Long, 
Optional<Map<String, String>>> newParamsOpt)
+  public void run(String paramTable, String idColumn, List<Long> ids, 
Map<Long, Optional<Map<String, String>>> newParamsOpt)
       throws MetaException {
     Map<Long, Map<String, String>> oldParams = getParams(paramTable, idColumn, 
ids);
 
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlUpdatePart.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlUpdatePart.java
similarity index 97%
rename from 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlUpdatePart.java
rename to 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlUpdatePart.java
index a2b1bf2ebd3..c22df8e54d9 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/DirectSqlUpdatePart.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/DirectSqlUpdatePart.java
@@ -16,10 +16,17 @@
  * limitations under the License.
  */
 
-package org.apache.hadoop.hive.metastore;
+package org.apache.hadoop.hive.metastore.directsql;
 
 import org.apache.commons.lang3.tuple.Pair;
 import org.apache.hadoop.hive.common.StatsSetupConst;
+import org.apache.hadoop.hive.metastore.Batchable;
+import org.apache.hadoop.hive.metastore.DatabaseProduct;
+import org.apache.hadoop.hive.metastore.MetaStoreListenerNotifier;
+import org.apache.hadoop.hive.metastore.ObjectStore;
+import org.apache.hadoop.hive.metastore.QueryWrapper;
+import org.apache.hadoop.hive.metastore.StatObjectConverter;
+import org.apache.hadoop.hive.metastore.TransactionalMetaStoreEventListener;
 import org.apache.hadoop.hive.metastore.api.ColumnStatistics;
 import org.apache.hadoop.hive.metastore.api.ColumnStatisticsDesc;
 import org.apache.hadoop.hive.metastore.api.ColumnStatisticsObj;
@@ -67,11 +74,10 @@
 
 import static 
org.apache.hadoop.hive.common.StatsSetupConst.COLUMN_STATS_ACCURATE;
 import static org.apache.hadoop.hive.metastore.HMSHandler.getPartValsFromName;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.executeWithArray;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.extractSqlClob;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.extractSqlInt;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.extractSqlLong;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.getModelIdentity;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.extractSqlClob;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.extractSqlInt;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.extractSqlLong;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.getModelIdentity;
 
 /**
  * This class contains the optimizations for MetaStore that rely on direct SQL 
access to
@@ -80,7 +86,7 @@
  *
  * This class separates out the update part from MetaStoreDirectSql class.
  */
-class DirectSqlUpdatePart extends DirectSqlBase{
+class DirectSqlUpdatePart extends DirectSqlBase {
   private static final Logger LOG = 
LoggerFactory.getLogger(DirectSqlUpdatePart.class.getName());
 
   private final Configuration conf;
@@ -567,14 +573,14 @@ public List<Long> run(List<Long> input) throws Exception {
             "from \"SKEWED_VALUES\" where \"SD_ID_OID\" in (" + idLists + ")";
         try (QueryWrapper query =
                  new QueryWrapper(pm.newQuery("javax.jdo.query.SQL", 
queryFromSkewedValues))) {
-          List<Long> sqlResult = executeWithArray(query.getInnerQuery(), null, 
queryFromSkewedValues);
+          List<Long> sqlResult = 
MetastoreDirectSqlUtils.executeWithArray(query.getInnerQuery(), null, 
queryFromSkewedValues);
           result.addAll(sqlResult);
         }
         String queryFromValueLoc = "select \"STRING_LIST_ID_KID\" " +
             "from \"SKEWED_COL_VALUE_LOC_MAP\" where \"SD_ID\" in (" + idLists 
+ ")";
         try (QueryWrapper query =
                  new QueryWrapper(pm.newQuery("javax.jdo.query.SQL", 
queryFromValueLoc))) {
-          List<Long> sqlResult = executeWithArray(query.getInnerQuery(), null, 
queryFromValueLoc);
+          List<Long> sqlResult = 
MetastoreDirectSqlUtils.executeWithArray(query.getInnerQuery(), null, 
queryFromValueLoc);
           result.addAll(sqlResult);
         }
         return result;
@@ -601,7 +607,7 @@ public List<Void> run(List<Long> input) throws Exception {
         String queryText = "select \"SD_ID\", \"CD_ID\", \"SERDE_ID\" from 
\"SDS\" " +
             "where \"SD_ID\" in (" + idLists + ")";
         try (QueryWrapper query = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", queryText))) {
-          List<Object[]> sqlResult = executeWithArray(query.getInnerQuery(), 
null, queryText);
+          List<Object[]> sqlResult = 
MetastoreDirectSqlUtils.executeWithArray(query.getInnerQuery(), null, 
queryText);
           for (Object[] row : sqlResult) {
             Long sdId = extractSqlLong(row[0]);
             Long cdId = extractSqlLong(row[1]);
@@ -663,7 +669,7 @@ public List<Long> run(List<Long> input) throws Exception {
             String queryText = "select DISTINCT \"CD_ID\" from \"SDS\" where 
\"CD_ID\" in ( " + idLists + ")";
             List<Long> cdIds = new ArrayList<>();
             try (QueryWrapper query = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", queryText))) {
-              List<Object> sqlResult = executeWithArray(query.getInnerQuery(), 
null, queryText);
+              List<Object[]> sqlResult = 
MetastoreDirectSqlUtils.executeWithArray(query.getInnerQuery(), null, 
queryText);
               if (sqlResult != null) {
                 for (Object cdId : sqlResult) {
                   cdIds.add(MetastoreDirectSqlUtils.extractSqlLong(cdId));
@@ -984,7 +990,7 @@ public List<Void> run(List<Long> input) throws Exception {
         String queryText = "select \"CD_ID\", \"COMMENT\", \"COLUMN_NAME\", 
\"TYPE_NAME\", " +
             "\"INTEGER_IDX\" from \"COLUMNS_V2\" where \"CD_ID\" in (" + 
idLists + ")";
         try (QueryWrapper query = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", queryText))) {
-          List<Object[]> sqlResult = executeWithArray(query.getInnerQuery(), 
null, queryText);
+          List<Object[]> sqlResult = 
MetastoreDirectSqlUtils.executeWithArray(query.getInnerQuery(), null, 
queryText);
           for (Object[] row : sqlResult) {
             Long id = extractSqlLong(row[0]);
             String comment = extractSqlClob(row[1]);
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/MetaStoreDirectSql.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/MetaStoreDirectSql.java
similarity index 97%
rename from 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/MetaStoreDirectSql.java
rename to 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/MetaStoreDirectSql.java
index 403355100c2..3a317104ad8 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/MetaStoreDirectSql.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/MetaStoreDirectSql.java
@@ -16,7 +16,7 @@
  * limitations under the License.
  */
 
-package org.apache.hadoop.hive.metastore;
+package org.apache.hadoop.hive.metastore.directsql;
 
 import static org.apache.commons.lang3.StringUtils.join;
 import static org.apache.commons.lang3.StringUtils.normalizeSpace;
@@ -29,9 +29,9 @@
 import static org.apache.hadoop.hive.metastore.ColumnType.TIMESTAMP_TYPE_NAME;
 import static org.apache.hadoop.hive.metastore.ColumnType.TINYINT_TYPE_NAME;
 import static org.apache.hadoop.hive.metastore.ColumnType.VARCHAR_TYPE_NAME;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.getFullyQualifiedName;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.makeParams;
-import static 
org.apache.hadoop.hive.metastore.MetastoreDirectSqlUtils.prepareParams;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.getFullyQualifiedName;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.makeParams;
+import static 
org.apache.hadoop.hive.metastore.directsql.MetastoreDirectSqlUtils.prepareParams;
 
 import java.sql.Connection;
 import java.sql.SQLException;
@@ -63,7 +63,19 @@
 import org.apache.commons.lang3.StringUtils;
 import org.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.hive.common.StatsSetupConst;
+import org.apache.hadoop.hive.metastore.AggregateStatsCache;
 import org.apache.hadoop.hive.metastore.AggregateStatsCache.AggrColStats;
+import org.apache.hadoop.hive.metastore.Batchable;
+import org.apache.hadoop.hive.metastore.ColumnType;
+import org.apache.hadoop.hive.metastore.DatabaseProduct;
+import org.apache.hadoop.hive.metastore.Deadline;
+import org.apache.hadoop.hive.metastore.ObjectStore;
+import org.apache.hadoop.hive.metastore.PartitionProjectionEvaluator;
+import org.apache.hadoop.hive.metastore.PersistenceManagerProvider;
+import org.apache.hadoop.hive.metastore.QueryWrapper;
+import org.apache.hadoop.hive.metastore.TableType;
+import org.apache.hadoop.hive.metastore.TransactionalMetaStoreEventListener;
+import org.apache.hadoop.hive.metastore.Warehouse;
 import org.apache.hadoop.hive.metastore.api.AggrStats;
 import org.apache.hadoop.hive.metastore.api.ColumnStatistics;
 import org.apache.hadoop.hive.metastore.api.ColumnStatisticsObj;
@@ -135,7 +147,7 @@
  * JDOQL partition retrieval is still present so as not to limit the ORM 
solution we have
  * to SQL stores only. There's always a way to do without direct SQL.
  */
-class MetaStoreDirectSql {
+public class MetaStoreDirectSql {
   private static final int NO_BATCHING = -1, DETECT_BATCHING = 0;
   private static final Set<String> ALLOWED_TABLES_TO_LOCK = 
Set.of("NOTIFICATION_SEQUENCE");
 
@@ -351,6 +363,10 @@ public boolean isCompatibleDatastore() {
     return isCompatibleDatastore;
   }
 
+  public int getDirectSqlBatchSize() {
+    return batchSize;
+  }
+
   private void executeNoResult(final String queryText) throws SQLException {
     JDOConnection jdoConn = pm.getDataStoreConnection();
     Statement statement = null;
@@ -868,10 +884,10 @@ public static class SqlFilterForPushdown {
     private String tableName;
     // whether should compact null elements in joins when generating sql 
filter.
     private boolean compactJoins = true;
-    SqlFilterForPushdown() {
+    public SqlFilterForPushdown() {
 
     }
-    SqlFilterForPushdown(Table table, boolean compactJoins) {
+    public SqlFilterForPushdown(Table table, boolean compactJoins) {
       this.catName = table.getCatName();
       this.dbName = table.getDbName();
       this.tableName = table.getTableName();
@@ -942,7 +958,7 @@ private boolean isViewTable(String catName, String dbName, 
String tblName) throw
    * @param max The maximum number of partitions to return.
    * @return List of partition objects.
    */
-  private List<Long> getPartitionIdsViaSqlFilter(
+   List<Long> getPartitionIdsViaSqlFilter(
       String catName, String dbName, String tblName, String sqlFilter,
       List<? extends Object> paramsForFilter, List<String> joinsForFilter, 
Integer max)
       throws MetaException {
@@ -2240,7 +2256,6 @@ public List<SQLCheckConstraint> 
getCheckConstraints(String catName, String db_na
    * @param dbName    Metastore db name.
    * @param tblName   Metastore table name.
    * @param partNames Partition names to get.
-   * @return List of partitions.
    */
   public void dropPartitionsViaSqlFilter(final String catName, final String 
dbName,
                                          final String tblName, List<String> 
partNames)
@@ -2706,62 +2721,6 @@ public void deleteColumnStatsState(long tbl_id) throws 
MetaException {
     }
   }
 
-  public boolean deleteTableColumnStatistics(long tableId, List<String> 
colNames, String engine) {
-    String deleteSql = "delete from " + TAB_COL_STATS + " where \"TBL_ID\" = 
?";
-    List<Object> params = new ArrayList<>(colNames == null ? 2 : 
colNames.size() + 2);
-    params.add(tableId);
-
-    if (colNames != null && !colNames.isEmpty()) {
-      deleteSql += " and \"COLUMN_NAME\" in (" + makeParams(colNames.size()) + 
")";
-      params.addAll(colNames);
-    }
-
-    if (engine != null) {
-      deleteSql += " and \"ENGINE\" = ?";
-      params.add(engine);
-    }
-
-    try (QueryWrapper queryParams = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", deleteSql))) {
-      executeWithArray(queryParams.getInnerQuery(), params.toArray(), 
deleteSql);
-    } catch (MetaException e) {
-      return false;
-    }
-
-    return true;
-  }
-
-  public boolean deletePartitionColumnStats(String catName, String dbName, 
String tblName,
-      List<String> partNames, List<String> colNames, String engine) throws 
MetaException {
-    Batchable.runBatched(batchSize, partNames, new Batchable<String, Void>() {
-      @Override
-      public List<Void> run(List<String> input) throws Exception {
-        String sqlFilter = PARTITIONS + ".\"PART_NAME\" in  (" + 
makeParams(input.size()) + ")";
-        List<Long> partitionIds = getPartitionIdsViaSqlFilter(catName, dbName, 
tblName, sqlFilter,
-            input, Collections.emptyList(), -1);
-        if (!partitionIds.isEmpty()) {
-          String deleteSql = "delete from " + PART_COL_STATS + " where 
\"PART_ID\" in ( " + getIdListForIn(partitionIds) + ")";
-          List<Object> params = new ArrayList<>(colNames == null ? 1 : 
colNames.size() + 1);
-
-          if (colNames != null && !colNames.isEmpty()) {
-            deleteSql += " and \"COLUMN_NAME\" in (" + 
makeParams(colNames.size()) + ")";
-            params.addAll(colNames);
-          }
-
-          if (engine != null) {
-            deleteSql += " and \"ENGINE\" = ?";
-            params.add(engine);
-          }
-
-          try (QueryWrapper queryParams = new 
QueryWrapper(pm.newQuery("javax.jdo.query.SQL", deleteSql))) {
-            executeWithArray(queryParams.getInnerQuery(), params.toArray(), 
deleteSql);
-          }
-        }
-        return null;
-      }
-    });
-    return true;
-  }
-
   public Map<String, Map<String, String>> updatePartitionColumnStatisticsBatch(
                                                       Map<String, 
ColumnStatistics> partColStatsMap,
                                                       Table tbl,
@@ -2874,7 +2833,7 @@ private List<Long> getFunctionIds(String catName) throws 
MetaException {
     }
   }
 
-  long updateTableParam(Table table, String key, String expectedValue, String 
newValue) {
+  public long updateTableParam(Table table, String key, String expectedValue, 
String newValue) {
     String statement = TxnUtils.createUpdatePreparedStmt(
         "\"TABLE_PARAMS\"",
         ImmutableList.of("\"PARAM_VALUE\""),
diff --git 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/MetastoreDirectSqlUtils.java
 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/MetastoreDirectSqlUtils.java
similarity index 92%
rename from 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/MetastoreDirectSqlUtils.java
rename to 
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/MetastoreDirectSqlUtils.java
index 19b92f5177c..a3c4523dc87 100644
--- 
a/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/MetastoreDirectSqlUtils.java
+++ 
b/standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/directsql/MetastoreDirectSqlUtils.java
@@ -17,10 +17,11 @@
  *  limitations under the License.
  */
 
-package org.apache.hadoop.hive.metastore;
+package org.apache.hadoop.hive.metastore.directsql;
 
 import org.apache.commons.lang3.BooleanUtils;
 import org.apache.commons.lang3.exception.ExceptionUtils;
+import org.apache.hadoop.hive.metastore.Deadline;
 import org.apache.hadoop.hive.metastore.api.FieldSchema;
 import org.apache.hadoop.hive.metastore.api.Function;
 import org.apache.hadoop.hive.metastore.api.MetaException;
@@ -41,6 +42,7 @@
 
 import javax.jdo.PersistenceManager;
 import javax.jdo.Query;
+import javax.jdo.datastore.JDOConnection;
 import java.math.BigDecimal;
 import java.math.MathContext;
 import java.sql.Blob;
@@ -59,13 +61,13 @@
 /**
  * Helper utilities used by DirectSQL code in HiveMetastore.
  */
-class MetastoreDirectSqlUtils {
+public class MetastoreDirectSqlUtils {
   private static final Logger LOG = 
LoggerFactory.getLogger(MetastoreDirectSqlUtils.class);
   private MetastoreDirectSqlUtils() {
 
   }
   @SuppressWarnings("unchecked")
-  static <T> T executeWithArray(Query query, Object[] params, String sql) 
throws MetaException {
+  public static <T> T executeWithArray(Query query, Object[] params, String 
sql) throws MetaException {
     return (T)executeWithArray(query, params, sql, -1);
   }
 
@@ -98,7 +100,7 @@ static List<Object[]> ensureList(Object result) throws 
MetaException {
     return (List<Object[]>)result;
   }
 
-  static Long extractSqlLong(Object obj) throws MetaException {
+  public static Long extractSqlLong(Object obj) throws MetaException {
     if (obj == null) return null;
     if (!(obj instanceof Number)) {
       throw new MetaException("Expected numeric type but got " + 
obj.getClass().getName());
@@ -106,7 +108,7 @@ static Long extractSqlLong(Object obj) throws MetaException 
{
     return ((Number)obj).longValue();
   }
 
-  static void timingTrace(boolean doTrace, String queryText, long start, long 
queryTime) {
+  public static void timingTrace(boolean doTrace, String queryText, long 
start, long queryTime) {
     if (!doTrace) return;
     LOG.debug("Direct SQL query in " + (queryTime - start) / 1000000.0 + "ms + 
" +
         (System.nanoTime() - queryTime) / 1000000.0 + "ms, the query is [" + 
queryText + "]");
@@ -187,7 +189,7 @@ static <T> int loopJoinOrderedResult(PersistenceManager pm, 
TreeMap<Long, T> tre
     return rv;
   }
 
-  static void setPartitionParametersWithFilter(String PARTITION_PARAMS,
+  public static void setPartitionParametersWithFilter(String PARTITION_PARAMS,
       boolean convertMapNullsToEmptyStrings, PersistenceManager pm, String 
partIds,
       TreeMap<Long, Partition> partitions, String includeParamKeyPattern, 
String excludeParamKeyPattern)
       throws MetaException {
@@ -220,7 +222,7 @@ public void apply(Partition t, Object[] fields) {
     }
   }
 
-  static void setPartitionValues(String PARTITION_KEY_VALS, PersistenceManager 
pm, String partIds,
+  public static void setPartitionValues(String PARTITION_KEY_VALS, 
PersistenceManager pm, String partIds,
       TreeMap<Long, Partition> partitions)
       throws MetaException {
     String queryText;
@@ -234,7 +236,7 @@ public void apply(Partition t, Object[] fields) {
       }});
   }
 
-  static String extractSqlClob(Object value) {
+  public static String extractSqlClob(Object value) {
     if (value == null) return null;
     try {
       if (value instanceof Clob) {
@@ -249,7 +251,7 @@ static String extractSqlClob(Object value) {
     }
   }
 
-  static void setSDParameters(String SD_PARAMS, boolean 
convertMapNullsToEmptyStrings,
+  public static void setSDParameters(String SD_PARAMS, boolean 
convertMapNullsToEmptyStrings,
       PersistenceManager pm, TreeMap<Long, StorageDescriptor> sds, String 
sdIds)
       throws MetaException {
     String queryText;
@@ -267,11 +269,11 @@ public void apply(StorageDescriptor t, Object[] fields) {
     }
   }
 
-  static int extractSqlInt(Object field) {
+  public static int extractSqlInt(Object field) {
     return ((Number)field).intValue();
   }
 
-  static void setSDSortCols(String SORT_COLS, List<String> columnNames, 
PersistenceManager pm,
+  public static void setSDSortCols(String SORT_COLS, List<String> columnNames, 
PersistenceManager pm,
       TreeMap<Long, StorageDescriptor> sds, String sdIds)
       throws MetaException {
     StringBuilder queryTextBuilder = new StringBuilder("select \"SD_ID\"");
@@ -309,7 +311,7 @@ public void apply(StorageDescriptor t, Object[] fields) {
       }});
   }
 
-  static void setSDSortCols(String SORT_COLS, PersistenceManager pm,
+  public static void setSDSortCols(String SORT_COLS, PersistenceManager pm,
       TreeMap<Long, StorageDescriptor> sds, String sdIds)
       throws MetaException {
     String queryText;
@@ -325,7 +327,7 @@ public void apply(StorageDescriptor t, Object[] fields) {
       }});
   }
 
-  static void setSDBucketCols(String BUCKETING_COLS, PersistenceManager pm,
+  public static void setSDBucketCols(String BUCKETING_COLS, PersistenceManager 
pm,
       TreeMap<Long, StorageDescriptor> sds, String sdIds)
       throws MetaException {
     String queryText;
@@ -339,7 +341,7 @@ public void apply(StorageDescriptor t, Object[] fields) {
       }});
   }
 
-  static boolean setSkewedColNames(String SKEWED_COL_NAMES, PersistenceManager 
pm,
+  public static boolean setSkewedColNames(String SKEWED_COL_NAMES, 
PersistenceManager pm,
       TreeMap<Long, StorageDescriptor> sds, String sdIds)
       throws MetaException {
     String queryText;
@@ -354,7 +356,7 @@ public void apply(StorageDescriptor t, Object[] fields) {
       }}) > 0;
   }
 
-  static void setSkewedColValues(String SKEWED_STRING_LIST_VALUES, String 
SKEWED_VALUES,
+  public static void setSkewedColValues(String SKEWED_STRING_LIST_VALUES, 
String SKEWED_VALUES,
       PersistenceManager pm, TreeMap<Long, StorageDescriptor> sds, String 
sdIds)
       throws MetaException {
     String queryText;
@@ -394,7 +396,7 @@ public void apply(StorageDescriptor t, Object[] fields) 
throws MetaException {
       }});
   }
 
-  static void setSkewedColLocationMaps(String SKEWED_COL_VALUE_LOC_MAP,
+  public static void setSkewedColLocationMaps(String SKEWED_COL_VALUE_LOC_MAP,
       String SKEWED_STRING_LIST_VALUES, PersistenceManager pm, TreeMap<Long, 
StorageDescriptor> sds,
       String sdIds)
       throws MetaException {
@@ -443,7 +445,7 @@ public void apply(StorageDescriptor t, Object[] fields) 
throws MetaException {
       }});
   }
 
-  static void setSDCols(String COLUMNS_V2, List<String> columnNames, 
PersistenceManager pm,
+  public static void setSDCols(String COLUMNS_V2, List<String> columnNames, 
PersistenceManager pm,
       TreeMap<Long, List<FieldSchema>> colss, String colIds)
       throws MetaException {
     StringBuilder queryTextBuilder = new StringBuilder("select \"CD_ID\"");
@@ -485,7 +487,7 @@ public void apply(List<FieldSchema> t, Object[] fields) {
       }});
   }
 
-  static void setSDCols(String COLUMNS_V2, PersistenceManager pm,
+  public static void setSDCols(String COLUMNS_V2, PersistenceManager pm,
       TreeMap<Long, List<FieldSchema>> colss, String colIds)
       throws MetaException {
     String queryText;
@@ -499,7 +501,7 @@ public void apply(List<FieldSchema> t, Object[] fields) {
       }});
   }
 
-  static void setSerdeParams(String SERDE_PARAMS, boolean 
convertMapNullsToEmptyStrings,
+  public static void setSerdeParams(String SERDE_PARAMS, boolean 
convertMapNullsToEmptyStrings,
       PersistenceManager pm, TreeMap<Long, SerDeInfo> serdes, String serdeIds) 
throws MetaException {
     String queryText;
     queryText = "select \"SERDE_ID\", \"PARAM_KEY\", \"PARAM_VALUE\" from " + 
SERDE_PARAMS + ""
@@ -542,7 +544,7 @@ static void setFunctionResourceUris(String FUNC_RU, 
PersistenceManager pm, Strin
    * @throws MetaException
    *           if the column value cannot be converted into a Boolean object
    */
-  static Boolean extractSqlBoolean(Object value) throws MetaException {
+  public static Boolean extractSqlBoolean(Object value) throws MetaException {
     if (value == null) {
       return null;
     }
@@ -577,12 +579,12 @@ static Boolean extractSqlBoolean(Object value) throws 
MetaException {
     throw new MetaException("Cannot extract boolean from column value " + 
value);
   }
 
-  static String extractSqlString(Object value) {
+ public static String extractSqlString(Object value) {
     if (value == null) return null;
     return value.toString();
   }
 
-  static Double extractSqlDouble(Object obj) throws MetaException {
+ public static Double extractSqlDouble(Object obj) throws MetaException {
     if (obj == null)
       return null;
     if (!(obj instanceof Number)) {
@@ -591,7 +593,7 @@ static Double extractSqlDouble(Object obj) throws 
MetaException {
     return ((Number) obj).doubleValue();
   }
 
-  static byte[] extractSqlBlob(Object value) throws MetaException {
+  public static byte[] extractSqlBlob(Object value) throws MetaException {
     if (value == null)
       return null;
     if (value instanceof Blob) {
@@ -717,4 +719,14 @@ public static Object max(Object first, Object second) {
     }
     return (new BigDecimal(first.toString())).max(new 
BigDecimal(second.toString()));
   }
+
+  static void closeDbConn(JDOConnection jdoConn) {
+    try {
+      if (jdoConn != null) {
+        jdoConn.close();
+      }
+    } catch (Exception e) {
+      LOG.warn("Failed to close db connection", e);
+    }
+  }
 }
diff --git 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/DummyCustomRDBMS.java
 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/DummyCustomRDBMS.java
index aa2ba60934a..005fdc99b77 100644
--- 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/DummyCustomRDBMS.java
+++ 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/DummyCustomRDBMS.java
@@ -54,7 +54,7 @@ public String getHiveSchemaPostfix() {
     return "DummyPostfix";
   }
   @Override
-  protected String toDate(String tableValue) {
+  public String toDate(String tableValue) {
     return "DummyDate";
   }
   @Override
diff --git 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestHiveMetaStore.java
 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestHiveMetaStore.java
index 939f09c3343..8f1ba073f1b 100644
--- 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestHiveMetaStore.java
+++ 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestHiveMetaStore.java
@@ -1849,6 +1849,11 @@ public void testColumnStatistics() throws Throwable {
       List<ColumnStatisticsObj> stats = client.getTableColumnStatistics(
           dbName, tblName, Lists.newArrayList(colName[1]), ENGINE);
       assertTrue("stats are not empty: " + stats, stats.isEmpty());
+      // test if all columns are deleted from parameter COLUMN_STATS_ACCURATE
+      Map<String, String> tableParams = client.getTable(dbName, 
tblName).getParameters();
+      String table_column_stats_accurate = 
tableParams.get(StatsSetupConst.COLUMN_STATS_ACCURATE);
+      assertTrue("parameter COLUMN_STATS_ACCURATE is not accurate in " + 
tblName, table_column_stats_accurate == null ||
+              (!table_column_stats_accurate.contains(colName[0]) && 
!table_column_stats_accurate.contains(colName[1])));
 
       colStats.setStatsDesc(statsDesc);
       colStats.setStatsObj(statsObjs);
@@ -1866,6 +1871,11 @@ public void testColumnStatistics() throws Throwable {
       // multiple columns
       request.setCol_names(Arrays.asList(colName));
       assertTrue(client.deleteColumnStatistics(request));
+      // test if the columns in colName array are deleted from parameter 
COLUMN_STATS_ACCURATE
+      tableParams = client.getTable(dbName, tblName).getParameters();
+      table_column_stats_accurate = 
tableParams.get(StatsSetupConst.COLUMN_STATS_ACCURATE);
+      assertTrue("parameter COLUMN_STATS_ACCURATE is not accurate in " + 
tblName, table_column_stats_accurate == null ||
+              (!table_column_stats_accurate.contains(colName[0]) && 
!table_column_stats_accurate.contains(colName[1])));
       colStats3 = client.getTableColumnStatistics(
           dbName, tblName, Lists.newArrayList(colName), ENGINE);
       assertTrue("stats are not empty: " + colStats3, colStats3.isEmpty());
@@ -1961,6 +1971,12 @@ public void testColumnStatistics() throws Throwable {
           Lists.newArrayList(partitions.get(0), partitions.get(1), 
partitions.get(2)), Lists.newArrayList(colName), ENGINE);
      assertEquals(1, stats2.size());
      assertEquals(2, stats2.get(partitions.get(2)).size());
+     // test if all columns are deleted from parameter COLUMN_STATS_ACCURATE
+     Partition partition_0 = client.getPartition(dbName, tblName, 
partitions.get(0));
+     Map<String, String> partitionParams = partition_0.getParameters();
+     String partition_column_stats_accurate = 
partitionParams.get(StatsSetupConst.COLUMN_STATS_ACCURATE);
+     assertTrue("parameter COLUMN_STATS_ACCURATE is not accurate in " + 
partitions.get(0),partition_column_stats_accurate == null ||
+             (!partition_column_stats_accurate.contains(colName[0]) && 
!partition_column_stats_accurate.contains(colName[1])));
 
      // no partition or column name is set
      request.unsetPart_names();
diff --git 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestObjectStore.java
 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestObjectStore.java
index 0d22719dd6a..afede2f768c 100644
--- 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestObjectStore.java
+++ 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/TestObjectStore.java
@@ -22,6 +22,7 @@
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableSet;
 import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.hive.common.StatsSetupConst;
 import org.apache.hadoop.hive.metastore.annotation.MetastoreUnitTest;
 import org.apache.hadoop.hive.metastore.api.Catalog;
 import org.apache.hadoop.hive.metastore.api.ColumnStatistics;
@@ -73,6 +74,7 @@
 import org.apache.hadoop.hive.metastore.columnstats.ColStatsBuilder;
 import org.apache.hadoop.hive.metastore.conf.MetastoreConf;
 import org.apache.hadoop.hive.metastore.conf.MetastoreConf.ConfVars;
+import org.apache.hadoop.hive.metastore.directsql.MetaStoreDirectSql;
 import org.apache.hadoop.hive.metastore.messaging.EventMessage;
 import org.apache.hadoop.hive.metastore.metrics.Metrics;
 import org.apache.hadoop.hive.metastore.metrics.MetricsConstants;
@@ -1947,6 +1949,106 @@ protected Object getJdoResult(ObjectStore.GetHelper 
ctx) throws MetaException, N
     }
   }
 
+  @Test
+  public void testDeleteColumnStatsAccurate() throws Exception {
+    createPartitionedTable(false, true, new HashSet<>());
+    // add extra table column statistics
+    ColumnStatistics stats = new ColumnStatistics();
+    ColumnStatisticsDesc desc = new ColumnStatisticsDesc();
+    desc.setCatName(DEFAULT_CATALOG_NAME);
+    desc.setDbName(DB1);
+    desc.setTableName(TABLE1);
+    desc.setIsTblLevel(true);
+    stats.setStatsDesc(desc);
+    List<ColumnStatisticsObj> statsObjList = new ArrayList<>(1);
+    stats.setStatsObj(statsObjList);
+    stats.setEngine(ENGINE);
+    ColumnStatisticsData data = new 
ColStatsBuilder<>(long.class).numNulls(1).numDVs(2)
+        .low(3L).high(4L).hll(3, 4).kll(3, 4).build();
+    statsObjList.add(new ColumnStatisticsObj("test_col1", "int", data));
+    statsObjList.add(new ColumnStatisticsObj("test_col' 2", "int", data));
+    objectStore.updateTableColumnStatistics(stats, null, -1);
+
+    Table table = objectStore.getTable(DEFAULT_CATALOG_NAME, DB1, TABLE1);
+    Set<String> statsAccurate = Set.of("test_col1", "test_col' 2");
+    List<String> columns = 
StatsSetupConst.getColumnsHavingStats(table.getParameters());
+    Assert.assertTrue(statsAccurate.containsAll(columns) && columns.size() == 
2);
+
+    try (AutoCloseable c = directsql(true)) {
+      objectStore.deleteTableColumnStatistics(DEFAULT_CATALOG_NAME, DB1, 
TABLE1, List.of("test_col1"), null);
+      table = objectStore.getTable(DEFAULT_CATALOG_NAME, DB1, TABLE1);
+      columns = StatsSetupConst.getColumnsHavingStats(table.getParameters());
+      Assert.assertTrue(columns.size() == 1 && 
columns.getFirst().equals("test_col' 2"));
+    }
+     // reset again
+    objectStore.updateTableColumnStatistics(stats, null, -1);
+    table = objectStore.getTable(DEFAULT_CATALOG_NAME, DB1, TABLE1);
+    columns = StatsSetupConst.getColumnsHavingStats(table.getParameters());
+    Assert.assertTrue(statsAccurate.containsAll(columns) && columns.size() == 
2);
+
+    try (AutoCloseable c = directsql(false)) {
+      objectStore.deleteTableColumnStatistics(DEFAULT_CATALOG_NAME, DB1, 
TABLE1, List.of("test_col1"), null);
+      table = objectStore.getTable(DEFAULT_CATALOG_NAME, DB1, TABLE1);
+      columns = StatsSetupConst.getColumnsHavingStats(table.getParameters());
+      Assert.assertTrue(columns.size() == 1 && 
columns.getFirst().equals("test_col' 2"));
+    }
+
+    try (AutoCloseable c = directsql(false)) {
+      objectStore.deleteTableColumnStatistics(DEFAULT_CATALOG_NAME, DB1, 
TABLE1, null, null);
+      table = objectStore.getTable(DEFAULT_CATALOG_NAME, DB1, TABLE1);
+      columns = StatsSetupConst.getColumnsHavingStats(table.getParameters());
+      Assert.assertTrue(columns.isEmpty());
+    }
+
+    // test partition
+    desc.setIsTblLevel(false);
+    List<Partition> partitions;
+    try (AutoCloseable c = deadline()) {
+      partitions = objectStore.getPartitions(DEFAULT_CATALOG_NAME, DB1, 
TABLE1, GetPartitionsArgs.getAllPartitions());
+    }
+    MTable mTable = objectStore.ensureGetMTable(DEFAULT_CATALOG_NAME, DB1, 
TABLE1);
+    for (Partition partition : partitions) {
+      desc.setPartName(Warehouse.makePartName(table.getPartitionKeys(), 
partition.getValues()));
+      objectStore.updatePartitionColumnStatistics(table, mTable, stats, 
partition.getValues(), null, -1);
+      Partition newPart = objectStore.getPartition(DEFAULT_CATALOG_NAME, DB1, 
TABLE1, partition.getValues());
+      columns = StatsSetupConst.getColumnsHavingStats(newPart.getParameters());
+      Set<String> statsColumns = Set.of("test_col1", "test_col' 2", 
"test_part_col");
+      Assert.assertTrue(statsColumns.containsAll(columns) && columns.size() == 
3);
+    }
+
+    try (AutoCloseable c = directsql(true)) {
+      objectStore.deletePartitionColumnStatistics(DEFAULT_CATALOG_NAME, DB1, 
TABLE1,
+          List.of("test_part_col=a0"), List.of("test_col1"), null);
+      Partition partition = objectStore.getPartition(DEFAULT_CATALOG_NAME, 
DB1, TABLE1, List.of("a0"));
+      columns = 
StatsSetupConst.getColumnsHavingStats(partition.getParameters());
+      Assert.assertTrue(columns.contains("test_part_col") && 
columns.contains("test_col' 2"));
+    }
+
+    try (AutoCloseable c = directsql(false)) {
+      objectStore.deletePartitionColumnStatistics(DEFAULT_CATALOG_NAME, DB1, 
TABLE1,
+          List.of("test_part_col=a1"), List.of("test_col1"), null);
+      Partition partition = objectStore.getPartition(DEFAULT_CATALOG_NAME, 
DB1, TABLE1, List.of("a1"));
+      columns = 
StatsSetupConst.getColumnsHavingStats(partition.getParameters());
+      Assert.assertTrue(columns.contains("test_part_col") && 
columns.contains("test_col' 2"));
+    }
+
+    try (AutoCloseable c = directsql(true)) {
+      objectStore.deletePartitionColumnStatistics(DEFAULT_CATALOG_NAME, DB1, 
TABLE1,
+          List.of("test_part_col=a1"), null, null);
+      Partition partition = objectStore.getPartition(DEFAULT_CATALOG_NAME, 
DB1, TABLE1, List.of("a1"));
+      columns = 
StatsSetupConst.getColumnsHavingStats(partition.getParameters());
+      Assert.assertTrue(columns.isEmpty());
+    }
+
+    try (AutoCloseable c = directsql(false)) {
+      objectStore.deletePartitionColumnStatistics(DEFAULT_CATALOG_NAME, DB1, 
TABLE1,
+          List.of("test_part_col=a2"), null, null);
+      Partition partition = objectStore.getPartition(DEFAULT_CATALOG_NAME, 
DB1, TABLE1, List.of("a2"));
+      columns = 
StatsSetupConst.getColumnsHavingStats(partition.getParameters());
+      Assert.assertTrue(columns.isEmpty());
+    }
+  }
+
   /**
    * Helper method to check whether the Java system properties were set 
correctly in {@link ObjectStore#configureSSL(Configuration)}
    * @param useSSL whether or not SSL is enabled
@@ -1980,5 +2082,13 @@ public void close() throws Exception {
       }
     };
   }
+
+  AutoCloseable directsql(boolean tryDirectSql) {
+    boolean directsql = MetastoreConf.getBoolVar(objectStore.getConf(), 
ConfVars.TRY_DIRECT_SQL);
+    MetastoreConf.setBoolVar(objectStore.getConf(), ConfVars.TRY_DIRECT_SQL, 
tryDirectSql);
+    return () -> {
+      MetastoreConf.setBoolVar(objectStore.getConf(), ConfVars.TRY_DIRECT_SQL, 
directsql);
+    };
+  }
 }
 
diff --git 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/VerifyingObjectStore.java
 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/VerifyingObjectStore.java
index 9575bd44b4c..31bc5635a40 100644
--- 
a/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/VerifyingObjectStore.java
+++ 
b/standalone-metastore/metastore-server/src/test/java/org/apache/hadoop/hive/metastore/VerifyingObjectStore.java
@@ -137,7 +137,13 @@ public List<Partition> alterPartitions(String catName, 
String dbName, String tbl
       for (List<String> partVal : part_vals) {
         partNames.add(Warehouse.makePartName(partCols, partVal));
       }
-      List<Partition> oldParts = getPartitionsByNames(catName, dbName, 
tblName, partNames);
+      // We might have made some changes over the partition due to deleting the
+      // stale column statistics earlier, and they haven't been committed yet, 
the cached instances
+      // could be different from that in the datastore.
+      // We cannot verify the partitions by getPartitionsByNames now.
+      GetPartitionsArgs args = new 
GetPartitionsArgs.GetPartitionsArgsBuilder().partNames(partNames).build();
+      List<Partition> oldParts = getPartitionsByNamesInternal(
+          catName, dbName, tblName, true, true, args);
       if (oldParts.size() != partNames.size()) {
         throw new MetaException("Some partitions to be altered are missing");
       }

Reply via email to