From 0b0c97abb3fd0e38571cdce849a0228fadeb9183 Mon Sep 17 00:00:00 2001 From: Minwoo Kang Date: Mon, 27 Jul 2026 13:52:15 +0900 Subject: [PATCH 1/3] HBASE-30285 BucketCache shutdown/disable does not release backingMap-owned BucketEntry references --- .../hbase/io/hfile/bucket/BucketCache.java | 223 ++++++++++---- .../io/hfile/bucket/BucketProtoUtils.java | 8 + .../TestBlockEvictionOnRegionMovement.java | 29 +- .../io/hfile/bucket/TestBucketCache.java | 4 + .../hfile/bucket/TestBucketCacheRefCnt.java | 291 ++++++++++++++++++ 5 files changed, 479 insertions(+), 76 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 19cf2af56ee2..36ddb71cc443 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 @@ -92,6 +92,7 @@ import org.apache.hadoop.hbase.util.IdReadWriteLockWithObjectPool; import org.apache.hadoop.hbase.util.IdReadWriteLockWithObjectPool.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; @@ -206,6 +207,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 @@ -416,12 +423,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() { @@ -643,21 +655,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 @@ -681,23 +696,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. @@ -719,11 +739,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); @@ -731,6 +749,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()); @@ -738,6 +759,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} */ @@ -1350,13 +1384,20 @@ 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, () -> { blockEvicted(key, previousEntry, false, false); return null; }); + bucketEntry.withWriteLock(offsetLock, () -> { + if (backingMap.get(key) == bucketEntry) { + blocksByHFile.add(key); + } + return null; + }); + } else { + blocksByHFile.add(key); } } @@ -1551,6 +1592,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()) { @@ -1560,7 +1607,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: {}," @@ -1740,13 +1787,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: {}", @@ -1807,59 +1855,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 cb3f7e5772be..0063d8960f5e 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 @@ -27,6 +27,7 @@ import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.hbase.HBaseTestingUtil; +import org.apache.hadoop.hbase.ServerName; import org.apache.hadoop.hbase.SingleProcessHBaseCluster; import org.apache.hadoop.hbase.StartTestingClusterOption; import org.apache.hadoop.hbase.TableName; @@ -126,20 +127,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 6a5eac5fe3d7..69ed2c90b6dc 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 @@ -385,7 +385,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..4a65a21d6c5b 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,29 @@ */ 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.CyclicBarrier; 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 +55,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 +322,280 @@ 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(); + } + } + } + /** *
    * This test is for HBASE-26281,
@@ -372,6 +662,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);

From c72699b8006861edab37e7158914db1e10370379 Mon Sep 17 00:00:00 2001
From: Minwoo Kang 
Date: Tue, 28 Jul 2026 10:58:03 +0900
Subject: [PATCH 2/3] HBASE-30285 Synchronize cache index publication with
 eviction

---
 .../hbase/io/hfile/bucket/BucketCache.java    | 14 ++---
 .../hfile/bucket/TestBucketCacheRefCnt.java   | 58 +++++++++++++++++++
 2 files changed, 64 insertions(+), 8 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 36ddb71cc443..e0d78a7dacea 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
@@ -1390,15 +1390,13 @@ protected void putIntoBackingMap(BlockCacheKey key, BucketEntry bucketEntry) {
         blockEvicted(key, previousEntry, false, false);
         return null;
       });
-      bucketEntry.withWriteLock(offsetLock, () -> {
-        if (backingMap.get(key) == bucketEntry) {
-          blocksByHFile.add(key);
-        }
-        return null;
-      });
-    } else {
-      blocksByHFile.add(key);
     }
+    bucketEntry.withWriteLock(offsetLock, () -> {
+      if (backingMap.get(key) == bucketEntry) {
+        blocksByHFile.add(key);
+      }
+      return null;
+    });
   }
 
   /**
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 4a65a21d6c5b..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
@@ -35,7 +35,12 @@
 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;
@@ -596,6 +601,59 @@ public void testPlainIOExceptionKeepsEntryAndCacheMetadata(@TempDir File testDir
     }
   }
 
+  @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,

From 35db4868feb5562ba510e1332393b7891fc438e8 Mon Sep 17 00:00:00 2001
From: Minwoo Kang 
Date: Wed, 29 Jul 2026 12:54:40 +0900
Subject: [PATCH 3/3] HBASE-30285 Fix TestBucketCache lock ordering

---
 .../io/hfile/bucket/TestBucketCache.java      | 34 ++++++++++++++-----
 1 file changed, 25 insertions(+), 9 deletions(-)

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 69ed2c90b6dc..2908f3603e41 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!");