From 24a38969f63eff6a07f2a36dfe20b4691c194c16 Mon Sep 17 00:00:00 2001 From: julianchandras Date: Mon, 5 Oct 2026 11:33:56 -0400 Subject: [PATCH] HBASE-30453 BucketCache region cached size is not decremented on eviction for non-persistent IOEngines --- .../hbase/io/hfile/bucket/BucketCache.java | 14 ++-- .../io/hfile/bucket/TestBucketCache.java | 25 ++++++ .../bucket/TestVerifyBucketCacheFile.java | 78 +++++++++++++++++++ 3 files changed, 110 insertions(+), 7 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 e0d78a7dacea..3f7621f1d968 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 @@ -779,10 +779,13 @@ void blockEvicted(BlockCacheKey cacheKey, BucketEntry bucketEntry, boolean decre boolean evictedByEvictionProcess) { bucketEntry.markAsEvicted(); blocksByHFile.remove(cacheKey); + // putIntoBackingMap accounts every entry added to the backingMap, regardless of the IOEngine, + // so every entry removed from it must be decremented as well. + updateRegionCachedSize(cacheKey, (bucketEntry.getLength() * -1)); if (decrementBlockNumber) { this.blockNumber.decrement(); if (ioEngine.isPersistent()) { - fileNotFullyCached(cacheKey, bucketEntry); + fileNotFullyCached(cacheKey.getHfileName()); } } if (evictedByEvictionProcess) { @@ -793,11 +796,8 @@ void blockEvicted(BlockCacheKey cacheKey, BucketEntry bucketEntry, boolean decre } } - private void fileNotFullyCached(BlockCacheKey key, BucketEntry entry) { - // Update the updateRegionCachedSize before removing the file from fullyCachedFiles. - // This computation should happen even if the file is not in fullyCachedFiles map. - updateRegionCachedSize(key, (entry.getLength() * -1)); - fullyCachedFiles.remove(key.getHfileName()); + private void fileNotFullyCached(String hfileName) { + fullyCachedFiles.remove(hfileName); } public void fileCacheCompleted(Path filePath, long size) { @@ -1743,7 +1743,7 @@ private void verifyFileIntegrity(BucketCacheProtos.BucketCacheEntry proto) { } catch (IOException e1) { LOG.debug("Check for key {} failed. Evicting.", keyEntry.getKey()); evictBlock(keyEntry.getKey()); - fileNotFullyCached(keyEntry.getKey(), keyEntry.getValue()); + fileNotFullyCached(keyEntry.getKey().getHfileName()); } } backingMapValidated.set(true); 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 2908f3603e41..2901aeec3e9a 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 @@ -1147,6 +1147,31 @@ public void testBlockPriority() throws Exception { assertEquals(cache.backingMap.get(block.getBlockName()).getPriority(), BlockPriority.MULTI); } + @TestTemplate + public void testRegionCachedSizeDecrementedOnEvictionForNonPersistentEngine() throws Exception { + assertFalse(cache.ioEngine.isPersistent()); + String hfileName = "testRegionCachedSizeDecrementedOnEviction"; + String regionName = "region"; + HFileBlockPair[] blocks = CacheTestUtils.generateHFileBlocks(BLOCK_SIZE, 2); + BlockCacheKey key1 = + new BlockCacheKey(hfileName, "cf", regionName, 0, true, BlockType.DATA, false); + BlockCacheKey key2 = + new BlockCacheKey(hfileName, "cf", regionName, BLOCK_SIZE, true, BlockType.DATA, false); + cacheAndWaitUntilFlushedToBucket(cache, key1, blocks[0].getBlock(), true); + cacheAndWaitUntilFlushedToBucket(cache, key2, blocks[1].getBlock(), true); + long length1 = cache.backingMap.get(key1).getLength(); + long length2 = cache.backingMap.get(key2).getLength(); + assertEquals(length1 + length2, (long) cache.getRegionCachedInfo().get().get(regionName)); + + // Single block eviction + assertTrue(cache.evictBlock(key1)); + assertEquals(length2, (long) cache.getRegionCachedInfo().get().get(regionName)); + + // Eviction by file, as done when a store file reader is closed + assertEquals(1, cache.evictBlocksByHfileName(hfileName)); + assertFalse(cache.getRegionCachedInfo().get().containsKey(regionName)); + } + @TestTemplate public void testIOTimePerHitReturnsZeroWhenNoHits() throws NoSuchFieldException, IllegalAccessException { diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestVerifyBucketCacheFile.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestVerifyBucketCacheFile.java index 8fcece5b9264..be18faa65367 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestVerifyBucketCacheFile.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/bucket/TestVerifyBucketCacheFile.java @@ -29,6 +29,7 @@ import java.io.File; import java.io.FileOutputStream; import java.io.OutputStreamWriter; +import java.io.RandomAccessFile; import java.nio.file.FileSystems; import java.nio.file.Files; import java.nio.file.attribute.FileTime; @@ -41,6 +42,7 @@ import org.apache.hadoop.hbase.HBaseTestingUtil; import org.apache.hadoop.hbase.Waiter; import org.apache.hadoop.hbase.io.hfile.BlockCacheKey; +import org.apache.hadoop.hbase.io.hfile.BlockType; import org.apache.hadoop.hbase.io.hfile.CacheConfig; import org.apache.hadoop.hbase.io.hfile.CacheTestUtils; import org.apache.hadoop.hbase.io.hfile.Cacheable; @@ -339,6 +341,82 @@ public void testModifiedBucketCacheFileTime() throws Exception { TEST_UTIL.cleanupTestDir(); } + /** + * Test whether the region cached size is still correct after the cache validation thread evicts + * an invalid block. First start BucketCache and add three blocks for the same region, then + * shutdown BucketCache and persist cache to file. Then overwrite the cached time recorded in the + * first block and modify the cache file's last modified time, so that the checksum verification + * fails when restarting BucketCache. The validation thread then goes through all the cached + * blocks and evicts the first one, which should be discounted from the region cached size only + * once. + * @throws Exception the exception + */ + @TestTemplate + public void testRegionCachedSizeAfterInvalidBlockEvictedOnValidation() throws Exception { + HBaseTestingUtil TEST_UTIL = new HBaseTestingUtil(); + Path testDir = TEST_UTIL.getDataTestDir(); + TEST_UTIL.getTestFileSystem().mkdirs(testDir); + Configuration conf = HBaseConfiguration.create(); + // Disables the persister thread by setting its interval to MAX_VALUE + conf.setLong(BUCKETCACHE_PERSIST_INTERVAL_KEY, Long.MAX_VALUE); + String regionName = "region"; + BucketCache bucketCache = null; + try { + bucketCache = new BucketCache("file:" + testDir + "/bucket.cache", capacitySize, + constructedBlockSize, constructedBlockSizes, writeThreads, writerQLen, + testDir + "/bucket.persistence", DEFAULT_ERROR_TOLERATION_DURATION, conf); + assertTrue(bucketCache.waitForCacheInitialization(10000)); + + CacheTestUtils.HFileBlockPair[] blocks = + CacheTestUtils.generateHFileBlocks(constructedBlockSize, 3); + // Add three blocks, all belonging to the same region + BlockCacheKey[] keys = new BlockCacheKey[blocks.length]; + for (int i = 0; i < blocks.length; i++) { + keys[i] = new BlockCacheKey("hfile", "cf", regionName, (long) i * constructedBlockSize, + true, BlockType.DATA, false); + cacheAndWaitUntilFlushedToBucket(bucketCache, keys[i], blocks[i].getBlock()); + } + long firstBlockOffset = bucketCache.backingMap.get(keys[0]).offset(); + // persist cache to file + bucketCache.shutdown(); + + // overwrite the cached time stored at the start of the first block + try (RandomAccessFile raf = new RandomAccessFile(testDir + "/bucket.cache", "rw")) { + raf.seek(firstBlockOffset); + raf.writeLong(-1L); + } + // modified bucket cache file LastModifiedTime, so the checksum verification fails + final java.nio.file.Path file = + FileSystems.getDefault().getPath(testDir.toString(), "bucket.cache"); + Files.setLastModifiedTime(file, FileTime.from(Instant.now().plusMillis(1_000))); + + bucketCache = new BucketCache("file:" + testDir + "/bucket.cache", capacitySize, + constructedBlockSize, constructedBlockSizes, writeThreads, writerQLen, + testDir + "/bucket.persistence", DEFAULT_ERROR_TOLERATION_DURATION, conf); + assertTrue(bucketCache.waitForCacheInitialization(10000)); + waitPersistentCacheValidation(conf, bucketCache); + + // The region cached size should match the blocks still in the cache. We sum what is left + // instead of expecting two blocks, as the validation thread may not evict the first block + // if it checks it before the cache is enabled. + long remainingSize = 0; + for (BlockCacheKey key : keys) { + BucketEntry entry = bucketCache.backingMap.get(key); + if (entry != null) { + remainingSize += entry.getLength(); + } + } + assertNotEquals(0, remainingSize); + assertEquals(remainingSize, + (long) bucketCache.getRegionCachedInfo().get().getOrDefault(regionName, 0L)); + } finally { + if (bucketCache != null) { + bucketCache.shutdown(); + } + } + TEST_UTIL.cleanupTestDir(); + } + /** * When using persistent bucket cache, there may be crashes between persisting the backing map and * syncing new blocks to the cache file itself, leading to an inconsistent state between the cache