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");
}