neatHyperTxt-meesho commented on code in PR #6785:
URL: https://github.com/apache/hive/pull/6785#discussion_r4039411233
##########
standalone-metastore/metastore-server/src/main/java/org/apache/hadoop/hive/metastore/cache/CachedStore.java:
##########
@@ -492,129 +517,242 @@ static void prewarm(RawStore rawStore) {
}
sharedCache.populateDatabasesInCache(databases);
LOG.info("Databases cache is now prewarmed. Now adding tables,
partitions and statistics to the cache");
- int numberOfDatabasesCachedSoFar = 0;
- for (Database db : databases) {
- String catName = StringUtils.normalizeIdentifier(db.getCatalogName());
- String dbName = StringUtils.normalizeIdentifier(db.getName());
- List<String> tblNames;
+ int prewarmThreads = Math.max(1,
+ MetastoreConf.getIntVar(rawStore.getConf(),
ConfVars.CACHED_RAW_STORE_PREWARM_THREADS));
+ ExecutorService prewarmPool = null;
+ List<RawStore> workerStores = new ArrayList<>();
+ if (prewarmThreads > 1) {
try {
- tblNames = rawStore.getAllTables(catName, dbName);
- } catch (MetaException e) {
- LOG.warn("Failed to cache tables for database " +
DatabaseName.getQualified(catName, dbName) + ", moving on");
- // Continue with next database
- continue;
+ // RawStore implementations (ObjectStore) are not thread safe, so
each worker gets its
+ // own instance and hence its own connection to the backing database
+ for (int i = 0; i < prewarmThreads; i++) {
+ workerStores.add(createRawStoreForPrewarm(rawStore.getConf()));
+ }
+ LOG.info("Prewarming table cache with {} threads", prewarmThreads);
+ prewarmPool = Executors.newFixedThreadPool(prewarmThreads, new
ThreadFactory() {
+ private final AtomicInteger threadCount = new AtomicInteger();
+ @Override public Thread newThread(Runnable r) {
+ Thread t = Executors.defaultThreadFactory().newThread(r);
+ t.setName("CachedStore-PrewarmWorker-" +
threadCount.getAndIncrement());
+ t.setDaemon(true);
+ return t;
+ }
+ });
+ } catch (RuntimeException e) {
+ LOG.warn("Failed to create RawStores for prewarm workers, falling
back to single threaded prewarm", e);
+ shutdownPrewarmWorkers(null, workerStores);
+ workerStores = new ArrayList<>();
}
- tblsPendingPrewarm.addTableNamesForPrewarming(tblNames);
- int totalTablesToCache = tblNames.size();
- int numberOfTablesCachedSoFar = 0;
- while (tblsPendingPrewarm.hasMoreTablesToPrewarm()) {
+ }
+ try {
+ int numberOfDatabasesCachedSoFar = 0;
+ for (Database db : databases) {
+ String catName =
StringUtils.normalizeIdentifier(db.getCatalogName());
+ String dbName = StringUtils.normalizeIdentifier(db.getName());
+ List<String> tblNames;
try {
- String tblName =
StringUtils.normalizeIdentifier(tblsPendingPrewarm.getNextTableNameToPrewarm());
- if (!shouldCacheTable(catName, dbName, tblName)) {
- continue;
- }
- Table table;
- try {
- table = rawStore.getTable(catName, dbName, tblName);
- } catch (MetaException e) {
- LOG.debug(ExceptionUtils.getStackTrace(e));
- // It is possible the table is deleted during fetching tables of
the database,
- // in that case, continue with the next table
- continue;
+ tblNames = rawStore.getAllTables(catName, dbName);
+ } catch (MetaException e) {
+ LOG.warn("Failed to cache tables for database {}, moving on",
DatabaseName.getQualified(catName, dbName));
+ // Continue with next database
+ continue;
+ }
+ tblsPendingPrewarm.addTableNamesForPrewarming(tblNames);
+ int totalTablesToCache = tblNames.size();
+ AtomicBoolean cacheMemoryFull = new AtomicBoolean(false);
+ AtomicInteger tablesCachedSoFar = new AtomicInteger();
+ if (prewarmPool != null) {
+ List<Future<?>> workers = new ArrayList<>(workerStores.size());
+ for (RawStore workerStore : workerStores) {
+ workers.add(prewarmPool.submit(
+ () -> drainTablesPendingPrewarm(workerStore, catName,
dbName, cacheMemoryFull, tablesCachedSoFar,
+ totalTablesToCache)));
}
- List<String> colNames =
MetaStoreUtils.getColumnNamesForTable(table);
- try {
- ColumnStatistics tableColStats = null;
- List<Partition> partitions = null;
- List<ColumnStatistics> partitionColStats = null;
- AggrStats aggrStatsAllPartitions = null;
- AggrStats aggrStatsAllButDefaultPartition = null;
- TableCacheObjects cacheObjects = new TableCacheObjects();
- if (!table.getPartitionKeys().isEmpty()) {
- Deadline.startTimer("getPartitions");
- partitions = rawStore.getPartitions(catName, dbName, tblName,
GetPartitionsArgs.getAllPartitions());
- Deadline.stopTimer();
- cacheObjects.setPartitions(partitions);
- List<String> partNames = new ArrayList<>(partitions.size());
- for (Partition p : partitions) {
-
partNames.add(Warehouse.makePartName(table.getPartitionKeys(), p.getValues()));
- }
- if (!partNames.isEmpty()) {
- // Get partition column stats for this table
- Deadline.startTimer("getPartitionColumnStatistics");
- partitionColStats =
- rawStore.getPartitionColumnStatistics(catName, dbName,
tblName, partNames, colNames, CacheUtils.HIVE_ENGINE);
- Deadline.stopTimer();
- cacheObjects.setPartitionColStats(partitionColStats);
- // Get aggregate stats for all partitions of a table and for
all but default
- // partition
- Deadline.startTimer("getAggrPartitionColumnStatistics");
- aggrStatsAllPartitions =
rawStore.get_aggr_stats_for(catName, dbName, tblName, partNames, colNames,
CacheUtils.HIVE_ENGINE);
- Deadline.stopTimer();
-
cacheObjects.setAggrStatsAllPartitions(aggrStatsAllPartitions);
- // Remove default partition from partition names and get
aggregate
- // stats again
- List<FieldSchema> partKeys = table.getPartitionKeys();
- String defaultPartitionValue =
- MetastoreConf.getVar(rawStore.getConf(),
ConfVars.DEFAULTPARTITIONNAME);
- List<String> partCols = new ArrayList<>();
- List<String> partVals = new ArrayList<>();
- for (FieldSchema fs : partKeys) {
- partCols.add(fs.getName());
- partVals.add(defaultPartitionValue);
- }
- String defaultPartitionName =
FileUtils.makePartName(partCols, partVals);
- partNames.remove(defaultPartitionName);
- Deadline.startTimer("getAggrPartitionColumnStatistics");
- aggrStatsAllButDefaultPartition =
- rawStore.get_aggr_stats_for(catName, dbName, tblName,
partNames, colNames, CacheUtils.HIVE_ENGINE);
- Deadline.stopTimer();
-
cacheObjects.setAggrStatsAllButDefaultPartition(aggrStatsAllButDefaultPartition);
- }
- } else {
- Deadline.startTimer("getTableColumnStatistics");
- tableColStats = rawStore.getTableColumnStatistics(catName,
dbName, tblName, colNames, CacheUtils.HIVE_ENGINE);
- Deadline.stopTimer();
- cacheObjects.setTableColStats(tableColStats);
- }
-
- Deadline.startTimer("getAllTableConstraints");
- SQLAllTableConstraints tableConstraints =
rawStore.getAllTableConstraints(
- new AllTableConstraintsRequest(catName, dbName, tblName));
- Deadline.stopTimer();
- cacheObjects.setTableConstraints(tableConstraints);
-
- // If the table could not cached due to memory limit, stop
prewarm
- boolean isSuccess = sharedCache
- .populateTableInCache(table, cacheObjects);
- if (isSuccess) {
- LOG.trace("Cached Database: {}'s Table: {}.", dbName, tblName);
- } else {
- LOG.info("Unable to cache Database: {}'s Table: {}, since the
cache memory is full. "
- + "Will stop attempting to cache any more tables.",
dbName, tblName);
+ for (Future<?> worker : workers) {
+ try {
+ worker.get();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ LOG.warn("Interrupted while waiting for prewarm workers on
database {}; "
+ + "completing prewarm with the metadata cached so far",
dbName);
completePrewarm(startTime, false);
return;
+ } catch (ExecutionException e) {
+ LOG.warn("Prewarm worker failed for database {}, moving on",
dbName, e);
}
- } catch (MetaException | NoSuchObjectException e) {
- LOG.debug(ExceptionUtils.getStackTrace(e));
- // Continue with next table
- continue;
}
- LOG.debug("Processed database: {}'s table: {}. Cached {} / {}
tables so far.", dbName, tblName,
- ++numberOfTablesCachedSoFar, totalTablesToCache);
- } catch (EmptyStackException e) {
- // We've prewarmed this database, continue with the next one
- continue;
+ } else {
+ drainTablesPendingPrewarm(rawStore, catName, dbName,
cacheMemoryFull, tablesCachedSoFar,
+ totalTablesToCache);
+ }
+ if (cacheMemoryFull.get()) {
+ // The shared cache is full; stop prewarm and serve with whatever
has been cached so far
+ completePrewarm(startTime, false);
+ return;
}
+ LOG.debug("Processed database: {}. Cached {} / {} databases so
far.", dbName, ++numberOfDatabasesCachedSoFar,
+ databases.size());
}
- LOG.debug("Processed database: {}. Cached {} / {} databases so far.",
dbName, ++numberOfDatabasesCachedSoFar,
- databases.size());
+ } finally {
+ shutdownPrewarmWorkers(prewarmPool, workerStores);
}
sharedCache.clearDirtyFlags();
completePrewarm(startTime, true);
}
}
+ /**
+ * Creates a fresh RawStore instance for a prewarm worker thread, mirroring
the way
+ * CacheUpdateMasterWork creates its own store.
+ */
+ private static RawStore createRawStoreForPrewarm(Configuration conf) {
+ String rawStoreClassName = MetastoreConf.getVar(conf,
ConfVars.CACHED_RAW_STORE_IMPL, ObjectStore.class.getName());
+ try {
+ RawStore rs = JavaUtils.getClass(rawStoreClassName,
RawStore.class).newInstance();
+ rs.setConf(conf);
+ return rs;
+ } catch (InstantiationException | IllegalAccessException | MetaException
e) {
+ throw new RuntimeException("Cannot instantiate " + rawStoreClassName, e);
+ }
+ }
+
+ private static void shutdownPrewarmWorkers(ExecutorService prewarmPool,
List<RawStore> workerStores) {
+ if (prewarmPool != null) {
+ prewarmPool.shutdownNow();
+ }
+ for (RawStore workerStore : workerStores) {
+ try {
+ workerStore.shutdown();
+ } catch (RuntimeException e) {
+ LOG.warn("Failed to shut down a prewarm worker RawStore", e);
+ }
+ }
+ }
+
+ /**
+ * Drains tables from tblsPendingPrewarm for the given database, caching
each one, until the
+ * pending list is empty or the shared cache reports its memory limit is
reached. Safe to run
+ * from multiple threads concurrently: the pending-table stack hands out
each table exactly once,
+ * and hot tables promoted by prioritizeTableForPrewarm are picked up by
whichever worker pops next.
+ */
+ private static void drainTablesPendingPrewarm(RawStore rawStore, String
catName, String dbName,
+ AtomicBoolean cacheMemoryFull, AtomicInteger tablesCachedSoFar, int
totalTablesToCache) {
+ // Deadline is thread local; register it for prewarm worker threads
+ Deadline.registerIfNot(1000000);
+ while (!cacheMemoryFull.get() &&
tblsPendingPrewarm.hasMoreTablesToPrewarm()) {
+ String tblName;
+ try {
+ tblName =
StringUtils.normalizeIdentifier(tblsPendingPrewarm.getNextTableNameToPrewarm());
+ } catch (EmptyStackException e) {
+ // Another worker drained the remaining tables between our check and
pop
+ break;
+ }
+ if (shouldCacheTable(catName, dbName, tblName)) {
+ if (!prewarmTable(rawStore, catName, dbName, tblName)) {
+ LOG.info("Unable to cache Database: {}'s Table: {}, since the cache
memory is full. "
+ + "Will stop attempting to cache any more tables.", dbName,
tblName);
+ cacheMemoryFull.set(true);
+ return;
+ }
+ LOG.debug("Processed database: {}'s table: {}. Cached {} / {} tables
so far.", dbName, tblName,
+ tablesCachedSoFar.incrementAndGet(), totalTablesToCache);
+ }
+ }
+ }
+
+ /**
+ * Fetches one table with its partitions, statistics and constraints from
the backing database
+ * and populates it in the shared cache. Returns false only when the shared
cache reports that
+ * its memory limit is reached; a table that vanished or failed to load is
skipped by returning
+ * true so that prewarm continues with the next table.
+ */
+ private static boolean prewarmTable(RawStore rawStore, String catName,
String dbName, String tblName) {
+ Table table;
+ try {
+ table = rawStore.getTable(catName, dbName, tblName);
+ } catch (MetaException e) {
+ LOG.debug(ExceptionUtils.getStackTrace(e));
+ // It is possible the table is deleted during fetching tables of the
database,
+ // in that case, continue with the next table
+ return true;
+ }
+ List<String> colNames = MetaStoreUtils.getColumnNamesForTable(table);
+ try {
+ ColumnStatistics tableColStats = null;
+ List<Partition> partitions = null;
+ List<ColumnStatistics> partitionColStats = null;
+ AggrStats aggrStatsAllPartitions = null;
+ AggrStats aggrStatsAllButDefaultPartition = null;
+ TableCacheObjects cacheObjects = new TableCacheObjects();
+ if (!table.getPartitionKeys().isEmpty()) {
+ Deadline.startTimer("getPartitions");
+ partitions = rawStore.getPartitions(catName, dbName, tblName,
GetPartitionsArgs.getAllPartitions());
+ Deadline.stopTimer();
+ cacheObjects.setPartitions(partitions);
+ List<String> partNames = new ArrayList<>(partitions.size());
+ for (Partition p : partitions) {
+ partNames.add(Warehouse.makePartName(table.getPartitionKeys(),
p.getValues()));
+ }
+ if (!partNames.isEmpty()) {
+ // Get partition column stats for this table
+ Deadline.startTimer("getPartitionColumnStatistics");
+ partitionColStats =
+ rawStore.getPartitionColumnStatistics(catName, dbName, tblName,
partNames, colNames, CacheUtils.HIVE_ENGINE);
+ Deadline.stopTimer();
+ cacheObjects.setPartitionColStats(partitionColStats);
+ // Get aggregate stats for all partitions of a table and for all but
default
+ // partition
+ Deadline.startTimer("getAggrPartitionColumnStatistics");
+ aggrStatsAllPartitions = rawStore.get_aggr_stats_for(catName,
dbName, tblName, partNames, colNames, CacheUtils.HIVE_ENGINE);
+ Deadline.stopTimer();
+ cacheObjects.setAggrStatsAllPartitions(aggrStatsAllPartitions);
+ // Remove default partition from partition names and get aggregate
+ // stats again
+ List<FieldSchema> partKeys = table.getPartitionKeys();
+ String defaultPartitionValue =
+ MetastoreConf.getVar(rawStore.getConf(),
ConfVars.DEFAULTPARTITIONNAME);
+ List<String> partCols = new ArrayList<>();
+ List<String> partVals = new ArrayList<>();
+ for (FieldSchema fs : partKeys) {
+ partCols.add(fs.getName());
+ partVals.add(defaultPartitionValue);
+ }
+ String defaultPartitionName = FileUtils.makePartName(partCols,
partVals);
+ partNames.remove(defaultPartitionName);
+ Deadline.startTimer("getAggrPartitionColumnStatistics");
+ aggrStatsAllButDefaultPartition =
+ rawStore.get_aggr_stats_for(catName, dbName, tblName, partNames,
colNames, CacheUtils.HIVE_ENGINE);
+ Deadline.stopTimer();
+
cacheObjects.setAggrStatsAllButDefaultPartition(aggrStatsAllButDefaultPartition);
+ }
+ } else {
+ Deadline.startTimer("getTableColumnStatistics");
+ tableColStats = rawStore.getTableColumnStatistics(catName, dbName,
tblName, colNames, CacheUtils.HIVE_ENGINE);
+ Deadline.stopTimer();
+ cacheObjects.setTableColStats(tableColStats);
+ }
+
+ Deadline.startTimer("getAllTableConstraints");
+ SQLAllTableConstraints tableConstraints =
rawStore.getAllTableConstraints(
+ new AllTableConstraintsRequest(catName, dbName, tblName));
+ Deadline.stopTimer();
+ cacheObjects.setTableConstraints(tableConstraints);
+
+ // If the table could not be cached due to memory limit, stop prewarm
+ boolean isSuccess = sharedCache
+ .populateTableInCache(table, cacheObjects);
Review Comment:
Fixed. sizeEstimators is now volatile and estimators are added copy-on-write
under a dedicated lock, so concurrent workers never read a map that is being
mutated. The estimate() walk itself stays lock-free
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]