Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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) {
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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
Expand Down
Loading