From 30e67604554e6689357bf454a66fa0bf5582cc5d Mon Sep 17 00:00:00 2001 From: Minwoo Kang Date: Fri, 31 Jul 2026 22:14:23 +0900 Subject: [PATCH] HBASE-30285 BucketCache shutdown/disable does not release backingMap-owned BucketEntry references (#8469) Signed-off-by: Wellington Chevreuil (cherry picked from commit f07a9fe1d89637898c8834596b122ce317d387d5) --- .../hbase/io/hfile/bucket/BucketCache.java | 221 +++++++---- .../io/hfile/bucket/BucketProtoUtils.java | 8 + .../TestBlockEvictionOnRegionMovement.java | 29 +- .../io/hfile/bucket/TestBucketCache.java | 38 +- .../hfile/bucket/TestBucketCacheRefCnt.java | 349 ++++++++++++++++++ 5 files changed, 560 insertions(+), 85 deletions(-) diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketCache.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketCache.java index d48abb363632..d4c6d82db798 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketCache.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketCache.java @@ -90,6 +90,7 @@ import org.apache.hadoop.hbase.util.IdReadWriteLock; import org.apache.hadoop.hbase.util.IdReadWriteLock.ReferenceType; import org.apache.hadoop.hbase.util.Pair; +import org.apache.hadoop.hbase.util.Threads; import org.apache.hadoop.util.StringUtils; import org.apache.yetus.audience.InterfaceAudience; import org.slf4j.Logger; @@ -200,6 +201,12 @@ protected enum CacheState { */ private volatile CacheState cacheState; + /** The single cleanup thread shared by disable and explicit shutdown calls. */ + private volatile Thread cacheCleanupThread; + + /** The thread restoring the persistent cache index during initialization. */ + private volatile Thread persistenceRetrieverThread; + /** * A list of writer queues. We have a queue per {@link WriterThread} we have running. In other * words, the work adding blocks to the BucketCache is divided up amongst the running @@ -403,12 +410,17 @@ private void startPersistenceRetriever(int[] bucketSizes, long capacity) { LOG.error("Exception during Bucket Allocation", allocatorException); } } finally { - this.cacheState = CacheState.ENABLED; - startWriterThreads(); + synchronized (BucketCache.this) { + if (cacheState == CacheState.INITIALIZING) { + cacheState = CacheState.ENABLED; + startWriterThreads(); + } + } } }; - Thread t = new Thread(persistentCacheRetriever); - t.start(); + persistenceRetrieverThread = new Thread(persistentCacheRetriever, + "BucketCachePersistenceRetriever-" + System.identityHashCode(this)); + persistenceRetrieverThread.start(); } private void sanityCheckConfigs() { @@ -630,21 +642,24 @@ protected void cacheBlockWithWaitInternal(BlockCacheKey cacheKey, Cacheable cach * the passed key doesn't relate to a reference. */ public BucketEntry getBlockForReference(BlockCacheKey key) { - BucketEntry foundEntry = null; - String referredFileName = null; - if (StoreFileInfo.isReference(key.getHfileName())) { - referredFileName = StoreFileInfo.getReferredToRegionAndFile(key.getHfileName()).getSecond(); - } - if (referredFileName != null) { - // Since we just need this key for a lookup, it's enough to use only name and offset - BlockCacheKey convertedCacheKey = new BlockCacheKey(referredFileName, key.getOffset()); - foundEntry = backingMap.get(convertedCacheKey); + BlockCacheKey referredKey = getBlockKeyForReference(key); + BucketEntry foundEntry = referredKey != null ? backingMap.get(referredKey) : null; + if (referredKey != null) { LOG.debug("Got a link/ref: {}. Related cacheKey: {}. Found entry: {}", key.getHfileName(), - convertedCacheKey, foundEntry); + referredKey, foundEntry); } return foundEntry; } + private BlockCacheKey getBlockKeyForReference(BlockCacheKey key) { + if (!StoreFileInfo.isReference(key.getHfileName())) { + return null; + } + String referredFileName = + StoreFileInfo.getReferredToRegionAndFile(key.getHfileName()).getSecond(); + return referredFileName != null ? new BlockCacheKey(referredFileName, key.getOffset()) : null; + } + /** * Get the buffer of the block with the specified key. * @param key block's cache key @@ -668,23 +683,28 @@ public Cacheable getBlock(BlockCacheKey key, boolean caching, boolean repeat, re.access(accessCount.incrementAndGet()); return re.getData(); } - BucketEntry bucketEntry = backingMap.get(key); + BlockCacheKey backingMapLookupKey = key; + BucketEntry bucketEntry = backingMap.get(backingMapLookupKey); LOG.debug("bucket entry for key {}: {}", key, bucketEntry == null ? null : bucketEntry.offset()); if (bucketEntry == null) { - bucketEntry = getBlockForReference(key); + backingMapLookupKey = getBlockKeyForReference(key); + if (backingMapLookupKey != null) { + bucketEntry = backingMap.get(backingMapLookupKey); + LOG.debug("Got a link/ref: {}. Related cacheKey: {}. Found entry: {}", key.getHfileName(), + backingMapLookupKey, bucketEntry); + } } if (bucketEntry != null) { long start = System.nanoTime(); ReentrantReadWriteLock lock = offsetLock.getLock(bucketEntry.offset()); + boolean inconsistentEntry = false; try { lock.readLock().lock(); // We can not read here even if backingMap does contain the given key because its offset // maybe changed. If we lock BlockCacheKey instead of offset, then we can only check // existence here. - if ( - bucketEntry.equals(backingMap.get(key)) || bucketEntry.equals(getBlockForReference(key)) - ) { + if (bucketEntry.equals(backingMap.get(backingMapLookupKey))) { // Read the block from IOEngine based on the bucketEntry's offset and length, NOTICE: the // block will use the refCnt of bucketEntry, which means if two HFileBlock mapping to // the same BucketEntry, then all of the three will share the same refCnt. @@ -706,11 +726,9 @@ public Cacheable getBlock(BlockCacheKey key, boolean caching, boolean repeat, return cachedBlock; } } catch (HBaseIOException hioex) { - // When using file io engine persistent cache, - // the cache map state might differ from the actual cache. If we reach this block, - // we should remove the cache key entry from the backing map - backingMap.remove(key); - fileNotFullyCached(key, bucketEntry); + // FileIOEngine throws this when its cached time differs from the persisted index. A plain + // IOException still follows the configured tolerance policy below. + inconsistentEntry = true; LOG.debug("Failed to fetch block for cache key: {}.", key, hioex); } catch (IOException ioex) { LOG.error("Failed reading block " + key + " from bucket cache", ioex); @@ -718,6 +736,9 @@ public Cacheable getBlock(BlockCacheKey key, boolean caching, boolean repeat, } finally { lock.readLock().unlock(); } + if (inconsistentEntry) { + evictInconsistentEntry(backingMapLookupKey, bucketEntry); + } } if (!repeat && updateCacheMetrics) { cacheStats.miss(caching, key.isPrimary(), key.getBlockType()); @@ -725,6 +746,19 @@ public Cacheable getBlock(BlockCacheKey key, boolean caching, boolean repeat, return null; } + private void evictInconsistentEntry(BlockCacheKey lookupKey, BucketEntry bucketEntry) { + BlockCacheKey storedKey = blocksByHFile.ceiling(lookupKey); + if (storedKey == null || !storedKey.equals(lookupKey)) { + return; + } + bucketEntry.withWriteLock(offsetLock, () -> { + if (backingMap.remove(storedKey, bucketEntry)) { + blockEvicted(storedKey, bucketEntry, true, false); + } + return null; + }); + } + /** * This method is invoked after the bucketEntry is removed from {@link BucketCache#backingMap} */ @@ -1342,7 +1376,6 @@ public void run() { */ protected void putIntoBackingMap(BlockCacheKey key, BucketEntry bucketEntry) { BucketEntry previousEntry = backingMap.put(key, bucketEntry); - blocksByHFile.add(key); updateRegionCachedSize(key, bucketEntry.getLength()); if (previousEntry != null && previousEntry != bucketEntry) { previousEntry.withWriteLock(offsetLock, () -> { @@ -1350,6 +1383,12 @@ protected void putIntoBackingMap(BlockCacheKey key, BucketEntry bucketEntry) { return null; }); } + bucketEntry.withWriteLock(offsetLock, () -> { + if (backingMap.get(key) == bucketEntry) { + blocksByHFile.add(key); + } + return null; + }); } /** @@ -1543,6 +1582,12 @@ static List getRAMQueueEntries(BlockingQueue q, @edu.umd.cs.findbugs.annotations.SuppressWarnings(value = "OBL_UNSATISFIED_OBLIGATION", justification = "false positive, try-with-resources ensures close is called.") void persistToFile() throws IOException { + persistToFile(entry -> { + }); + } + + private void persistToFile(Consumer> entryCopiedAction) + throws IOException { LOG.debug("Thread {} started persisting bucket cache to file", Thread.currentThread().getName()); if (!isCachePersistent()) { @@ -1552,7 +1597,7 @@ void persistToFile() throws IOException { try (FileOutputStream fos = new FileOutputStream(tempPersistencePath, false)) { LOG.debug("Persist in new chunked persistence format."); - persistChunkedBackingMap(fos); + persistChunkedBackingMap(fos, entryCopiedAction); LOG.debug( "PersistToFile: after persisting backing map size: {}, fullycachedFiles size: {}," @@ -1732,13 +1777,14 @@ private void parsePB(BucketCacheProtos.BucketCacheEntry proto) throws IOExceptio verifyCapacityAndClasses(proto.getCacheCapacity(), proto.getIoClass(), proto.getMapClass()); } - private void persistChunkedBackingMap(FileOutputStream fos) throws IOException { + private void persistChunkedBackingMap(FileOutputStream fos, + Consumer> entryCopiedAction) throws IOException { LOG.debug( "persistToFile: before persisting backing map size: {}, " + "fullycachedFiles size: {}, chunkSize: {}", backingMap.size(), fullyCachedFiles.size(), persistenceChunkSize); - BucketProtoUtils.serializeAsPB(this, fos, persistenceChunkSize); + BucketProtoUtils.serializeAsPB(this, fos, persistenceChunkSize, entryCopiedAction); LOG.debug( "persistToFile: after persisting backing map size: {}, " + "fullycachedFiles size: {}", @@ -1799,59 +1845,106 @@ private void checkIOErrorIsTolerated() { /** * Used to shut down the cache -or- turn it off in the case of something broken. + * @return whether explicit shutdown should wait for cleanup */ - private void disableCache() { - if (!isCacheEnabled()) { - return; + private synchronized boolean disableCache() { + if (cacheState == CacheState.DISABLED) { + return false; } + boolean waitForCleanup = cacheState == CacheState.ENABLED && isCachePersistent(); LOG.info("Disabling cache"); cacheState = CacheState.DISABLED; - ioEngine.shutdown(); this.scheduleThreadPool.shutdown(); - for (int i = 0; i < writerThreads.length; ++i) - writerThreads[i].interrupt(); - this.ramCache.clear(); - if (!ioEngine.isPersistent() || persistencePath == null) { - // If persistent ioengine and a path, we will serialize out the backingMap. - this.backingMap.clear(); - this.blocksByHFile.clear(); - this.fullyCachedFiles.clear(); - this.regionCachedSize.clear(); + for (WriterThread writerThread : writerThreads) { + writerThread.interrupt(); } + // Closing the IO engine helps unblock an in-flight writer before the cleanup thread joins it. + // FileIOEngine can reopen a channel, so cleanup closes the engine again after writers stop. + ioEngine.shutdown(); if (cacheStats.getMetricsRollerScheduler() != null) { cacheStats.getMetricsRollerScheduler().shutdownNow(); } + cacheCleanupThread = Threads.setDaemonThreadRunning(new Thread(this::cleanupCache), + "BucketCacheCleanup-" + System.identityHashCode(this), Threads.LOGGING_EXCEPTION_HANDLER); + return waitForCleanup; } - private void join() throws InterruptedException { - for (int i = 0; i < writerThreads.length; ++i) - writerThreads[i].join(); - } - - @Override - public void shutdown() { - if (isCacheEnabled()) { - disableCache(); - LOG.info("Shutdown bucket cache: IO persistent=" + ioEngine.isPersistent() - + "; path to write=" + persistencePath); - if (ioEngine.isPersistent() && persistencePath != null) { + private void cleanupCache() { + try { + Threads.shutdown(persistenceRetrieverThread); + for (WriterThread writerThread : writerThreads) { + Threads.shutdown(writerThread); + } + for (BlockingQueue writerQueue : writerQueues) { + writerQueue.clear(); + } + ramCache.clear(); + if (cachePersister != null) { + LOG.info("Shutting down cache persister thread."); + cachePersister.shutdown(); + Threads.shutdown(cachePersister); + } + if (isCachePersistent()) { try { - join(); - if (cachePersister != null) { - LOG.info("Shutting down cache persister thread."); - cachePersister.shutdown(); - while (cachePersister.isAlive()) { - Thread.sleep(10); - } - } - persistToFile(); + // The serializer already visits every entry. Release owner references in the same pass. + persistToFile(this::cleanupBackingMapEntry); } catch (IOException ex) { LOG.error("Unable to persist data on exit: " + ex.toString(), ex); - } catch (InterruptedException e) { - LOG.warn("Failed to persist data on exit", e); } } + } finally { + try { + cleanupCacheIndex(); + } finally { + ioEngine.shutdown(); + } + } + } + + private void cleanupCacheIndex() { + // A successful persistent cleanup emptied the map during serialization. Avoid creating a + // second iterator over a large ConcurrentHashMap in that case. + if (!backingMap.isEmpty()) { + for (Map.Entry entry : backingMap.entrySet()) { + cleanupBackingMapEntry(entry); + } + } + blocksByHFile.clear(); + fullyCachedFiles.clear(); + regionCachedSize.clear(); + } + + private void cleanupBackingMapEntry(Map.Entry entry) { + BlockCacheKey cacheKey = entry.getKey(); + BucketEntry bucketEntry = entry.getValue(); + bucketEntry.withWriteLock(offsetLock, () -> { + if (backingMap.remove(cacheKey, bucketEntry)) { + bucketEntry.markAsEvicted(); + } + return null; + }); + } + + private void waitForCacheCleanup() throws InterruptedException { + Thread cleanupThread = cacheCleanupThread; + if (cleanupThread == null || cleanupThread == Thread.currentThread()) { + return; + } + cleanupThread.join(); + } + + @Override + public void shutdown() { + if (disableCache()) { + try { + waitForCacheCleanup(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + LOG.warn("Interrupted while waiting for bucket cache cleanup", e); + } } + LOG.info("Shutdown bucket cache: IO persistent=" + ioEngine.isPersistent() + "; path to write=" + + persistencePath); } /** diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketProtoUtils.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketProtoUtils.java index 95808a7a8554..1488d949deff 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketProtoUtils.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/bucket/BucketProtoUtils.java @@ -25,6 +25,7 @@ import java.util.NavigableSet; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentSkipListSet; +import java.util.function.Consumer; import java.util.function.Function; import org.apache.hadoop.hbase.io.ByteBuffAllocator; import org.apache.hadoop.hbase.io.ByteBuffAllocator.Recycler; @@ -64,6 +65,12 @@ static BucketCacheProtos.BucketCacheEntry toPB(BucketCache cache, public static void serializeAsPB(BucketCache cache, FileOutputStream fos, long chunkSize) throws IOException { + serializeAsPB(cache, fos, chunkSize, entry -> { + }); + } + + static void serializeAsPB(BucketCache cache, FileOutputStream fos, long chunkSize, + Consumer> entryCopiedAction) throws IOException { // Write the new version of magic number. fos.write(PB_MAGIC_V2); @@ -79,6 +86,7 @@ public static void serializeAsPB(BucketCache cache, FileOutputStream fos, long c for (Map.Entry entry : cache.backingMap.entrySet()) { blockCount++; addEntryToBuilder(entry, entryBuilder, builder); + entryCopiedAction.accept(entry); if (blockCount % chunkSize == 0) { builder.build().writeDelimitedTo(fos); builder.clear(); diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestBlockEvictionOnRegionMovement.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestBlockEvictionOnRegionMovement.java index 979a5f9c641d..05247e9e1267 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestBlockEvictionOnRegionMovement.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestBlockEvictionOnRegionMovement.java @@ -28,6 +28,7 @@ import org.apache.hadoop.fs.Path; import org.apache.hadoop.hbase.HBaseTestingUtility; import org.apache.hadoop.hbase.MiniHBaseCluster; +import org.apache.hadoop.hbase.ServerName; import org.apache.hadoop.hbase.StartMiniClusterOption; import org.apache.hadoop.hbase.TableName; import org.apache.hadoop.hbase.Waiter; @@ -125,20 +126,24 @@ public void testBlockEvictionOnGracefulStop() throws Exception { : cluster.getRegionServer(0); assertTrue(regionServingRS.getBlockCache().isPresent()); - long oldUsedCacheSize = - regionServingRS.getBlockCache().get().getBlockCaches()[1].getCurrentSize(); - assertNotEquals(0, regionServingRS.getBlockCache().get().getBlockCaches()[1].getBlockCount()); + BlockCache oldBucketCache = regionServingRS.getBlockCache().get().getBlockCaches()[1]; + long oldUsedCacheSize = oldBucketCache.getCurrentSize(); + assertNotEquals(0, oldUsedCacheSize); + assertNotEquals(0, oldBucketCache.getBlockCount()); - cluster.stopRegionServer(regionServingRS.getServerName()); - Thread.sleep(500); - cluster.startRegionServer(); - Thread.sleep(500); + ServerName serverName = regionServingRS.getServerName(); + cluster.stopRegionServer(serverName); + cluster.waitForRegionServerToStop(serverName, 10000); - regionServingRS.getBlockCache().get().waitForCacheInitialization(10000); - long newUsedCacheSize = - regionServingRS.getBlockCache().get().getBlockCaches()[1].getCurrentSize(); - assertEquals(oldUsedCacheSize, newUsedCacheSize); - assertNotEquals(0, regionServingRS.getBlockCache().get().getBlockCaches()[1].getBlockCount()); + assertEquals(0, oldBucketCache.getCurrentSize()); + + HRegionServer restartedRegionServer = cluster.startRegionServer().getRegionServer(); + assertTrue(restartedRegionServer.getBlockCache().isPresent()); + BlockCache restoredBucketCache = + restartedRegionServer.getBlockCache().get().getBlockCaches()[1]; + assertTrue(restoredBucketCache.waitForCacheInitialization(10000)); + assertEquals(oldUsedCacheSize, restoredBucketCache.getCurrentSize()); + assertNotEquals(0, restoredBucketCache.getBlockCount()); } public TableName writeDataToTable(String testName) throws IOException, InterruptedException { diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCache.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCache.java index 3446d78ee17d..08dcbd8e29c5 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCache.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCache.java @@ -35,6 +35,7 @@ import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; import static org.mockito.Mockito.mock; @@ -275,23 +276,36 @@ public void testMemoryLeak() throws Exception { final BlockCacheKey cacheKey = new BlockCacheKey("dummy", 1L); cacheAndWaitUntilFlushedToBucket(cache, cacheKey, new CacheTestUtils.ByteArrayCacheable(new byte[10]), true); - long lockId = cache.backingMap.get(cacheKey).offset(); + BucketEntry oldEntry = cache.backingMap.get(cacheKey); + long lockId = oldEntry.offset(); ReentrantReadWriteLock lock = cache.offsetLock.getLock(lockId); - lock.writeLock().lock(); Thread evictThread = new Thread("evict-block") { @Override public void run() { cache.evictBlock(cacheKey); } }; - evictThread.start(); - cache.offsetLock.waitForWaiters(lockId, 1); - cache.blockEvicted(cacheKey, cache.backingMap.remove(cacheKey), true, true); - assertEquals(0, cache.getBlockCount()); - cacheAndWaitUntilFlushedToBucket(cache, cacheKey, - new CacheTestUtils.ByteArrayCacheable(new byte[10]), true); + BucketEntry replacementEntry; + lock.writeLock().lock(); + try { + evictThread.start(); + cache.offsetLock.waitForWaiters(lockId, 1); + assertTrue(cache.backingMap.remove(cacheKey, oldEntry)); + cache.blockEvicted(cacheKey, oldEntry, true, true); + assertEquals(0, cache.getBlockCount()); + cache.cacheBlock(cacheKey, new CacheTestUtils.ByteArrayCacheable(new byte[10]), false, true); + // The replacement is published before taking its offset lock. Full flush must wait until the + // old entry's lock is released. + Waiter.waitFor(HBaseConfiguration.create(), 10000, () -> { + BucketEntry currentEntry = cache.backingMap.get(cacheKey); + return currentEntry != null && currentEntry != oldEntry; + }); + replacementEntry = cache.backingMap.get(cacheKey); + } finally { + lock.writeLock().unlock(); + } + waitUntilFlushedToBucket(cache, cacheKey); assertEquals(1, cache.getBlockCount()); - lock.writeLock().unlock(); evictThread.join(); /** *
@@ -310,6 +324,8 @@ public void run() {
      * it had seen should not be evicted.
      * 
*/ + assertSame(replacementEntry, cache.backingMap.get(cacheKey)); + assertTrue(cache.blocksByHFile.contains(cacheKey)); assertEquals(1L, cache.getBlockCount()); assertTrue(cache.getCurrentSize() > 0L); assertTrue(cache.iterator().hasNext(), "We should have a block!"); @@ -385,7 +401,11 @@ private void testRetrievalUtils(Path testDir, String ioEngineName) } usedSize = bucketCache.getAllocator().getUsedSize(); assertNotEquals(0, usedSize); + BucketEntry persistedEntry = bucketCache.backingMap.values().iterator().next(); + assertEquals(1, persistedEntry.refCnt()); bucketCache.shutdown(); + assertEquals(0, persistedEntry.refCnt()); + assertTrue(bucketCache.backingMap.isEmpty()); assertTrue(new File(persistencePath).exists()); bucketCache = new BucketCache(ioEngineName, capacitySize, constructedBlockSize, constructedBlockSizes, writeThreads, writerQLen, persistencePath); diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCacheRefCnt.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCacheRefCnt.java index b466719136c0..d6adb35d9cca 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCacheRefCnt.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestBucketCacheRefCnt.java @@ -17,19 +17,34 @@ */ package org.apache.hadoop.hbase.io.hfile.bucket; +import static org.apache.hadoop.hbase.io.hfile.CacheConfig.BUCKETCACHE_PERSIST_INTERVAL_KEY; +import static org.apache.hadoop.hbase.io.hfile.bucket.BucketCache.DEFAULT_ERROR_TOLERATION_DURATION; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; +import java.io.DataInputStream; +import java.io.File; +import java.io.FileInputStream; import java.io.IOException; import java.nio.ByteBuffer; +import java.nio.channels.FileChannel; +import java.nio.file.Files; import java.util.Arrays; import java.util.List; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; import org.apache.hadoop.hbase.HBaseConfiguration; import org.apache.hadoop.hbase.Waiter; import org.apache.hadoop.hbase.io.ByteBuffAllocator; @@ -45,9 +60,15 @@ import org.apache.hadoop.hbase.nio.RefCnt; import org.apache.hadoop.hbase.testclassification.IOTests; import org.apache.hadoop.hbase.testclassification.SmallTests; +import org.apache.hadoop.hbase.util.EnvironmentEdgeManager; +import org.apache.hadoop.hbase.util.ManualEnvironmentEdge; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.mockito.Mockito; + +import org.apache.hadoop.hbase.shaded.protobuf.generated.BucketCacheProtos; @Tag(IOTests.TAG) @Tag(SmallTests.TAG) @@ -306,6 +327,333 @@ public void testMarkStaleAsEvicted() throws Exception { } } + @Test + public void testShutdownReleasesOnlyBackingMapReference() throws Exception { + ByteBuffAllocator allocator = ByteBuffAllocator.create(HBaseConfiguration.create(), true); + HFileBlock blockToCache = createBlock(200, 1020, allocator); + HFileBlock retainedBlock = null; + try { + cache = create(1, 1000); + BlockCacheKey key = createKey("testShutdownReleasesOnlyBackingMapReference", 200); + cache.cacheBlock(key, blockToCache); + waitUntilFlushedToCache(cache, key); + + retainedBlock = (HFileBlock) cache.getBlock(key, false, false, false); + assertNotNull(retainedBlock); + assertEquals(2, retainedBlock.refCnt()); + + HFileBlock callerReference = retainedBlock; + cache.shutdown(); + cache.shutdown(); + + Waiter.waitFor(HBaseConfiguration.create(), 10000, + () -> cache.backingMap.isEmpty() && callerReference.refCnt() == 1); + assertEquals(1, retainedBlock.refCnt()); + assertTrue(retainedBlock.release()); + assertEquals(0, retainedBlock.refCnt()); + } finally { + if (cache != null) { + cache.shutdown(); + cache = null; + } + if (retainedBlock != null) { + while (retainedBlock.refCnt() > 0) { + retainedBlock.release(); + } + } + while (blockToCache.refCnt() > 0) { + blockToCache.release(); + } + allocator.clean(); + } + } + + @Test + public void testShutdownPersistsEmptyMapOverPreviousCheckpoint(@TempDir File testDir) + throws Exception { + HFileBlock blockToCache = createBlock(200, 1020); + BucketCache bucketCache = null; + BucketCache recoveredCache = null; + String cachePath = new File(testDir, "bucket.cache").getAbsolutePath(); + String persistencePath = new File(testDir, "bucket.persistence").getAbsolutePath(); + BlockCacheKey key = createKey("testShutdownPersistsEmptyMapOverPreviousCheckpoint", 200); + Configuration conf = HBaseConfiguration.create(); + conf.setLong(BUCKETCACHE_PERSIST_INTERVAL_KEY, Long.MAX_VALUE); + try { + bucketCache = new BucketCache("file:" + cachePath, CAPACITY_SIZE, BLOCK_SIZE, + BLOCK_SIZE_ARRAY, 1, 1000, persistencePath, DEFAULT_ERROR_TOLERATION_DURATION, conf); + assertTrue(bucketCache.waitForCacheInitialization(10000)); + bucketCache.cacheBlock(key, blockToCache); + waitUntilFlushedToCache(bucketCache, key); + + bucketCache.persistToFile(); + assertTrue(new File(persistencePath).isFile()); + assertTrue(bucketCache.backingMap.containsKey(key)); + assertTrue(bucketCache.evictBlock(key)); + assertTrue(bucketCache.backingMap.isEmpty()); + bucketCache.shutdown(); + try ( + DataInputStream in = new DataInputStream(new FileInputStream(new File(persistencePath)))) { + byte[] magic = new byte[BucketProtoUtils.PB_MAGIC_V2.length]; + in.readFully(magic); + assertArrayEquals(BucketProtoUtils.PB_MAGIC_V2, magic); + assertNotNull(BucketCacheProtos.BucketCacheEntry.parseDelimitedFrom(in)); + assertEquals(-1, in.read()); + } + bucketCache = null; + + recoveredCache = new BucketCache("file:" + cachePath, CAPACITY_SIZE, BLOCK_SIZE, + BLOCK_SIZE_ARRAY, 1, 1000, persistencePath, DEFAULT_ERROR_TOLERATION_DURATION, conf); + assertTrue(recoveredCache.waitForCacheInitialization(10000)); + assertTrue(recoveredCache.backingMap.isEmpty()); + assertEquals(0, recoveredCache.getAllocator().getUsedSize()); + assertEquals(0, recoveredCache.getBlockCount()); + Cacheable staleBlock = recoveredCache.getBlock(key, false, false, false); + if (staleBlock != null) { + staleBlock.release(); + } + assertNull(staleBlock); + } finally { + if (bucketCache != null) { + bucketCache.shutdown(); + } + if (recoveredCache != null) { + recoveredCache.shutdown(); + } + while (blockToCache.refCnt() > 0) { + blockToCache.release(); + } + } + } + + @Test + public void testFinalPersistFailureReleasesEntriesAndKeepsPreviousCheckpoint( + @TempDir File testDir) throws Exception { + HFileBlock firstBlock = createBlock(200, 1020); + HFileBlock secondBlock = createBlock(400, 1020); + BucketCache bucketCache = null; + BucketCache recoveredCache = null; + String cachePath = new File(testDir, "bucket.cache").getAbsolutePath(); + String persistencePath = new File(testDir, "bucket.persistence").getAbsolutePath(); + BlockCacheKey firstKey = createKey("first", 200); + BlockCacheKey secondKey = createKey("second", 400); + Configuration conf = HBaseConfiguration.create(); + conf.setLong(BUCKETCACHE_PERSIST_INTERVAL_KEY, Long.MAX_VALUE); + try { + bucketCache = new BucketCache("file:" + cachePath, CAPACITY_SIZE, BLOCK_SIZE, + BLOCK_SIZE_ARRAY, 1, 1000, persistencePath, DEFAULT_ERROR_TOLERATION_DURATION, conf); + assertTrue(bucketCache.waitForCacheInitialization(10000)); + bucketCache.cacheBlock(firstKey, firstBlock); + waitUntilFlushedToCache(bucketCache, firstKey); + bucketCache.persistToFile(); + byte[] previousCheckpoint = Files.readAllBytes(new File(persistencePath).toPath()); + + bucketCache.cacheBlock(secondKey, secondBlock); + waitUntilFlushedToCache(bucketCache, secondKey); + BucketEntry firstEntry = bucketCache.backingMap.get(firstKey); + BucketEntry secondEntry = bucketCache.backingMap.get(secondKey); + assertNotNull(firstEntry); + assertNotNull(secondEntry); + assertEquals(1, firstEntry.refCnt()); + assertEquals(1, secondEntry.refCnt()); + + long failureTime = 123456789L; + File tempPersistencePath = new File(persistencePath + failureTime); + assertTrue(tempPersistencePath.mkdir()); + ManualEnvironmentEdge edge = new ManualEnvironmentEdge(); + edge.setValue(failureTime); + EnvironmentEdgeManager.injectEdge(edge); + try { + bucketCache.shutdown(); + } finally { + EnvironmentEdgeManager.reset(); + } + + assertTrue(bucketCache.backingMap.isEmpty()); + assertEquals(0, firstEntry.refCnt()); + assertEquals(0, secondEntry.refCnt()); + assertArrayEquals(previousCheckpoint, Files.readAllBytes(new File(persistencePath).toPath())); + bucketCache = null; + + recoveredCache = new BucketCache("file:" + cachePath, CAPACITY_SIZE, BLOCK_SIZE, + BLOCK_SIZE_ARRAY, 1, 1000, persistencePath, DEFAULT_ERROR_TOLERATION_DURATION, conf); + assertTrue(recoveredCache.waitForCacheInitialization(10000)); + BucketCache cacheToValidate = recoveredCache; + Waiter.waitFor(HBaseConfiguration.create(), 10000, + () -> cacheToValidate.getBackingMapValidated().get()); + assertEquals(1, recoveredCache.backingMap.size()); + Cacheable recoveredBlock = recoveredCache.getBlock(firstKey, false, false, false); + assertNotNull(recoveredBlock); + recoveredBlock.release(); + Cacheable uncheckpointedBlock = recoveredCache.getBlock(secondKey, false, false, false); + if (uncheckpointedBlock != null) { + uncheckpointedBlock.release(); + } + assertNull(uncheckpointedBlock); + } finally { + EnvironmentEdgeManager.reset(); + if (bucketCache != null) { + bucketCache.shutdown(); + } + if (recoveredCache != null) { + recoveredCache.shutdown(); + } + while (firstBlock.refCnt() > 0) { + firstBlock.release(); + } + while (secondBlock.refCnt() > 0) { + secondBlock.release(); + } + } + } + + @Test + public void testHBaseIOExceptionThroughReferenceEvictsStoredEntry(@TempDir File testDir) + throws Exception { + HFileBlock blockToCache = createBlock(200, 1020); + BucketEntry bucketEntry = null; + String cachePath = new File(testDir, "bucket.cache").getAbsolutePath(); + String persistencePath = new File(testDir, "bucket.persistence").getAbsolutePath(); + BucketCache bucketCache = new BucketCache("file:" + cachePath, CAPACITY_SIZE, BLOCK_SIZE, + BLOCK_SIZE_ARRAY, 1, 1000, persistencePath); + try { + assertTrue(bucketCache.waitForCacheInitialization(10000)); + String hfileName = "0123456789abcdef"; + String regionName = "region"; + BlockCacheKey storedKey = + new BlockCacheKey(hfileName, "cf", regionName, 200, true, BlockType.DATA, false); + BlockCacheKey referenceKey = createKey(hfileName + ".parent", 200); + bucketCache.cacheBlock(storedKey, blockToCache); + waitUntilFlushedToCache(bucketCache, storedKey); + bucketCache.fileCacheCompleted(new Path("/table/" + regionName + "/cf/" + hfileName), 1020); + + bucketEntry = bucketCache.backingMap.get(storedKey); + assertNotNull(bucketEntry); + assertTrue(bucketCache.regionCachedSize.containsKey(regionName)); + assertTrue(bucketCache.fullyCachedFiles.containsKey(hfileName)); + + ByteBuffer invalidCachedTime = ByteBuffer.allocate(Long.BYTES); + invalidCachedTime.putLong(bucketEntry.getCachedTime() + 1).flip(); + bucketCache.ioEngine.write(invalidCachedTime, bucketEntry.offset()); + bucketCache.ioEngine.sync(); + + assertNull(bucketCache.getBlock(referenceKey, false, false, false)); + assertFalse(bucketCache.backingMap.containsKey(storedKey)); + assertFalse(bucketCache.blocksByHFile.contains(storedKey)); + assertFalse(bucketCache.regionCachedSize.containsKey(regionName)); + assertFalse(bucketCache.fullyCachedFiles.containsKey(hfileName)); + assertEquals(0, bucketEntry.refCnt()); + assertEquals(0, bucketCache.getAllocator().getUsedSize()); + } finally { + bucketCache.shutdown(); + if (bucketEntry != null && bucketEntry.refCnt() > 0) { + bucketEntry.markAsEvicted(); + } + while (blockToCache.refCnt() > 0) { + blockToCache.release(); + } + } + } + + @Test + public void testPlainIOExceptionKeepsEntryAndCacheMetadata(@TempDir File testDir) + throws Exception { + HFileBlock blockToCache = createBlock(200, 1020); + String cachePath = new File(testDir, "bucket.cache").getAbsolutePath(); + String persistencePath = new File(testDir, "bucket.persistence").getAbsolutePath(); + BucketCache bucketCache = new BucketCache("file:" + cachePath, CAPACITY_SIZE, BLOCK_SIZE, + BLOCK_SIZE_ARRAY, 1, 1000, persistencePath); + FileChannel originalChannel = null; + try { + assertTrue(bucketCache.waitForCacheInitialization(10000)); + String hfileName = "testPlainIOExceptionKeepsEntryAndCacheMetadata"; + String regionName = "region"; + BlockCacheKey key = + new BlockCacheKey(hfileName, "cf", regionName, 200, true, BlockType.DATA, false); + bucketCache.cacheBlock(key, blockToCache); + waitUntilFlushedToCache(bucketCache, key); + bucketCache.fileCacheCompleted(new Path("/table/" + regionName + "/cf/" + hfileName), 1020); + + BucketEntry bucketEntry = bucketCache.backingMap.get(key); + assertNotNull(bucketEntry); + FileIOEngine fileIOEngine = (FileIOEngine) bucketCache.ioEngine; + originalChannel = fileIOEngine.getFileChannels()[0]; + FileChannel failingChannel = Mockito.mock(FileChannel.class); + Mockito.when(failingChannel.read(Mockito.any(ByteBuffer.class), Mockito.anyLong())) + .thenThrow(new IOException("Injected read failure")); + fileIOEngine.getFileChannels()[0] = failingChannel; + + assertNull(bucketCache.getBlock(key, false, false, false)); + assertTrue(bucketCache.isCacheEnabled()); + assertEquals(bucketEntry, bucketCache.backingMap.get(key)); + assertTrue(bucketCache.blocksByHFile.contains(key)); + assertTrue(bucketCache.regionCachedSize.containsKey(regionName)); + assertTrue(bucketCache.fullyCachedFiles.containsKey(hfileName)); + assertEquals(1, bucketEntry.refCnt()); + } finally { + if (originalChannel != null) { + ((FileIOEngine) bucketCache.ioEngine).getFileChannels()[0] = originalChannel; + } + bucketCache.shutdown(); + while (blockToCache.refCnt() > 0) { + blockToCache.release(); + } + } + } + + @Test + public void testInitialPublicationDoesNotRestoreIndexAfterConcurrentEviction() throws Exception { + BucketCache bucketCache = create(1, 1000); + ExecutorService executor = Executors.newSingleThreadExecutor(); + CountDownLatch entryPublished = new CountDownLatch(1); + CountDownLatch continuePublication = new CountDownLatch(1); + BucketEntry publishedEntry = null; + try { + BlockCacheKey key = + createKey("testInitialPublicationDoesNotRestoreIndexAfterConcurrentEviction", 200); + publishedEntry = new BucketEntry(8192, 1020, 1020, 0, false, entry -> ByteBuffAllocator.NONE, + ByteBuffAllocator.HEAP); + BucketEntry entryToPublish = publishedEntry; + bucketCache.backingMap = Mockito.spy(bucketCache.backingMap); + Mockito.doAnswer(invocation -> { + Object previousEntry = invocation.callRealMethod(); + entryPublished.countDown(); + if (!continuePublication.await(10, TimeUnit.SECONDS)) { + throw new AssertionError("Timed out waiting to resume publication"); + } + return previousEntry; + }).when(bucketCache.backingMap).put(key, entryToPublish); + Future publication = + executor.submit(() -> bucketCache.putIntoBackingMap(key, entryToPublish)); + + assertTrue(entryPublished.await(10, TimeUnit.SECONDS)); + assertTrue(entryToPublish.withWriteLock(bucketCache.offsetLock, () -> { + if (bucketCache.backingMap.remove(key, entryToPublish)) { + bucketCache.blockEvicted(key, entryToPublish, false, false); + return true; + } + return false; + })); + continuePublication.countDown(); + publication.get(10, TimeUnit.SECONDS); + + assertFalse(bucketCache.backingMap.containsKey(key)); + assertFalse(bucketCache.blocksByHFile.contains(key)); + assertEquals(0, publishedEntry.refCnt()); + } finally { + continuePublication.countDown(); + executor.shutdownNow(); + try { + assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS)); + } finally { + bucketCache.shutdown(); + if (publishedEntry != null && publishedEntry.refCnt() > 0) { + publishedEntry.markAsEvicted(); + } + } + } + } + /** *
    * This test is for HBASE-26281,
@@ -372,6 +720,7 @@ public void testReplacingBlockAndGettingBlockConcurrently() throws Exception {
 
       cacheBlockThread.join();
       assertTrue(exceptionRef.get() == null);
+      assertTrue(myBucketCache.blocksByHFile.contains(blockCacheKey));
       assertEquals(1, gotHFileBlock.refCnt());
       assertTrue(gotHFileBlock.equals(hfileBlock));
       assertTrue(myBucketCache.overwiteByteBuff == null);