From efe9e4da202fa1e99a53ca43c9c08cce885daafd Mon Sep 17 00:00:00 2001 From: peiwangdb Date: Thu, 5 Aug 2021 06:44:36 -0700 Subject: [PATCH] Heuristic Indexer: automatically refresh index cache up-to-date --- .../hetu/core/heuristicindex/TestHindex.java | 43 +++++++++ .../prestosql/heuristicindex/IndexCache.java | 94 +++++++++++++------ 2 files changed, 110 insertions(+), 27 deletions(-) diff --git a/hetu-heuristic-index/src/test/java/io/hetu/core/heuristicindex/TestHindex.java b/hetu-heuristic-index/src/test/java/io/hetu/core/heuristicindex/TestHindex.java index afd187df6..7072d21f6 100644 --- a/hetu-heuristic-index/src/test/java/io/hetu/core/heuristicindex/TestHindex.java +++ b/hetu-heuristic-index/src/test/java/io/hetu/core/heuristicindex/TestHindex.java @@ -126,6 +126,49 @@ public class TestHindex } } + @Test(dataProvider = "tableData1") + public void testIndexAutoloadCache(String indexType, String queryVariable, String queryValue) + throws Exception + { + System.out.println("Running testIndexAutoloadUpdateCache[indexType: " + indexType + + ", queryVariable: " + queryVariable + ", queryValue: " + queryValue + "]"); + + String tableName = getNewTableName(); + createTable1(tableName); + String testerQuery = "SELECT * FROM " + tableName + " WHERE " + queryVariable + "=" + queryValue; + String indexName = getNewIndexName(); + long threadRefreshRate = 5000; // this is the rate that background thread gets executed + int splitsBeforeIndex = getSplitAndMaterializedResult(testerQuery).getFirst(); + + //create index + assertQuerySucceeds("CREATE INDEX " + indexName + " USING " + + indexType + " ON " + tableName + " (" + queryVariable + ")"); + + Thread.sleep(threadRefreshRate); + int splitsAfterIndex = getSplitAndMaterializedResult(testerQuery).getFirst(); + + //update index + assertQuerySucceeds("INSERT INTO " + tableName + " VALUES(7, 'new1'), (8, 'new2')"); + assertQuerySucceeds("INSERT INTO " + tableName + " VALUES(1, 'test')"); + assertQuerySucceeds("INSERT INTO " + tableName + " VALUES(2, '123'), (3, 'temp')"); + assertQuerySucceeds("INSERT INTO " + tableName + " VALUES(3, 'data'), (9, 'ttt'), (5, 'num')"); + assertQuerySucceeds("UPDATE INDEX " + indexName); + Thread.sleep(threadRefreshRate); + int splitsAfterIndexUpdate = getSplitAndMaterializedResult(testerQuery).getFirst(); + + //drop index + assertQuerySucceeds("DROP INDEX " + indexName); + Thread.sleep(threadRefreshRate); + int splitsDropIndex = getSplitAndMaterializedResult(testerQuery).getFirst(); + + assertTrue(splitsBeforeIndex > splitsAfterIndex, "splits should be fewer after index creation"); + assertTrue(splitsBeforeIndex > splitsAfterIndexUpdate, "splits should be more after adding more data"); + assertTrue(splitsBeforeIndex < splitsDropIndex, "splits be more even without any index after dropping index due to data addition"); + assertTrue(splitsAfterIndex < splitsAfterIndexUpdate, "splits should be more after adding data and update index"); + assertTrue(splitsAfterIndex < splitsDropIndex, "splits should be more anyway after dropping index because data volume increased"); + assertTrue(splitsAfterIndexUpdate < splitsDropIndex, "the number of splits after dropping index is the largest"); + } + // Tests data consistency and splits for which table data is changed after index creation. @Test(dataProvider = "tableData1") public void testDataConsistencyWithAdditionChange(String indexType, String queryVariable, String queryValue) diff --git a/presto-main/src/main/java/io/prestosql/heuristicindex/IndexCache.java b/presto-main/src/main/java/io/prestosql/heuristicindex/IndexCache.java index 5f2c064f9..0c4e2a9b6 100644 --- a/presto-main/src/main/java/io/prestosql/heuristicindex/IndexCache.java +++ b/presto-main/src/main/java/io/prestosql/heuristicindex/IndexCache.java @@ -38,8 +38,10 @@ import java.nio.file.Path; import java.nio.file.Paths; import java.util.ArrayList; import java.util.Collections; +import java.util.HashMap; import java.util.LinkedList; import java.util.List; +import java.util.Locale; import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.concurrent.Executors; @@ -98,35 +100,14 @@ public class IndexCache if (PropertyService.getBooleanProperty(HetuConstant.FILTER_CACHE_SOFT_REFERENCE)) { cacheBuilder.softValues(); } - // Refresh cache according to index records in the background. Evict index from cache if it's dropped. + executor.scheduleAtFixedRate(() -> { + // This thread automatically keep the cache updated every 5 secs in the background. try { - if (cache.size() > 0) { - // only refresh cache is it's not empty - List newRecords = indexClient.getAllIndexRecords(); - - if (indexRecords != null) { - for (IndexRecord old : indexRecords) { - boolean found = false; - for (IndexRecord now : newRecords) { - if (now.name.equals(old.name)) { - found = true; - if (now.lastModifiedTime != old.lastModifiedTime) { - // index record has been updated. evict - evictFromCache(old); - LOG.debug("Index for {%s} has been evicted from cache because the index has been updated.", old); - } - } - } - // old record is gone. evict from cache - if (!found) { - evictFromCache(old); - LOG.debug("Index for {%s} has been evicted from cache because the index has been dropped.", old); - } - } - } - - indexRecords = newRecords; + List newRecords = indexClient.getAllIndexRecords(); + boolean success = autoUpdateCache(newRecords); + if (success) { + LOG.debug("Cache refreshed"); } } catch (Exception e) { @@ -137,6 +118,65 @@ public class IndexCache } } + private boolean autoUpdateCache(List newRecords) + { + // Three cases are checked: + // 1. if new index record is created, add it to cache + // 2. if the index in cache is updated, update the cache + // 3. if the index in cache is outdated, evict it from cache + + if (newRecords == null && indexRecords == null) { + return false; + } + + HashMap oldIndexMap = new HashMap<>(); + if (indexRecords != null) { + for (IndexRecord oldIndexRecord : indexRecords) { + oldIndexMap.put(oldIndexRecord.name, oldIndexRecord.lastModifiedTime); + } + } + boolean dropped = false; + boolean created = false; + boolean updated = false; + HashMap newIndexMap = new HashMap<>(); + if (newRecords != null) { + for (IndexRecord newIndexRecord : newRecords) { + newIndexMap.put(newIndexRecord.name, newIndexRecord.lastModifiedTime); + if (oldIndexMap.containsKey(newIndexRecord.name)) { + if (oldIndexMap.get(newIndexRecord.name) != newIndexRecord.lastModifiedTime) { + // update operation + updated = true; + //cache.refresh(newIndexRecord); + evictFromCache(newIndexRecord); + CreateIndexMetadata.Level indexLevel = CreateIndexMetadata.Level.valueOf(newIndexRecord.getProperty(CreateIndexMetadata.LEVEL_PROP_KEY).toUpperCase(Locale.ROOT)); + preloadIndex(newIndexRecord.qualifiedTable, String.join(",", newIndexRecord.columns), newIndexRecord.indexType, indexLevel); + LOG.debug("Index {%s} has been updated in cache.", newIndexRecord); + } + } + else { + // create operation + created = true; + CreateIndexMetadata.Level indexLevel = CreateIndexMetadata.Level.valueOf(newIndexRecord.getProperty(CreateIndexMetadata.LEVEL_PROP_KEY).toUpperCase(Locale.ROOT)); + preloadIndex(newIndexRecord.qualifiedTable, String.join(",", newIndexRecord.columns), newIndexRecord.indexType, indexLevel); + LOG.debug("Index {%s} has been inserted to cache.", newIndexRecord); + } + } + } + + if (indexRecords != null) { + for (IndexRecord oldIndexRecord : indexRecords) { + if (!newIndexMap.containsKey(oldIndexRecord.name)) { + // drop operation + dropped = true; + evictFromCache(oldIndexRecord); + LOG.debug("Index {%s} has been evicted from cache because the index has been dropped.", oldIndexRecord); + } + } + } + indexRecords = newRecords; + return (dropped || created || updated); + } + public void preloadIndex(String table, String column, String type, CreateIndexMetadata.Level level) { String filterKeyPath = table + "/" + column + "/" + type;