diff --git a/dev-support/spotbugs-exclude.xml b/dev-support/spotbugs-exclude.xml index 17b8d2cbdedd..b90b6bf99297 100644 --- a/dev-support/spotbugs-exclude.xml +++ b/dev-support/spotbugs-exclude.xml @@ -247,6 +247,11 @@ + + + + + diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheFactory.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheFactory.java index debfc09442a1..c3dd8947ce87 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheFactory.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheFactory.java @@ -27,6 +27,9 @@ import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hbase.HConstants; import org.apache.hadoop.hbase.io.hfile.bucket.BucketCache; +import org.apache.hadoop.hbase.io.hfile.cache.CacheEngine; +import org.apache.hadoop.hbase.io.hfile.cache.CacheEngines; +import org.apache.hadoop.hbase.io.hfile.cache.LruCacheEngine; import org.apache.hadoop.hbase.io.util.MemorySizeUtil; import org.apache.hadoop.hbase.regionserver.HRegion; import org.apache.hadoop.hbase.util.ReflectionUtils; @@ -71,8 +74,8 @@ public final class BlockCacheFactory { */ public static final String BLOCKCACHE_BLOCKSIZE_KEY = "hbase.blockcache.minblocksize"; - private static final String EXTERNAL_BLOCKCACHE_KEY = "hbase.blockcache.use.external"; - private static final boolean EXTERNAL_BLOCKCACHE_DEFAULT = false; + public static final String EXTERNAL_BLOCKCACHE_KEY = "hbase.blockcache.use.external"; + public static final boolean EXTERNAL_BLOCKCACHE_DEFAULT = false; private static final String EXTERNAL_BLOCKCACHE_CLASS_KEY = "hbase.blockcache.external.class"; @@ -135,6 +138,62 @@ public static BlockCache createBlockCache(Configuration conf) { return createBlockCache(conf, null); } + /** + * Creates the configured first-level cache as a cache engine. + *

+ * Cache implementations that have been migrated to {@link CacheEngine} are instantiated directly. + * Legacy {@link FirstLevelBlockCache} implementations are adapted until their migration is + * complete. + * @param c cache configuration + * @return first-level cache engine, or {@code null} when the on-heap cache is disabled + */ + public static CacheEngine createFirstLevelCacheEngine(final Configuration c) { + final long cacheSize = MemorySizeUtil.getOnHeapCacheSize(c); + if (cacheSize < 0) { + return null; + } + + String policy = c.get(BLOCKCACHE_POLICY_KEY, BLOCKCACHE_POLICY_DEFAULT); + int blockSize = c.getInt(BLOCKCACHE_BLOCKSIZE_KEY, HConstants.DEFAULT_BLOCKSIZE); + LOG.info("Allocating CacheEngine size=" + StringUtils.byteDesc(cacheSize) + ", blockSize=" + + StringUtils.byteDesc(blockSize)); + + if (policy.equalsIgnoreCase("LRU")) { + return new LruCacheEngine(cacheSize, blockSize, true, c); + } else if (policy.equalsIgnoreCase("IndexOnlyLRU")) { + return CacheEngines.fromBlockCache(new IndexOnlyLruBlockCache(cacheSize, blockSize, true, c)); + } else if (policy.equalsIgnoreCase("TinyLFU")) { + return CacheEngines + .fromBlockCache(new TinyLfuBlockCache(cacheSize, blockSize, ForkJoinPool.commonPool(), c)); + } else if (policy.equalsIgnoreCase("AdaptiveLRU")) { + return CacheEngines.fromBlockCache(new LruAdaptiveBlockCache(cacheSize, blockSize, true, c)); + } else { + throw new IllegalArgumentException("Unknown policy: " + policy); + } + } + + /** + * Creates the configured external second-level cache as a cache engine. + * @param c cache configuration + * @return external cache engine, or {@code null} when no external cache can be created + */ + public static CacheEngine createExternalCacheEngine(Configuration c) { + BlockCache blockCache = createExternalBlockcache(c); + return blockCache == null ? null : CacheEngines.fromBlockCache(blockCache); + } + + /** + * Creates the configured bucket cache as a cache engine. + * @param c cache configuration + * @param onlineRegions currently online regions + * @return bucket cache engine, or {@code null} when BucketCache is disabled + */ + public static CacheEngine createBucketCacheEngine(Configuration c, + Map onlineRegions) { + BucketCache bucketCache = createBucketCache(c, onlineRegions); + return bucketCache == null ? null : CacheEngines.fromBlockCache(bucketCache); + } + private static FirstLevelBlockCache createFirstLevelCache(final Configuration c) { final long cacheSize = MemorySizeUtil.getOnHeapCacheSize(c); if (cacheSize < 0) { diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheUtil.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheUtil.java index f48da92f515d..d9f4c39afdb8 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheUtil.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/BlockCacheUtil.java @@ -30,6 +30,7 @@ import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.ConcurrentSkipListSet; import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.io.hfile.cache.CacheEngine; import org.apache.hadoop.hbase.metrics.impl.FastLongHistogram; import org.apache.hadoop.hbase.nio.ByteBuff; import org.apache.hadoop.hbase.regionserver.HRegion; @@ -248,6 +249,51 @@ public static boolean shouldReplaceExistingCacheBlock(BlockCache blockCache, } } + /** + * Because of region splitting, it is possible that the split key is located in the middle of a + * block. As a result, both daughter regions may load the same block from their parent HFile. + *

+ * When using positional reads, HBase does not force the read to include the complete next-block + * header. Therefore, when two threads try to cache the same block, one thread may have read the + * complete next-block header while the other did not. If the already cached block does not + * contain {@code nextBlockOnDiskSize} but the new block does, replacing the existing block with + * the new one improves subsequent read performance. See HBASE-20447. + *

+ * @param cacheEngine cache engine to check + * @param cacheKey block cache key + * @param newBlock new block being considered for insertion + * @return {@code true} if the existing cached block should be replaced by {@code newBlock}; + * {@code false} if the existing cached block should be retained + */ + public static boolean shouldReplaceExistingCacheBlock(CacheEngine cacheEngine, + BlockCacheKey cacheKey, Cacheable newBlock) { + // NOTICE: getBlock retains the existingBlock before returning it. + Cacheable existingBlock = cacheEngine.getBlock(cacheKey, false, false, false); + if (existingBlock == null) { + return true; + } + + try { + int comparison = BlockCacheUtil.validateBlockAddition(existingBlock, newBlock, cacheKey); + if (comparison < 0) { + LOG.warn("Cached block contents differ by nextBlockOnDiskSize, the new block has " + + "nextBlockOnDiskSize set. Caching new block."); + return true; + } else if (comparison > 0) { + LOG.warn("Cached block contents differ by nextBlockOnDiskSize, the existing block has " + + "nextBlockOnDiskSize set. Keeping cached block."); + return false; + } else { + LOG.debug("Caching an already cached block: {}. This is harmless and can happen in rare " + + "cases (see HBASE-8547)", cacheKey); + return false; + } + } finally { + // Release the reference retained by CacheEngine#getBlock. + existingBlock.release(); + } + } + public static Set listAllFilesNames(Map onlineRegions) { Set files = new HashSet<>(); onlineRegions.values().forEach(r -> { diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/CacheConfig.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/CacheConfig.java index 9ab39b00c90d..0950f58d29b8 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/CacheConfig.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/CacheConfig.java @@ -320,6 +320,7 @@ private void initFromConf(Configuration conf, ColumnFamilyDescriptor family) { /** * Constructs a cache configuration copied from the specified configuration. + * @param cacheConf cache configuration to copy */ public CacheConfig(CacheConfig cacheConf) { this.cacheDataOnRead = cacheConf.cacheDataOnRead; @@ -337,9 +338,7 @@ public CacheConfig(CacheConfig cacheConf) { this.blockCache = cacheConf.blockCache; this.byteBuffAllocator = cacheConf.byteBuffAllocator; this.heapUsageThreshold = cacheConf.heapUsageThreshold; - this.cacheAccessService = blockCache != null - ? CacheAccessServices.fromBlockCache(blockCache) - : CacheAccessServices.disabled(); + this.cacheAccessService = cacheConf.cacheAccessService; } private CacheConfig() { @@ -551,8 +550,12 @@ public boolean shouldLockOnCacheMiss(BlockType blockType) { } /** - * Returns the block cache. - * @return the block cache, or null if caching is completely disabled + * Returns the legacy block cache when this configuration exposes one. + *

+ * Native cache-engine configurations do not necessarily expose a {@link BlockCache}. Callers + * should use {@link #getCacheAccessService()} for cache access and diagnostics. + *

+ * @return legacy block cache when available */ public Optional getBlockCache() { return Optional.ofNullable(this.blockCache); @@ -561,10 +564,9 @@ public Optional getBlockCache() { /** * Returns the cache access service used by HFile read/write path callers. *

- * This service is the migration-facing cache abstraction. For now it is backed by the existing - * {@link BlockCache} when block cache is configured, or by a disabled no-op implementation when - * block cache is unavailable. This keeps cache construction unchanged while allowing callers such - * as {@code HFileReaderImpl} to depend on {@link CacheAccessService}. + * The cache access service is the primary cache abstraction exposed by this configuration. It may + * be backed by native {@link CacheEngine} implementations, legacy {@link BlockCache} + * implementations adapted as cache engines, or a combination of both. *

* @return cache access service */ @@ -583,10 +585,10 @@ private static boolean isCombinedBlockCacheCompatible(CacheAccessService cacheAc if (!(cacheAccessService instanceof TopologyBackedCacheAccessService)) { return false; } - TopologyBackedCacheAccessService service = (TopologyBackedCacheAccessService) cacheAccessService; - return service.getTopology().getType() == CacheTopologyType.TIERED_EXCLUSIVE; + CacheTopology topology = service.getTopology(); + return topology.getTiers().contains(CacheTier.L1) && topology.getTiers().contains(CacheTier.L2); } public ByteBuffAllocator getByteBuffAllocator() { diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/HFileBlock.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/HFileBlock.java index 087c619a7e99..47632910e838 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/HFileBlock.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/HFileBlock.java @@ -2084,7 +2084,7 @@ public boolean equals(Object comparison) { return true; } - DataBlockEncoding getDataBlockEncoding() { + public DataBlockEncoding getDataBlockEncoding() { if (blockType == BlockType.ENCODED_DATA) { return DataBlockEncoding.getEncodingById(getDataBlockEncodingId()); } @@ -2211,7 +2211,7 @@ private static HFileBlock shallowClone(HFileBlock blk, ByteBuff newBuf) { return createBuilder(blk, newBuf).build(); } - static HFileBlock deepCloneOnHeap(HFileBlock blk) { + public static HFileBlock deepCloneOnHeap(HFileBlock blk) { ByteBuff deepCloned = ByteBuff .wrap(ByteBuffer.wrap(blk.bufWithoutChecksum.toBytes(0, blk.bufWithoutChecksum.limit()))); return createBuilder(blk, deepCloned).build(); diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/AggregateCacheStats.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/AggregateCacheStats.java new file mode 100644 index 000000000000..5dbdb022a494 --- /dev/null +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/AggregateCacheStats.java @@ -0,0 +1,386 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache; + +import java.util.Objects; +import java.util.function.ToLongFunction; +import org.apache.hadoop.hbase.io.hfile.CacheStats; +import org.apache.yetus.audience.InterfaceAudience; + +/** + * Aggregate cache statistics backed by multiple cache-engine statistics instances. + */ +@InterfaceAudience.Private +final class AggregateCacheStats extends CacheStats { + + private final CacheStats[] delegates; + + /** + * Creates aggregate cache statistics. + * @param name aggregate statistics name + * @param delegates cache statistics to aggregate + */ + AggregateCacheStats(String name, CacheStats... delegates) { + super(name); + Objects.requireNonNull(delegates, "delegates must not be null"); + + int count = 0; + for (CacheStats delegate : delegates) { + if (delegate != null) { + count++; + } + } + + this.delegates = new CacheStats[count]; + int index = 0; + for (CacheStats delegate : delegates) { + if (delegate != null) { + this.delegates[index++] = delegate; + } + } + } + + /** + * Sums a value across all underlying cache statistics. + * @param extractor value extractor + * @return aggregate value + */ + private long sum(ToLongFunction extractor) { + long result = 0L; + for (CacheStats delegate : delegates) { + result += extractor.applyAsLong(delegate); + } + return result; + } + + /** + * Returns the aggregate data-block miss count. + * @return aggregate data-block miss count + */ + @Override + public long getDataMissCount() { + return sum(CacheStats::getDataMissCount); + } + + /** + * Returns the aggregate leaf-index miss count. + * @return aggregate leaf-index miss count + */ + @Override + public long getLeafIndexMissCount() { + return sum(CacheStats::getLeafIndexMissCount); + } + + /** + * Returns the aggregate bloom-chunk miss count. + * @return aggregate bloom-chunk miss count + */ + @Override + public long getBloomChunkMissCount() { + return sum(CacheStats::getBloomChunkMissCount); + } + + /** + * Returns the aggregate metadata miss count. + * @return aggregate metadata miss count + */ + @Override + public long getMetaMissCount() { + return sum(CacheStats::getMetaMissCount); + } + + /** + * Returns the aggregate root-index miss count. + * @return aggregate root-index miss count + */ + @Override + public long getRootIndexMissCount() { + return sum(CacheStats::getRootIndexMissCount); + } + + /** + * Returns the aggregate intermediate-index miss count. + * @return aggregate intermediate-index miss count + */ + @Override + public long getIntermediateIndexMissCount() { + return sum(CacheStats::getIntermediateIndexMissCount); + } + + /** + * Returns the aggregate file-info miss count. + * @return aggregate file-info miss count + */ + @Override + public long getFileInfoMissCount() { + return sum(CacheStats::getFileInfoMissCount); + } + + /** + * Returns the aggregate general-bloom metadata miss count. + * @return aggregate general-bloom metadata miss count + */ + @Override + public long getGeneralBloomMetaMissCount() { + return sum(CacheStats::getGeneralBloomMetaMissCount); + } + + /** + * Returns the aggregate delete-family bloom miss count. + * @return aggregate delete-family bloom miss count + */ + @Override + public long getDeleteFamilyBloomMissCount() { + return sum(CacheStats::getDeleteFamilyBloomMissCount); + } + + /** + * Returns the aggregate trailer miss count. + * @return aggregate trailer miss count + */ + @Override + public long getTrailerMissCount() { + return sum(CacheStats::getTrailerMissCount); + } + + /** + * Returns the aggregate data-block hit count. + * @return aggregate data-block hit count + */ + @Override + public long getDataHitCount() { + return sum(CacheStats::getDataHitCount); + } + + /** + * Returns the aggregate leaf-index hit count. + * @return aggregate leaf-index hit count + */ + @Override + public long getLeafIndexHitCount() { + return sum(CacheStats::getLeafIndexHitCount); + } + + /** + * Returns the aggregate bloom-chunk hit count. + * @return aggregate bloom-chunk hit count + */ + @Override + public long getBloomChunkHitCount() { + return sum(CacheStats::getBloomChunkHitCount); + } + + /** + * Returns the aggregate metadata hit count. + * @return aggregate metadata hit count + */ + @Override + public long getMetaHitCount() { + return sum(CacheStats::getMetaHitCount); + } + + /** + * Returns the aggregate root-index hit count. + * @return aggregate root-index hit count + */ + @Override + public long getRootIndexHitCount() { + return sum(CacheStats::getRootIndexHitCount); + } + + /** + * Returns the aggregate intermediate-index hit count. + * @return aggregate intermediate-index hit count + */ + @Override + public long getIntermediateIndexHitCount() { + return sum(CacheStats::getIntermediateIndexHitCount); + } + + /** + * Returns the aggregate file-info hit count. + * @return aggregate file-info hit count + */ + @Override + public long getFileInfoHitCount() { + return sum(CacheStats::getFileInfoHitCount); + } + + /** + * Returns the aggregate general-bloom metadata hit count. + * @return aggregate general-bloom metadata hit count + */ + @Override + public long getGeneralBloomMetaHitCount() { + return sum(CacheStats::getGeneralBloomMetaHitCount); + } + + /** + * Returns the aggregate delete-family bloom hit count. + * @return aggregate delete-family bloom hit count + */ + @Override + public long getDeleteFamilyBloomHitCount() { + return sum(CacheStats::getDeleteFamilyBloomHitCount); + } + + /** + * Returns the aggregate trailer hit count. + * @return aggregate trailer hit count + */ + @Override + public long getTrailerHitCount() { + return sum(CacheStats::getTrailerHitCount); + } + + /** + * Returns the aggregate miss count. + * @return aggregate miss count + */ + @Override + public long getMissCount() { + return sum(CacheStats::getMissCount); + } + + /** + * Returns the aggregate primary-replica miss count. + * @return aggregate primary-replica miss count + */ + @Override + public long getPrimaryMissCount() { + return sum(CacheStats::getPrimaryMissCount); + } + + /** + * Returns the aggregate caching-request miss count. + * @return aggregate caching-request miss count + */ + @Override + public long getMissCachingCount() { + return sum(CacheStats::getMissCachingCount); + } + + /** + * Returns the aggregate hit count. + * @return aggregate hit count + */ + @Override + public long getHitCount() { + return sum(CacheStats::getHitCount); + } + + /** + * Returns the aggregate primary-replica hit count. + * @return aggregate primary-replica hit count + */ + @Override + public long getPrimaryHitCount() { + return sum(CacheStats::getPrimaryHitCount); + } + + /** + * Returns the aggregate caching-request hit count. + * @return aggregate caching-request hit count + */ + @Override + public long getHitCachingCount() { + return sum(CacheStats::getHitCachingCount); + } + + /** + * Returns the aggregate eviction count. + * @return aggregate eviction count + */ + @Override + public long getEvictionCount() { + return sum(CacheStats::getEvictionCount); + } + + /** + * Returns the aggregate number of evicted blocks. + * @return aggregate evicted-block count + */ + @Override + public long getEvictedCount() { + return sum(CacheStats::getEvictedCount); + } + + /** + * Returns the aggregate number of evicted primary-replica blocks. + * @return aggregate primary evicted-block count + */ + @Override + public long getPrimaryEvictedCount() { + return sum(CacheStats::getPrimaryEvictedCount); + } + + /** + * Returns the aggregate failed-insert count. + * @return aggregate failed-insert count + */ + @Override + public long getFailedInserts() { + return sum(CacheStats::getFailedInserts); + } + + /** + * Rolls the metrics period for all underlying cache statistics. + */ + @Override + public void rollMetricsPeriod() { + for (CacheStats delegate : delegates) { + delegate.rollMetricsPeriod(); + } + } + + /** + * Returns the aggregate hit count over the configured rolling window. + * @return aggregate rolling hit count + */ + @Override + public long getSumHitCountsPastNPeriods() { + return sum(CacheStats::getSumHitCountsPastNPeriods); + } + + /** + * Returns the aggregate request count over the configured rolling window. + * @return aggregate rolling request count + */ + @Override + public long getSumRequestCountsPastNPeriods() { + return sum(CacheStats::getSumRequestCountsPastNPeriods); + } + + /** + * Returns the aggregate caching-hit count over the configured rolling window. + * @return aggregate rolling caching-hit count + */ + @Override + public long getSumHitCachingCountsPastNPeriods() { + return sum(CacheStats::getSumHitCachingCountsPastNPeriods); + } + + /** + * Returns the aggregate caching-request count over the configured rolling window. + * @return aggregate rolling caching-request count + */ + @Override + public long getSumRequestCachingCountsPastNPeriods() { + return sum(CacheStats::getSumRequestCachingCountsPastNPeriods); + } +} diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServices.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServices.java index 344cfc19b501..050ab87f3dde 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServices.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServices.java @@ -17,6 +17,7 @@ */ package org.apache.hadoop.hbase.io.hfile.cache; +import java.util.Map; import java.util.Objects; import java.util.Optional; import org.apache.hadoop.conf.Configuration; @@ -25,20 +26,15 @@ import org.apache.hadoop.hbase.io.hfile.CachedBlock; import org.apache.hadoop.hbase.io.hfile.CombinedBlockCache; import org.apache.hadoop.hbase.io.hfile.InclusiveCombinedBlockCache; +import org.apache.hadoop.hbase.regionserver.HRegion; import org.apache.yetus.audience.InterfaceAudience; /** * Utility methods for creating {@link CacheAccessService} instances. *

- * This class keeps service construction centralized without introducing a full factory or plugin - * loader in the initial {@code CacheAccessService} ticket. The first supported construction modes - * are a legacy {@link BlockCache}-backed service, a topology-backed service, and a disabled no-op - * service. - *

- *

- * A later integration step can move construction into {@code BlockCacheFactory} once HBase runtime - * wiring starts returning {@link CacheAccessService} instead of, or alongside, raw - * {@link BlockCache}. + * This class is the construction entry point for cache access services. Native {@link CacheEngine} + * implementations are used directly, while legacy {@link BlockCache} implementations are adapted to + * cache engines during migration. *

*/ @InterfaceAudience.Private @@ -81,33 +77,59 @@ public static CacheAccessService fromBlockCache(BlockCache blockCache) { DefaultHBaseCachePlacementAdmissionPolicy.INSTANCE); } + /** + * Creates a {@link CacheAccessService} from the block cache configuration. + * @param conf cache configuration + * @return configured cache access service, or a disabled service when block caching is disabled + * @throws NullPointerException if {@code conf} is {@code null} + */ + public static CacheAccessService fromConfiguration(Configuration conf) { + return fromConfiguration(conf, null); + } + /** * Creates a {@link CacheAccessService} from the block cache configuration. *

- * This method is a compatibility factory for tests and transitional code paths that want to - * obtain a {@link CacheAccessService} directly from {@link Configuration}, while still using the - * existing {@link BlockCacheFactory} and legacy {@link BlockCache} implementations underneath. - *

- *

- * The method delegates block-cache construction to - * {@link BlockCacheFactory#createBlockCache(Configuration)}. If the legacy factory creates a - * {@link BlockCache}, the returned service is backed by that cache through - * {@link TopologyBackedCacheAccessService}. If the legacy factory does not create a cache, this - * method returns the disabled/no-op cache access service. - *

- *

- * This method does not introduce new cache-engine or topology-based runtime wiring. It is - * intended only as a bridge while existing HBase tests and integration paths migrate from direct - * {@link BlockCache} usage to {@link CacheAccessService}. + * Cache implementations that implement {@link CacheEngine} natively are used directly. Legacy + * {@link BlockCache} implementations are adapted to {@link CacheEngine} until their migration is + * complete. *

- * @param conf configuration used by {@link BlockCacheFactory} - * @return cache access service created from the configured legacy block cache, or disabled when - * no block cache is configured + * @param conf cache configuration + * @param onlineRegions currently online regions, or {@code null} when unavailable + * @return configured cache access service, or a disabled service when block caching is disabled * @throws NullPointerException if {@code conf} is {@code null} */ - public static CacheAccessService fromConfiguration(Configuration conf) { + public static CacheAccessService fromConfiguration(Configuration conf, + Map onlineRegions) { Objects.requireNonNull(conf, "conf must not be null"); - return fromBlockCache(BlockCacheFactory.createBlockCache(conf)); + + CacheEngine l1 = BlockCacheFactory.createFirstLevelCacheEngine(conf); + if (l1 == null) { + return disabled(); + } + + CachePlacementAdmissionPolicy policy = DefaultHBaseCachePlacementAdmissionPolicy.INSTANCE; + + boolean useExternal = conf.getBoolean(BlockCacheFactory.EXTERNAL_BLOCKCACHE_KEY, + BlockCacheFactory.EXTERNAL_BLOCKCACHE_DEFAULT); + + if (useExternal) { + CacheEngine l2 = BlockCacheFactory.createExternalCacheEngine(conf); + if (l2 == null) { + return TopologyBackedCacheAccessServices.fromSingleCacheEngine("single", l1, policy); + } + + return TopologyBackedCacheAccessServices.fromTieredInclusiveCacheEngines("inclusive", l1, l2, + policy); + } + + CacheEngine l2 = BlockCacheFactory.createBucketCacheEngine(conf, onlineRegions); + if (l2 == null) { + return TopologyBackedCacheAccessServices.fromSingleCacheEngine("single", l1, policy); + } + + return TopologyBackedCacheAccessServices.fromTieredExclusiveCacheEngines("combined", l1, l2, + policy); } /** diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheEngine.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheEngine.java index 04020725b570..55dc85b37163 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheEngine.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheEngine.java @@ -344,4 +344,17 @@ default Optional>> getFullyCachedFiles() { default Optional> asCachedBlockIterable() { return Optional.empty(); } + + /** + * Sets the listener that receives capacity-driven block eviction events. + *

+ * Cache engines that can expose pressure evictions should invoke the listener before the evicted + * block becomes unavailable. Explicit invalidation operations must not be reported through this + * listener. + *

+ * @param listener eviction listener, or {@code null} to remove the current listener + */ + default void setEvictionListener(CacheEvictionListener listener) { + // noop + } } diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheEvictionListener.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheEvictionListener.java new file mode 100644 index 000000000000..248082877b20 --- /dev/null +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheEvictionListener.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache; + +import org.apache.hadoop.hbase.io.hfile.BlockCacheKey; +import org.apache.hadoop.hbase.io.hfile.Cacheable; +import org.apache.yetus.audience.InterfaceAudience; + +/** + * Listener for cache-engine eviction events. + *

+ * The listener is intended for topology-level handling of capacity-driven evictions. In particular, + * an exclusive tiered topology may use an eviction from L1 as a request to demote the block to L2. + *

+ *

+ * The supplied block remains valid for the duration of the callback. An implementation that needs + * to retain the block after the callback returns must establish its own reference. + *

+ */ +@InterfaceAudience.Private +public interface CacheEvictionListener { + + /** + * Handles a block evicted by a cache engine because of cache pressure. + * @param sourceEngine engine that evicted the block + * @param cacheKey key identifying the evicted block + * @param block evicted block + */ + void onEviction(CacheEngine sourceEngine, BlockCacheKey cacheKey, Cacheable block); +} diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheTopology.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheTopology.java index 16e139148739..0aa69135e92d 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheTopology.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/CacheTopology.java @@ -199,4 +199,20 @@ default boolean demote(BlockCacheKey cacheKey, Cacheable block, CacheEngine sour */ CacheTopologyView getView(); + /** + * Handles a block evicted by one of this topology's cache engines because of cache pressure. + *

+ * The default implementation discards the event. Topologies that support eviction-driven demotion + * may override this method. + *

+ * @param cacheKey key identifying the evicted block + * @param block evicted block + * @param sourceEngine engine that evicted the block + * @return {@code true} if the topology placed the block in another engine + */ + default boolean handleEviction(BlockCacheKey cacheKey, Cacheable block, + CacheEngine sourceEngine) { + return false; + } + } diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/LruCacheEngine.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/LruCacheEngine.java new file mode 100644 index 000000000000..2717a5553140 --- /dev/null +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/LruCacheEngine.java @@ -0,0 +1,1623 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache; + +import java.lang.ref.WeakReference; +import java.util.EnumMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.PriorityQueue; +import java.util.SortedSet; +import java.util.TreeSet; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.atomic.LongAdder; +import java.util.concurrent.locks.ReentrantLock; +import org.apache.commons.lang3.mutable.MutableBoolean; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.HBaseInterfaceAudience; +import org.apache.hadoop.hbase.io.HeapSize; +import org.apache.hadoop.hbase.io.encoding.DataBlockEncoding; +import org.apache.hadoop.hbase.io.hfile.BlockCacheKey; +import org.apache.hadoop.hbase.io.hfile.BlockCacheUtil; +import org.apache.hadoop.hbase.io.hfile.BlockPriority; +import org.apache.hadoop.hbase.io.hfile.BlockType; +import org.apache.hadoop.hbase.io.hfile.CacheStats; +import org.apache.hadoop.hbase.io.hfile.Cacheable; +import org.apache.hadoop.hbase.io.hfile.CachedBlock; +import org.apache.hadoop.hbase.io.hfile.HFileBlock; +import org.apache.hadoop.hbase.io.hfile.LruCachedBlock; +import org.apache.hadoop.hbase.io.hfile.LruCachedBlockQueue; +import org.apache.hadoop.hbase.util.ClassSize; +import org.apache.hadoop.util.StringUtils; +import org.apache.yetus.audience.InterfaceAudience; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.hbase.thirdparty.com.google.common.base.MoreObjects; +import org.apache.hbase.thirdparty.com.google.common.base.Objects; +import org.apache.hbase.thirdparty.com.google.common.util.concurrent.ThreadFactoryBuilder; + +/** + * Native LRU {@link CacheEngine} implementation. + *

+ * This cache is memory-aware using {@link HeapSize}, memory-bound using an LRU eviction algorithm, + * and concurrent. It is backed by a {@link ConcurrentHashMap} and can use a non-blocking eviction + * thread, providing constant-time {@link #cacheBlock(BlockCacheKey, Cacheable, boolean)} and + * {@link #getBlock(BlockCacheKey, boolean, boolean, boolean)} operations. + *

+ *

+ * The cache maintains three block-priority levels to provide scan resistance and support in-memory + * column families: + *

+ *
    + *
  • single-access blocks
  • + *
  • multiple-access blocks
  • + *
  • in-memory blocks
  • + *
+ *

+ * Each priority is assigned a portion of the total cache capacity. During eviction the cache tries + * to preserve the configured relative sizes while allowing unused capacity in one priority to be + * consumed by another. + *

+ *

+ * This class is a storage engine only. It does not perform L1/L2 orchestration, victim-cache + * delegation, tier placement, admission control, promotion, or demotion. Those responsibilities + * belong to the cache topology and policy layers. + *

+ */ +@InterfaceAudience.Private +public class LruCacheEngine implements CacheEngine, HeapSize, Iterable { + + private static final Logger LOG = LoggerFactory.getLogger(LruCacheEngine.class); + + /** + * Percentage of total size that eviction will evict until. + */ + private static final String LRU_MIN_FACTOR_CONFIG_NAME = "hbase.lru.blockcache.min.factor"; + + /** + * Acceptable cache size above which eviction is triggered. + */ + private static final String LRU_ACCEPTABLE_FACTOR_CONFIG_NAME = + "hbase.lru.blockcache.acceptable.factor"; + + /** + * Hard capacity limit. Inserts are rejected once the cache exceeds this factor multiplied by the + * acceptable size. + */ + static final String LRU_HARD_CAPACITY_LIMIT_FACTOR_CONFIG_NAME = + "hbase.lru.blockcache.hard.capacity.limit.factor"; + + private static final String LRU_SINGLE_PERCENTAGE_CONFIG_NAME = + "hbase.lru.blockcache.single.percentage"; + + private static final String LRU_MULTI_PERCENTAGE_CONFIG_NAME = + "hbase.lru.blockcache.multi.percentage"; + + private static final String LRU_MEMORY_PERCENTAGE_CONFIG_NAME = + "hbase.lru.blockcache.memory.percentage"; + + /** + * Configuration key that gives data blocks from in-memory HFiles higher eviction priority. + */ + private static final String LRU_IN_MEMORY_FORCE_MODE_CONFIG_NAME = + "hbase.lru.rs.inmemoryforcemode"; + + static final float DEFAULT_LOAD_FACTOR = 0.75f; + static final int DEFAULT_CONCURRENCY_LEVEL = 16; + + private static final float DEFAULT_MIN_FACTOR = 0.95f; + static final float DEFAULT_ACCEPTABLE_FACTOR = 0.99f; + + private static final float DEFAULT_SINGLE_FACTOR = 0.25f; + private static final float DEFAULT_MULTI_FACTOR = 0.50f; + private static final float DEFAULT_MEMORY_FACTOR = 0.25f; + + private static final float DEFAULT_HARD_CAPACITY_LIMIT_FACTOR = 1.2f; + + private static final boolean DEFAULT_IN_MEMORY_FORCE_MODE = false; + + private static final int STAT_THREAD_PERIOD = 60 * 5; + + private static final String LRU_MAX_BLOCK_SIZE = "hbase.lru.max.block.size"; + + private static final long DEFAULT_MAX_BLOCK_SIZE = 16L * 1024L * 1024L; + + /** + * Fixed heap overhead of an LRU cache engine instance. + */ + public static final long CACHE_FIXED_OVERHEAD = + ClassSize.estimateBase(LruCacheEngine.class, false); + + /** + * Cached blocks keyed by their HFile block cache key. + *

+ * A {@link ConcurrentHashMap} is required because {@link #getBlock} and eviction depend on the + * atomicity guarantees of {@code computeIfPresent}. + *

+ */ + private transient final ConcurrentHashMap map; + + /** Lock protecting the eviction process. */ + private transient final ReentrantLock evictionLock = new ReentrantLock(true); + + /** Maximum size of an individual block accepted by this cache. */ + private final long maxBlockSize; + + /** Whether an eviction pass is currently running. */ + private volatile boolean evictionInProgress; + + /** Optional background eviction thread. */ + private transient final EvictionThread evictionThread; + + /** + * Listener notified about capacity-driven block evictions. + */ + private volatile CacheEvictionListener evictionListener; + + /** Executor used to periodically report cache statistics. */ + private transient final ScheduledExecutorService scheduleThreadPool = + Executors.newScheduledThreadPool(1, new ThreadFactoryBuilder() + .setNameFormat("LruCacheEngineStatsExecutor").setDaemon(true).build()); + + /** Current total heap size used by the cache. */ + private final AtomicLong size; + + /** Current heap size of data blocks. */ + private final LongAdder dataBlockSize = new LongAdder(); + + /** Current heap size of index blocks. */ + private final LongAdder indexBlockSize = new LongAdder(); + + /** Current heap size of bloom blocks. */ + private final LongAdder bloomBlockSize = new LongAdder(); + + /** Current number of cached blocks. */ + private final AtomicLong elements; + + /** Current number of cached data blocks. */ + private final LongAdder dataBlockElements = new LongAdder(); + + /** Current number of cached index blocks. */ + private final LongAdder indexBlockElements = new LongAdder(); + + /** Current number of cached bloom blocks. */ + private final LongAdder bloomBlockElements = new LongAdder(); + + /** Sequential cache access identifier. */ + private final AtomicLong count; + + /** Hard cache capacity limit factor. */ + private float hardCapacityLimitFactor; + + /** Cache statistics. */ + private final CacheStats stats; + + /** Maximum cache size in bytes. */ + private long maxSize; + + /** Expected average block size. */ + private long blockSize; + + /** Cache size factor at which eviction is triggered. */ + private float acceptableFactor; + + /** Cache size factor to which an eviction pass should reduce the cache. */ + private float minFactor; + + /** Fraction of capacity assigned to single-access blocks. */ + private float singleFactor; + + /** Fraction of capacity assigned to multiple-access blocks. */ + private float multiFactor; + + /** Fraction of capacity assigned to in-memory blocks. */ + private float memoryFactor; + + /** Heap overhead of the cache structure itself. */ + private long overhead; + + /** Whether data blocks from in-memory HFiles receive stronger retention priority. */ + private boolean forceInMemory; + + /** + * Creates an LRU cache engine using default configuration values. + * @param maxSize maximum size of the cache, in bytes + * @param blockSize expected average block size, in bytes + */ + public LruCacheEngine(long maxSize, long blockSize) { + this(maxSize, blockSize, true); + } + + /** + * Creates an LRU cache engine and optionally enables the background eviction thread. + * @param maxSize maximum size of the cache, in bytes + * @param blockSize expected average block size, in bytes + * @param evictionThread whether background eviction should be enabled + */ + public LruCacheEngine(long maxSize, long blockSize, boolean evictionThread) { + this(maxSize, blockSize, evictionThread, (int) Math.ceil(1.2 * maxSize / blockSize), + DEFAULT_LOAD_FACTOR, DEFAULT_CONCURRENCY_LEVEL, DEFAULT_MIN_FACTOR, DEFAULT_ACCEPTABLE_FACTOR, + DEFAULT_SINGLE_FACTOR, DEFAULT_MULTI_FACTOR, DEFAULT_MEMORY_FACTOR, + DEFAULT_HARD_CAPACITY_LIMIT_FACTOR, false, DEFAULT_MAX_BLOCK_SIZE); + } + + /** + * Creates an LRU cache engine using values from the supplied configuration. + * @param maxSize maximum size of the cache, in bytes + * @param blockSize expected average block size, in bytes + * @param evictionThread whether background eviction should be enabled + * @param conf cache configuration + */ + public LruCacheEngine(long maxSize, long blockSize, boolean evictionThread, Configuration conf) { + this(maxSize, blockSize, evictionThread, (int) Math.ceil(1.2 * maxSize / blockSize), + DEFAULT_LOAD_FACTOR, DEFAULT_CONCURRENCY_LEVEL, + conf.getFloat(LRU_MIN_FACTOR_CONFIG_NAME, DEFAULT_MIN_FACTOR), + conf.getFloat(LRU_ACCEPTABLE_FACTOR_CONFIG_NAME, DEFAULT_ACCEPTABLE_FACTOR), + conf.getFloat(LRU_SINGLE_PERCENTAGE_CONFIG_NAME, DEFAULT_SINGLE_FACTOR), + conf.getFloat(LRU_MULTI_PERCENTAGE_CONFIG_NAME, DEFAULT_MULTI_FACTOR), + conf.getFloat(LRU_MEMORY_PERCENTAGE_CONFIG_NAME, DEFAULT_MEMORY_FACTOR), + conf.getFloat(LRU_HARD_CAPACITY_LIMIT_FACTOR_CONFIG_NAME, DEFAULT_HARD_CAPACITY_LIMIT_FACTOR), + conf.getBoolean(LRU_IN_MEMORY_FORCE_MODE_CONFIG_NAME, DEFAULT_IN_MEMORY_FORCE_MODE), + conf.getLong(LRU_MAX_BLOCK_SIZE, DEFAULT_MAX_BLOCK_SIZE)); + } + + /** + * Creates an LRU cache engine using the supplied configuration and background eviction. + * @param maxSize maximum size of the cache, in bytes + * @param blockSize expected average block size, in bytes + * @param conf cache configuration + */ + public LruCacheEngine(long maxSize, long blockSize, Configuration conf) { + this(maxSize, blockSize, true, conf); + } + + /** + * Creates a fully configured LRU cache engine. + * @param maxSize maximum size of this cache, in bytes + * @param blockSize expected average size of blocks, in bytes + * @param evictionThread whether to run eviction in a background thread + * @param mapInitialSize initial size of the backing map + * @param mapLoadFactor load factor of the backing map + * @param mapConcurrencyLevel concurrency level of the backing map + * @param minFactor fraction of maximum size retained after eviction + * @param acceptableFactor fraction of maximum size that triggers eviction + * @param singleFactor fraction assigned to single-access blocks + * @param multiFactor fraction assigned to multiple-access blocks + * @param memoryFactor fraction assigned to in-memory blocks + * @param hardLimitFactor hard capacity limit factor + * @param forceInMemory whether in-memory HFile blocks receive stronger retention priority + * @param maxBlockSize largest individual block accepted by this cache + */ + public LruCacheEngine(long maxSize, long blockSize, boolean evictionThread, int mapInitialSize, + float mapLoadFactor, int mapConcurrencyLevel, float minFactor, float acceptableFactor, + float singleFactor, float multiFactor, float memoryFactor, float hardLimitFactor, + boolean forceInMemory, long maxBlockSize) { + this.maxBlockSize = maxBlockSize; + + if ( + singleFactor + multiFactor + memoryFactor != 1 || singleFactor < 0 || multiFactor < 0 + || memoryFactor < 0 + ) { + throw new IllegalArgumentException( + "Single, multi, and memory factors should be non-negative and total 1.0"); + } + + if (minFactor >= acceptableFactor) { + throw new IllegalArgumentException("minFactor must be smaller than acceptableFactor"); + } + + if (minFactor >= 1.0f || acceptableFactor >= 1.0f) { + throw new IllegalArgumentException("all factors must be < 1"); + } + + this.maxSize = maxSize; + this.blockSize = blockSize; + this.forceInMemory = forceInMemory; + this.map = new ConcurrentHashMap<>(mapInitialSize, mapLoadFactor, mapConcurrencyLevel); + this.minFactor = minFactor; + this.acceptableFactor = acceptableFactor; + this.singleFactor = singleFactor; + this.multiFactor = multiFactor; + this.memoryFactor = memoryFactor; + this.stats = new CacheStats(getClass().getSimpleName()); + this.count = new AtomicLong(0); + this.elements = new AtomicLong(0); + this.overhead = calculateOverhead(maxSize, blockSize, mapConcurrencyLevel); + this.size = new AtomicLong(this.overhead); + this.hardCapacityLimitFactor = hardLimitFactor; + + if (evictionThread) { + this.evictionThread = new EvictionThread(this); + this.evictionThread.start(); + } else { + this.evictionThread = null; + } + + this.scheduleThreadPool.scheduleAtFixedRate(new StatisticsThread(this), STAT_THREAD_PERIOD, + STAT_THREAD_PERIOD, TimeUnit.SECONDS); + } + + /** + * Returns the human-readable name of this cache engine. + * @return cache engine name + */ + @Override + public String getName() { + return getClass().getSimpleName(); + } + + /** + * Updates the maximum size of this cache. + *

+ * If the cache is already larger than the new acceptable size, an eviction pass is started. + *

+ * @param maxSize new maximum size, in bytes + */ + public void setMaxSize(long maxSize) { + this.maxSize = maxSize; + if (size.get() > acceptableSize() && !evictionInProgress) { + runEviction(); + } + } + + /** + * Returns a heap-backed reference suitable for storage in this cache. + *

+ * Shared-memory {@link HFileBlock}s are cloned onto the heap. Other blocks are retained before + * being referenced by the cache. + *

+ * @param buf block to convert + * @return heap-backed retained block + */ + private Cacheable asReferencedHeapBlock(Cacheable buf) { + if (buf instanceof HFileBlock) { + HFileBlock block = (HFileBlock) buf; + if (block.isSharedMem()) { + return HFileBlock.deepCloneOnHeap(block); + } + } + + return buf.retain(); + } + + /** + * Caches the specified block. + * @param cacheKey block cache key + * @param buf block contents + * @param inMemory whether the block should receive in-memory priority + */ + @Override + public void cacheBlock(BlockCacheKey cacheKey, Cacheable buf, boolean inMemory) { + if (buf.heapSize() > maxBlockSize) { + if (stats.failInsert() % 50 == 0) { + LOG.warn("Trying to cache too large a block " + cacheKey.getHfileName() + " @ " + + cacheKey.getOffset() + " is " + buf.heapSize() + " which is larger than " + + maxBlockSize); + } + return; + } + + LruCachedBlock existingBlock = map.get(cacheKey); + if ( + existingBlock != null && !BlockCacheUtil.shouldReplaceExistingCacheBlock(this, cacheKey, buf) + ) { + return; + } + + long currentSize = size.get(); + long currentAcceptableSize = acceptableSize(); + long hardLimitSize = (long) (hardCapacityLimitFactor * currentAcceptableSize); + + if (currentSize >= hardLimitSize) { + stats.failInsert(); + if (LOG.isTraceEnabled()) { + LOG.trace("LruCacheEngine current size " + StringUtils.byteDesc(currentSize) + + " has exceeded acceptable size " + StringUtils.byteDesc(currentAcceptableSize) + "." + + " The hard limit size is " + StringUtils.byteDesc(hardLimitSize) + + ", failed to put cacheKey:" + cacheKey + " into LruCacheEngine."); + } + if (!evictionInProgress) { + runEviction(); + } + return; + } + + Cacheable referencedBlock = asReferencedHeapBlock(buf); + LruCachedBlock newBlock = + new LruCachedBlock(cacheKey, referencedBlock, count.incrementAndGet(), inMemory); + + long newSize; + long elementCount; + + if (existingBlock != null) { + if (!replaceBlock(cacheKey, existingBlock, newBlock)) { + referencedBlock.release(); + return; + } + + newSize = size.get(); + elementCount = elements.get(); + } else { + newSize = updateSizeMetrics(newBlock, false); + map.put(cacheKey, newBlock); + + elementCount = elements.incrementAndGet(); + incrementBlockTypeElementCount(referencedBlock); + } + + if (LOG.isTraceEnabled()) { + assertCounterSanity(map.size(), elementCount); + } + + if (newSize > currentAcceptableSize && !evictionInProgress) { + runEviction(); + } + } + + /** + * Atomically replaces the expected cached block and updates cache accounting. + *

+ * Replacement does not change the total block count because the map cardinality remains + * unchanged. The cache-owned reference held by the replaced block is released before the + * replacement is published. + *

+ * @param cacheKey block cache key + * @param expected block expected to be currently mapped + * @param replacement replacement block + * @return {@code true} if the expected block was replaced + */ + private boolean replaceBlock(BlockCacheKey cacheKey, LruCachedBlock expected, + LruCachedBlock replacement) { + MutableBoolean replaced = new MutableBoolean(false); + + map.computeIfPresent(cacheKey, (key, current) -> { + if (current != expected) { + return current; + } + + updateSizeMetrics(current, true); + decrementBlockTypeElementCount(current.getBuffer()); + current.getBuffer().release(); + + updateSizeMetrics(replacement, false); + incrementBlockTypeElementCount(replacement.getBuffer()); + + replaced.setTrue(); + return replacement; + }); + + return replaced.booleanValue(); + } + + /** + * Increments the cached-element counter associated with the specified block type. + * @param block cached block + */ + private void incrementBlockTypeElementCount(Cacheable block) { + BlockType blockType = block.getBlockType(); + if (blockType.isBloom()) { + bloomBlockElements.increment(); + } else if (blockType.isIndex()) { + indexBlockElements.increment(); + } else if (blockType.isData()) { + dataBlockElements.increment(); + } + } + + /** + * Decrements the cached-element counter associated with the specified block type. + * @param block cached block + */ + private void decrementBlockTypeElementCount(Cacheable block) { + BlockType blockType = block.getBlockType(); + if (blockType.isBloom()) { + bloomBlockElements.decrement(); + } else if (blockType.isIndex()) { + indexBlockElements.decrement(); + } else if (blockType.isData()) { + dataBlockElements.decrement(); + } + } + + /** + * Caches the specified block with normal cache priority. + * @param cacheKey block cache key + * @param buf block contents + */ + @Override + public void cacheBlock(BlockCacheKey cacheKey, Cacheable buf) { + cacheBlock(cacheKey, buf, false); + } + + /** + * Caches the specified block. + *

+ * LRU insertion is synchronous, so {@code waitWhenCache} has no effect. + *

+ * @param cacheKey block cache key + * @param buf block contents + * @param inMemory whether the block should receive in-memory priority + * @param waitWhenCache whether the caller requests synchronous completion + */ + @Override + public void cacheBlock(BlockCacheKey cacheKey, Cacheable buf, boolean inMemory, + boolean waitWhenCache) { + cacheBlock(cacheKey, buf, inMemory); + } + + /** + * Checks consistency between the backing-map size and the element counter. + *

+ * This method is intended for TRACE-level diagnostics and assertion-enabled JVMs. + *

+ * @param mapSize current backing-map size + * @param counterVal current element-counter value + */ + private static void assertCounterSanity(long mapSize, long counterVal) { + if (counterVal < 0) { + LOG.trace("counterVal overflow. Assertions unreliable. counterVal=" + counterVal + + ", mapSize=" + mapSize); + return; + } + + if (mapSize < Integer.MAX_VALUE) { + double percentageDifference = Math.abs((((double) counterVal) / ((double) mapSize)) - 1.0); + if (percentageDifference > 0.05) { + LOG.trace("delta between reported and actual size > 5%. counterVal=" + counterVal + + ", mapSize=" + mapSize); + } + } + } + + /** + * Updates total and block-type-specific size metrics. + * @param cachedBlock cached block whose size should be applied + * @param evict whether this operation represents removal + * @return new total cache size + */ + private long updateSizeMetrics(LruCachedBlock cachedBlock, boolean evict) { + long heapSize = cachedBlock.heapSize(); + BlockType blockType = cachedBlock.getBuffer().getBlockType(); + + if (evict) { + heapSize *= -1; + } + + if (blockType != null) { + if (blockType.isBloom()) { + bloomBlockSize.add(heapSize); + } else if (blockType.isIndex()) { + indexBlockSize.add(heapSize); + } else if (blockType.isData()) { + dataBlockSize.add(heapSize); + } + } + + return size.addAndGet(heapSize); + } + + /** + * Returns the cached block associated with the specified key. + *

+ * Lookup is strictly local to this cache engine. A miss is returned to the topology layer rather + * than being delegated to another cache tier. + *

+ * @param cacheKey block cache key + * @param caching whether the caller caches blocks on misses + * @param repeat whether this is a repeated lookup for the same block + * @param updateCacheMetrics whether cache statistics should be updated + * @return cached block, or {@code null} if not present + */ + @Override + public Cacheable getBlock(BlockCacheKey cacheKey, boolean caching, boolean repeat, + boolean updateCacheMetrics) { + LruCachedBlock cachedBlock = map.computeIfPresent(cacheKey, (key, value) -> { + value.getBuffer().retain(); + return value; + }); + + if (cachedBlock == null) { + if (!repeat && updateCacheMetrics) { + stats.miss(caching, cacheKey.isPrimary(), cacheKey.getBlockType()); + } + return null; + } + + if (updateCacheMetrics) { + stats.hit(caching, cacheKey.isPrimary(), cacheKey.getBlockType()); + } + + cachedBlock.access(count.incrementAndGet()); + return cachedBlock.getBuffer(); + } + + /** + * Returns whether the specified block is currently cached. + * @param cacheKey block cache key + * @return {@code true} if the block is present + */ + public boolean containsBlock(BlockCacheKey cacheKey) { + return map.containsKey(cacheKey); + } + + /** + * Returns whether the specified block is currently cached. + * @param cacheKey block cache key + * @return optional containing the local cache-membership result + */ + @Override + public Optional isAlreadyCached(BlockCacheKey cacheKey) { + return Optional.of(containsBlock(cacheKey)); + } + + /** + * Evicts the specified block. + * @param cacheKey block cache key + * @return {@code true} if a block was found and evicted + */ + @Override + public boolean evictBlock(BlockCacheKey cacheKey) { + LruCachedBlock cachedBlock = map.get(cacheKey); + return cachedBlock != null && evictBlock(cachedBlock, false) > 0; + } + + /** + * Evicts all cached blocks belonging to the specified HFile. + *

+ * This is a linear scan over the cache contents. + *

+ * @param hfileName HFile name + * @return number of blocks evicted + */ + @Override + public int evictBlocksByHfileName(String hfileName) { + int numEvicted = 0; + + for (BlockCacheKey key : map.keySet()) { + if (key.getHfileName().equals(hfileName) && evictBlock(key)) { + numEvicted++; + } + } + + return numEvicted; + } + + /** + * Evicts the specified block from this cache. + *

+ * For capacity-driven evictions, the configured eviction listener is notified while a temporary + * reference to the evicted block is retained. Explicit invalidations do not generate eviction + * notifications. + *

+ * @param block block to evict + * @param evictedByEvictionProcess whether the eviction was caused by cache pressure + * @return heap size of the evicted block, or {@code 0} if the block was not present + */ + protected long evictBlock(LruCachedBlock block, boolean evictedByEvictionProcess) { + final AtomicReference removedBlock = new AtomicReference<>(); + final CacheEvictionListener listener = evictedByEvictionProcess ? evictionListener : null; + + map.computeIfPresent(block.getCacheKey(), (key, value) -> { + if (value != block) { + return value; + } + + Cacheable buffer = value.getBuffer(); + + // Keep the block alive after removing the cache-owned reference. This retain must happen + // while the value is still protected by the map operation. + buffer.retain(); + removedBlock.set(value); + + // Release the reference owned by this cache entry. + buffer.release(); + return null; + }); + + LruCachedBlock removed = removedBlock.get(); + if (removed == null) { + return 0; + } + + Cacheable buffer = removed.getBuffer(); + try { + updateSizeMetrics(removed, true); + + long elementCount = elements.decrementAndGet(); + if (LOG.isTraceEnabled()) { + assertCounterSanity(map.size(), elementCount); + } + + decrementBlockTypeElementCount(buffer); + + if (evictedByEvictionProcess) { + stats.evicted(removed.getCachedTime(), removed.getCacheKey().isPrimary()); + } + + if (listener != null) { + listener.onEviction(this, removed.getCacheKey(), buffer); + } + + return removed.heapSize(); + } finally { + // Release the temporary reference acquired inside computeIfPresent(). + buffer.release(); + } + } + + /** + * Starts an eviction pass synchronously or notifies the background eviction thread. + */ + private void runEviction() { + if (evictionThread == null || !evictionThread.isGo()) { + evict(); + } else { + evictionThread.evict(); + } + } + + /** + * Returns whether an eviction pass is currently executing. + * @return {@code true} if eviction is in progress + */ + boolean isEvictionInProgress() { + return evictionInProgress; + } + + /** + * Returns the calculated fixed and map overhead for this cache. + * @return cache overhead, in bytes + */ + long getOverhead() { + return overhead; + } + + /** + * Performs an LRU eviction pass. + */ + void evict() { + if (!evictionLock.tryLock()) { + return; + } + + try { + evictionInProgress = true; + + long currentSize = size.get(); + long bytesToFree = currentSize - minSize(); + + if (LOG.isTraceEnabled()) { + LOG.trace("LRU cache eviction started; Attempting to free " + + StringUtils.byteDesc(bytesToFree) + " of total=" + StringUtils.byteDesc(currentSize)); + } + + if (bytesToFree <= 0) { + return; + } + + BlockBucket bucketSingle = new BlockBucket("single", bytesToFree, blockSize, singleSize()); + BlockBucket bucketMulti = new BlockBucket("multi", bytesToFree, blockSize, multiSize()); + BlockBucket bucketMemory = new BlockBucket("memory", bytesToFree, blockSize, memorySize()); + + for (LruCachedBlock cachedBlock : map.values()) { + switch (cachedBlock.getPriority()) { + case SINGLE: + bucketSingle.add(cachedBlock); + break; + case MULTI: + bucketMulti.add(cachedBlock); + break; + case MEMORY: + bucketMemory.add(cachedBlock); + break; + default: + throw new IllegalStateException( + "Unsupported block priority " + cachedBlock.getPriority()); + } + } + + long bytesFreed = 0; + + if (forceInMemory || memoryFactor > 0.999f) { + long singleSize = bucketSingle.totalSize(); + long multiSize = bucketMulti.totalSize(); + + if (bytesToFree > singleSize + multiSize) { + bytesFreed = bucketSingle.free(singleSize); + bytesFreed += bucketMulti.free(multiSize); + + if (LOG.isTraceEnabled()) { + LOG.trace( + "freed " + StringUtils.byteDesc(bytesFreed) + " from single and multi buckets"); + } + + bytesFreed += bucketMemory.free(bytesToFree - bytesFreed); + + if (LOG.isTraceEnabled()) { + LOG + .trace("freed " + StringUtils.byteDesc(bytesFreed) + " total from all three buckets"); + } + } else { + long bytesRemain = singleSize + multiSize - bytesToFree; + + if (3 * singleSize <= bytesRemain) { + bytesFreed = bucketMulti.free(bytesToFree); + } else if (3 * multiSize <= 2 * bytesRemain) { + bytesFreed = bucketSingle.free(bytesToFree); + } else { + bytesFreed = bucketSingle.free(singleSize - bytesRemain / 3); + if (bytesFreed < bytesToFree) { + bytesFreed += bucketMulti.free(bytesToFree - bytesFreed); + } + } + } + } else { + PriorityQueue bucketQueue = new PriorityQueue<>(3); + + bucketQueue.add(bucketSingle); + bucketQueue.add(bucketMulti); + bucketQueue.add(bucketMemory); + + int remainingBuckets = bucketQueue.size(); + + BlockBucket bucket; + while ((bucket = bucketQueue.poll()) != null) { + long overflow = bucket.overflow(); + + if (overflow > 0) { + long bucketBytesToFree = + Math.min(overflow, (bytesToFree - bytesFreed) / remainingBuckets); + bytesFreed += bucket.free(bucketBytesToFree); + } + + remainingBuckets--; + } + } + + if (LOG.isTraceEnabled()) { + long single = bucketSingle.totalSize(); + long multi = bucketMulti.totalSize(); + long memory = bucketMemory.totalSize(); + + LOG.trace("LRU cache eviction completed; freed=" + StringUtils.byteDesc(bytesFreed) + + ", total=" + StringUtils.byteDesc(size.get()) + ", single=" + + StringUtils.byteDesc(single) + ", multi=" + StringUtils.byteDesc(multi) + ", memory=" + + StringUtils.byteDesc(memory)); + } + } finally { + stats.evict(); + evictionInProgress = false; + evictionLock.unlock(); + } + } + + /** + * Returns a diagnostic string describing this cache. + * @return cache description + */ + @Override + public String toString() { + return MoreObjects.toStringHelper(this).add("blockCount", getBlockCount()) + .add("currentSize", StringUtils.byteDesc(getCurrentSize())) + .add("freeSize", StringUtils.byteDesc(getFreeSize())) + .add("maxSize", StringUtils.byteDesc(getMaxSize())) + .add("heapSize", StringUtils.byteDesc(heapSize())) + .add("minSize", StringUtils.byteDesc(minSize())).add("minFactor", minFactor) + .add("multiSize", StringUtils.byteDesc(multiSize())).add("multiFactor", multiFactor) + .add("singleSize", StringUtils.byteDesc(singleSize())).add("singleFactor", singleFactor) + .toString(); + } + + /** + * Groups cached blocks belonging to the same LRU priority. + */ + private class BlockBucket implements Comparable { + + private final String name; + private final LruCachedBlockQueue queue; + private long totalSize; + private final long bucketSize; + + /** + * Creates an LRU priority bucket. + * @param name bucket name + * @param bytesToFree number of bytes the eviction pass attempts to free + * @param blockSize expected average block size + * @param bucketSize target size of this priority bucket + */ + BlockBucket(String name, long bytesToFree, long blockSize, long bucketSize) { + this.name = name; + this.bucketSize = bucketSize; + this.queue = new LruCachedBlockQueue(bytesToFree, blockSize); + } + + /** + * Adds a block to this priority bucket. + * @param block cached block + */ + void add(LruCachedBlock block) { + totalSize += block.heapSize(); + queue.add(block); + } + + /** + * Evicts least-recently-used blocks until at least the requested number of bytes has been + * released or the bucket is exhausted. + * @param toFree requested number of bytes to release + * @return number of bytes actually released + */ + long free(long toFree) { + if (LOG.isTraceEnabled()) { + LOG.trace("freeing " + StringUtils.byteDesc(toFree) + " from " + this); + } + + LruCachedBlock cachedBlock; + long freedBytes = 0; + + while ((cachedBlock = queue.pollLast()) != null) { + freedBytes += evictBlock(cachedBlock, true); + if (freedBytes >= toFree) { + return freedBytes; + } + } + + if (LOG.isTraceEnabled()) { + LOG.trace("freed " + StringUtils.byteDesc(freedBytes) + " from " + this); + } + + return freedBytes; + } + + /** + * Returns the number of bytes by which this bucket exceeds its target size. + * @return bucket overflow in bytes + */ + long overflow() { + return totalSize - bucketSize; + } + + /** + * Returns the total size of blocks assigned to this bucket. + * @return bucket size in bytes + */ + long totalSize() { + return totalSize; + } + + /** + * Compares this bucket with another bucket by overflow. + * @param that other bucket + * @return comparison result + */ + @Override + public int compareTo(BlockBucket that) { + return Long.compare(overflow(), that.overflow()); + } + + /** + * Returns whether another object represents a bucket with the same overflow. + * @param that object to compare + * @return {@code true} when the buckets compare equally + */ + @Override + public boolean equals(Object that) { + if (!(that instanceof LruCacheEngine.BlockBucket)) { + return false; + } + return compareTo((BlockBucket) that) == 0; + } + + /** + * Returns the hash code of this bucket. + * @return bucket hash code + */ + @Override + public int hashCode() { + return Objects.hashCode(name, bucketSize, queue, totalSize); + } + + /** + * Returns a diagnostic string describing this bucket. + * @return bucket description + */ + @Override + public String toString() { + return MoreObjects.toStringHelper(this).add("name", name) + .add("totalSize", StringUtils.byteDesc(totalSize)) + .add("bucketSize", StringUtils.byteDesc(bucketSize)).toString(); + } + } + + /** + * Returns the maximum cache size. + * @return maximum cache size in bytes + */ + @Override + public long getMaxSize() { + return maxSize; + } + + /** + * Returns the current total heap size consumed by this cache. + * @return current size in bytes + */ + @Override + public long getCurrentSize() { + return size.get(); + } + + /** + * Returns the current heap size of cached data blocks. + * @return data block size in bytes + */ + @Override + public long getCurrentDataSize() { + return dataBlockSize.sum(); + } + + /** + * Returns the current heap size of cached index blocks. + * @return index block size in bytes + */ + public long getCurrentIndexSize() { + return indexBlockSize.sum(); + } + + /** + * Returns the current heap size of cached bloom blocks. + * @return bloom block size in bytes + */ + public long getCurrentBloomSize() { + return bloomBlockSize.sum(); + } + + /** + * Returns unused cache capacity. + * @return free size in bytes + */ + @Override + public long getFreeSize() { + return getMaxSize() - getCurrentSize(); + } + + /** + * Returns the configured capacity of this cache. + *

+ * This preserves the existing LRU cache behavior where {@code size()} reports configured cache + * capacity rather than current occupancy. + *

+ * @return configured cache capacity in bytes + */ + @Override + public long size() { + return getMaxSize(); + } + + /** + * Returns the total number of blocks currently cached. + * @return cached block count + */ + @Override + public long getBlockCount() { + return elements.get(); + } + + /** + * Returns the number of data blocks currently cached. + * @return cached data block count + */ + @Override + public long getDataBlockCount() { + return dataBlockElements.sum(); + } + + /** + * Returns the number of index blocks currently cached. + * @return cached index block count + */ + public long getIndexBlockCount() { + return indexBlockElements.sum(); + } + + /** + * Returns the number of bloom blocks currently cached. + * @return cached bloom block count + */ + public long getBloomBlockCount() { + return bloomBlockElements.sum(); + } + + /** + * Returns the background eviction thread. + * @return eviction thread, or {@code null} if background eviction is disabled + */ + EvictionThread getEvictionThread() { + return evictionThread; + } + + /** + * Background thread responsible for initiating LRU eviction passes. + */ + static class EvictionThread extends Thread { + + private final WeakReference cache; + + private volatile boolean go = true; + + private boolean enteringRun; + + /** + * Creates an eviction thread for the specified cache. + * @param cache owning cache engine + */ + EvictionThread(LruCacheEngine cache) { + super(Thread.currentThread().getName() + ".LruCacheEngine.EvictionThread"); + setDaemon(true); + this.cache = new WeakReference<>(cache); + } + + /** + * Waits for eviction notifications and executes eviction passes. + */ + @Override + public void run() { + enteringRun = true; + + while (go) { + synchronized (this) { + try { + wait(1000 * 10); + } catch (InterruptedException e) { + LOG.warn("Interrupted eviction thread", e); + Thread.currentThread().interrupt(); + } + } + + LruCacheEngine cacheEngine = cache.get(); + if (cacheEngine == null) { + go = false; + break; + } + + cacheEngine.evict(); + } + } + + /** + * Wakes this thread so that it can execute an eviction pass. + */ + @edu.umd.cs.findbugs.annotations.SuppressWarnings(value = "NN_NAKED_NOTIFY", + justification = "This is what we want") + void evict() { + synchronized (this) { + notifyAll(); + } + } + + /** + * Stops the background eviction thread. + */ + synchronized void shutdown() { + go = false; + notifyAll(); + } + + /** + * Returns whether this thread should continue running. + * @return {@code true} while the thread is active + */ + boolean isGo() { + return go; + } + + /** + * Returns whether this thread has entered its run method. + * @return {@code true} after {@link #run()} has started + */ + boolean isEnteringRun() { + return enteringRun; + } + } + + /** + * Periodically logs LRU cache statistics. + */ + static class StatisticsThread extends Thread { + + private final LruCacheEngine lru; + + /** + * Creates a statistics thread for the specified cache. + * @param lru cache whose statistics should be logged + */ + StatisticsThread(LruCacheEngine lru) { + super("LruCacheEngineStats"); + setDaemon(true); + this.lru = lru; + } + + /** + * Logs the current cache statistics. + */ + @Override + public void run() { + lru.logStats(); + } + } + + /** + * Logs current LRU cache size and access statistics. + */ + public void logStats() { + long usedSize = heapSize(); + long freeSize = maxSize - usedSize; + + LOG.info("totalSize=" + StringUtils.byteDesc(maxSize) + ", usedSize=" + + StringUtils.byteDesc(usedSize) + ", freeSize=" + StringUtils.byteDesc(freeSize) + ", max=" + + StringUtils.byteDesc(maxSize) + ", blockCount=" + getBlockCount() + ", accesses=" + + stats.getRequestCount() + ", hits=" + stats.getHitCount() + ", hitRatio=" + + (stats.getHitCount() == 0 ? "0" : StringUtils.formatPercent(stats.getHitRatio(), 2) + ", ") + + ", cachingAccesses=" + stats.getRequestCachingCount() + ", cachingHits=" + + stats.getHitCachingCount() + ", cachingHitsRatio=" + + (stats.getHitCachingCount() == 0 + ? "0," + : StringUtils.formatPercent(stats.getHitCachingRatio(), 2) + ", ") + + "evictions=" + stats.getEvictionCount() + ", evicted=" + stats.getEvictedCount() + + ", evictedPerRun=" + stats.evictedPerEviction()); + } + + /** + * Returns cache access and eviction statistics. + * @return cache statistics + */ + @Override + public CacheStats getStats() { + return stats; + } + + /** + * Returns the total heap size currently consumed by this cache. + * @return current heap size in bytes + */ + @Override + public long heapSize() { + return getCurrentSize(); + } + + /** + * Calculates the estimated fixed and backing-map overhead for a cache. + * @param maxSize maximum cache size + * @param blockSize expected average block size + * @param concurrency backing-map concurrency level + * @return estimated overhead in bytes + */ + private static long calculateOverhead(long maxSize, long blockSize, int concurrency) { + return CACHE_FIXED_OVERHEAD + ClassSize.CONCURRENT_HASHMAP + + ((long) Math.ceil(maxSize * 1.2 / blockSize) * ClassSize.CONCURRENT_HASHMAP_ENTRY) + + ((long) concurrency * ClassSize.CONCURRENT_HASHMAP_SEGMENT); + } + + /** + * Returns an iterator over the cached blocks. + * @return cached block iterator + */ + @Override + public Iterator iterator() { + final Iterator iterator = map.values().iterator(); + + return new Iterator() { + + private final long now = System.nanoTime(); + + /** + * Returns whether another cached block is available. + * @return {@code true} if another cached block is available + */ + @Override + public boolean hasNext() { + return iterator.hasNext(); + } + + /** + * Returns the next cached block. + * @return next cached block + */ + @Override + public CachedBlock next() { + final LruCachedBlock block = iterator.next(); + + return new CachedBlock() { + + /** + * Returns a diagnostic representation of this cached block. + * @return cached block description + */ + @Override + public String toString() { + return BlockCacheUtil.toString(this, now); + } + + /** + * Returns the LRU priority assigned to this cached block. + * @return block priority + */ + @Override + public BlockPriority getBlockPriority() { + return block.getPriority(); + } + + /** + * Returns the HFile block type. + * @return block type + */ + @Override + public BlockType getBlockType() { + return block.getBuffer().getBlockType(); + } + + /** + * Returns the block offset in its HFile. + * @return HFile offset + */ + @Override + public long getOffset() { + return block.getCacheKey().getOffset(); + } + + /** + * Returns the heap size of the cached block. + * @return block size in bytes + */ + @Override + public long getSize() { + return block.getBuffer().heapSize(); + } + + /** + * Returns the time at which the block was cached. + * @return cached time + */ + @Override + public long getCachedTime() { + return block.getCachedTime(); + } + + /** + * Returns the name of the HFile containing this block. + * @return HFile name + */ + @Override + public String getFilename() { + return block.getCacheKey().getHfileName(); + } + + /** + * Compares cached blocks by filename, offset, and cache time. + * @param other cached block to compare + * @return comparison result + */ + @Override + public int compareTo(CachedBlock other) { + int difference = getFilename().compareTo(other.getFilename()); + if (difference != 0) { + return difference; + } + + difference = Long.compare(getOffset(), other.getOffset()); + if (difference != 0) { + return difference; + } + + if (other.getCachedTime() < 0 || getCachedTime() < 0) { + throw new IllegalStateException(getCachedTime() + ", " + other.getCachedTime()); + } + + return Long.compare(other.getCachedTime(), getCachedTime()); + } + + /** + * Returns the hash code of the underlying cached block. + * @return cached block hash code + */ + @Override + public int hashCode() { + return block.hashCode(); + } + + /** + * Returns whether another object represents the same cached block. + * @param object object to compare + * @return {@code true} when the objects represent the same cached block + */ + @Override + public boolean equals(Object object) { + if (!(object instanceof CachedBlock)) { + return false; + } + + return compareTo((CachedBlock) object) == 0; + } + }; + } + + /** + * Removal through this iterator is unsupported. + * @throws UnsupportedOperationException always + */ + @Override + public void remove() { + throw new UnsupportedOperationException(); + } + }; + } + + /** + * Returns an iterable view over cached blocks. + * @return optional containing this cache as a cached-block iterable + */ + @Override + public Optional> asCachedBlockIterable() { + return Optional.of(this); + } + + /** + * Returns the acceptable size above which eviction is triggered. + * @return acceptable size in bytes + */ + @InterfaceAudience.LimitedPrivate(HBaseInterfaceAudience.UNITTEST) + public long acceptableSize() { + return (long) Math.floor(maxSize * acceptableFactor); + } + + /** + * Returns the target size below which an eviction pass should reduce the cache. + * @return minimum target size in bytes + */ + private long minSize() { + return (long) Math.floor(maxSize * minFactor); + } + + /** + * Returns the target capacity assigned to single-access blocks. + * @return single-access bucket size in bytes + */ + private long singleSize() { + return (long) Math.floor(maxSize * singleFactor * minFactor); + } + + /** + * Returns the target capacity assigned to multiple-access blocks. + * @return multiple-access bucket size in bytes + */ + private long multiSize() { + return (long) Math.floor(maxSize * multiFactor * minFactor); + } + + /** + * Returns the target capacity assigned to in-memory blocks. + * @return in-memory bucket size in bytes + */ + private long memorySize() { + return (long) Math.floor(maxSize * memoryFactor * minFactor); + } + + /** + * Shuts down this cache engine and its background threads. + */ + @Override + public void shutdown() { + scheduleThreadPool.shutdown(); + + for (int i = 0; i < 10; i++) { + if (!scheduleThreadPool.isShutdown()) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + LOG.warn("Interrupted while sleeping"); + Thread.currentThread().interrupt(); + break; + } + } + } + + if (!scheduleThreadPool.isShutdown()) { + List runnables = scheduleThreadPool.shutdownNow(); + LOG.debug("Still running " + runnables); + } + + if (evictionThread != null) { + evictionThread.shutdown(); + } + } + + /** + * Clears all cached blocks. + *

+ * This method is intended for tests. + *

+ */ + public void clearCache() { + map.clear(); + elements.set(0); + } + + /** + * Sets the listener that receives capacity-driven block eviction events. + * @param listener eviction listener, or {@code null} to clear the current listener + */ + @Override + public void setEvictionListener(CacheEvictionListener listener) { + this.evictionListener = listener; + } + + /** + * Returns the names of files that currently have blocks in this cache. + *

+ * This method is intended for tests and performs a full cache scan. + *

+ * @return sorted set of cached HFile names + */ + SortedSet getCachedFileNamesForTest() { + SortedSet fileNames = new TreeSet<>(); + + for (BlockCacheKey cacheKey : map.keySet()) { + fileNames.add(cacheKey.getHfileName()); + } + + return fileNames; + } + + /** + * Returns counts of cached blocks grouped by data block encoding. + *

+ * This method is intended for tests. + *

+ * @return block counts grouped by data block encoding + */ + public Map getEncodingCountsForTest() { + Map counts = new EnumMap<>(DataBlockEncoding.class); + + for (LruCachedBlock cachedBlock : map.values()) { + DataBlockEncoding encoding = ((HFileBlock) cachedBlock.getBuffer()).getDataBlockEncoding(); + Integer currentCount = counts.get(encoding); + counts.put(encoding, currentCount == null ? 1 : currentCount + 1); + } + + return counts; + } + + /** + * Returns the internal cached-block map. + *

+ * This method is intended for tests. + *

+ * @return internal cached-block map + */ + Map getMapForTests() { + return map; + } + +} diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/SingleTierTopology.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/SingleTierTopology.java index 9dfec5507678..7967fdc01e0c 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/SingleTierTopology.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/SingleTierTopology.java @@ -177,4 +177,19 @@ public boolean demote(BlockCacheKey cacheKey, Cacheable block, CacheEngine sourc public void shutdown() { engine.shutdown(); } + + /** + * Handles a capacity-driven eviction from the single cache engine. + *

+ * A single-tier topology has no lower tier to which a block can be demoted. + *

+ * @param cacheKey key identifying the evicted block + * @param block evicted block + * @param sourceEngine engine that evicted the block + * @return {@code false} + */ + @Override + public boolean handleEviction(BlockCacheKey cacheKey, Cacheable block, CacheEngine sourceEngine) { + return false; + } } diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredExclusiveTopology.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredExclusiveTopology.java index d20d5c8c34e9..a1aa6561ab67 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredExclusiveTopology.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredExclusiveTopology.java @@ -43,12 +43,15 @@ public class TieredExclusiveTopology implements CacheTopology { private final CacheEngine l1; private final CacheEngine l2; private final CacheTopologyView view; + private final CacheStats stats; public TieredExclusiveTopology(String name, CacheEngine l1, CacheEngine l2) { this.name = name; this.l1 = l1; this.l2 = l2; this.view = new CacheTopologyView(this); + this.stats = new AggregateCacheStats(name, l1.getStats(), l2.getStats()); + } @Override @@ -90,8 +93,7 @@ public CacheTopologyView getView() { @Override public CacheStats getStats() { - // TODO: replace with aggregate topology stats in follow-up metrics ticket. - return l1.getStats(); + return stats; } @Override @@ -123,4 +125,25 @@ public void shutdown() { l1.shutdown(); l2.shutdown(); } + + /** + * Handles a capacity-driven eviction from an engine in this exclusive topology. + *

+ * An L1 pressure eviction is demoted to L2. An L2 pressure eviction leaves the topology entirely, + * because there is no lower tier. + *

+ * @param cacheKey key identifying the evicted block + * @param block evicted block + * @param sourceEngine engine that evicted the block + * @return {@code true} if the block was demoted to L2 + */ + @Override + public boolean handleEviction(BlockCacheKey cacheKey, Cacheable block, CacheEngine sourceEngine) { + if (sourceEngine != l1) { + return false; + } + + l2.cacheBlock(cacheKey, block); + return true; + } } diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredInclusiveTopology.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredInclusiveTopology.java index 984e7ca98ca0..6cc55b37c4db 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredInclusiveTopology.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TieredInclusiveTopology.java @@ -43,12 +43,14 @@ public class TieredInclusiveTopology implements CacheTopology { private final CacheEngine l1; private final CacheEngine l2; private final CacheTopologyView view; + private final CacheStats stats; public TieredInclusiveTopology(String name, CacheEngine l1, CacheEngine l2) { this.name = name; this.l1 = l1; this.l2 = l2; this.view = new CacheTopologyView(this); + this.stats = new AggregateCacheStats(name, l1.getStats(), l2.getStats()); } @Override @@ -90,8 +92,7 @@ public CacheTopologyView getView() { @Override public CacheStats getStats() { - // TODO: replace with aggregate topology stats in follow-up metrics ticket. - return l1.getStats(); + return stats; } @Override @@ -110,4 +111,26 @@ public void shutdown() { l1.shutdown(); l2.shutdown(); } + + /** + * Handles a capacity-driven eviction from this inclusive topology. + *

+ * L1 eviction does not require demotion because inclusive placement normally maintains a + * corresponding block in L2. Cache placement is best-effort and is not atomic across tiers, so + * there may be short windows where an L2 copy is not present. The topology deliberately does not + * perform an L2 membership check on every L1 eviction to avoid adding cross-tier lookup overhead + * to the eviction path. A missing copy results only in a subsequent cache miss. + *

+ *

+ * L2 pressure eviction likewise does not cause movement to another tier. + *

+ * @param cacheKey key identifying the evicted block + * @param block evicted block + * @param sourceEngine engine that evicted the block + * @return {@code false}, because no additional placement is required + */ + @Override + public boolean handleEviction(BlockCacheKey cacheKey, Cacheable block, CacheEngine sourceEngine) { + return false; + } } diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessService.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessService.java index f98ddc6f262f..80f5e901b7aa 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessService.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessService.java @@ -73,6 +73,10 @@ public TopologyBackedCacheAccessService(CacheTopology topology, this.policy = Objects.requireNonNull(policy, "policy must not be null"); this.topologyView = Objects.requireNonNull(topology.getView(), "topology view must not be null"); + for (CacheTier tier : topology.getTiers()) { + topology.getEngine(tier) + .ifPresent(engine -> engine.setEvictionListener(this::handleEngineEviction)); + } } /** @@ -138,6 +142,7 @@ public Cacheable getBlock(BlockCacheKey cacheKey, CacheRequestContext context) { private Cacheable getBlockFromTieredExclusiveTopology(BlockCacheKey cacheKey, CacheRequestContext context) { + Optional l1 = topology.getEngine(CacheTier.L1); Optional l2 = topology.getEngine(CacheTier.L2); @@ -163,11 +168,6 @@ private Cacheable getBlockFromTieredExclusiveTopology(BlockCacheKey cacheKey, } Cacheable block = getBlockFromEngine(selectedEngine, cacheKey, context); - boolean updateCacheMetrics = context.isUpdateCacheMetrics(); - boolean caching = context.isCaching(); - if (updateCacheMetrics) { - updateBlockMetrics(block, cacheKey, selectedEngine, caching); - } if (block != null) { maybePromote(cacheKey, block, selectedTier, selectedEngine, context); @@ -195,19 +195,6 @@ private Cacheable getBlockFromAllTiers(BlockCacheKey cacheKey, CacheRequestConte return null; } - private void updateBlockMetrics(Cacheable block, BlockCacheKey key, CacheEngine engine, - boolean caching) { - CacheStats stats = engine.getStats(); - if (stats == null) { - return; - } - if (block == null) { - stats.miss(caching, key.isPrimary(), key.getBlockType()); - } else { - stats.hit(caching, key.isPrimary(), key.getBlockType()); - } - } - /** * Caches a block using the topology-backed cache access service. *

@@ -743,4 +730,15 @@ public CachedBlock next() { public Iterator iterator() { return asCachedBlockIterable().orElse(Collections.emptyList()).iterator(); } + + /** + * Handles a capacity-driven eviction reported by a cache engine. + * @param sourceEngine engine that evicted the block + * @param cacheKey key identifying the evicted block + * @param block evicted block + */ + private void handleEngineEviction(CacheEngine sourceEngine, BlockCacheKey cacheKey, + Cacheable block) { + topology.handleEviction(cacheKey, block, sourceEngine); + } } diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessServices.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessServices.java index 22810f526607..fada7eb2d8bd 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessServices.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/TopologyBackedCacheAccessServices.java @@ -86,16 +86,12 @@ public static TopologyBackedCacheAccessService fromCombinedBlockCache( } /** - * Creates a topology-backed cache access service from existing L1 and L2 block caches. - *

- * The resulting service uses {@link TieredExclusiveTopology}, which models the current - * CombinedBlockCache-compatible L1/L2 behavior where promotion can move a block from one tier to - * another. - *

- * @param name human-readable topology/service name - * @param l1 L1 block cache - * @param l2 L2 block cache - * @param policy placement and admission policy + * Creates a topology-backed cache access service by adapting two legacy block caches as an + * exclusive tiered topology. + * @param name topology name + * @param l1 first-level legacy block cache + * @param l2 second-level legacy block cache + * @param policy cache placement and admission policy * @return topology-backed cache access service */ public static TopologyBackedCacheAccessService fromTieredExclusiveBlockCaches(String name, @@ -107,10 +103,8 @@ public static TopologyBackedCacheAccessService fromTieredExclusiveBlockCaches(St if (l1 instanceof FirstLevelBlockCache) { ((FirstLevelBlockCache) l1).unsetVictimCache(); } - CacheEngine l1Engine = CacheEngines.fromBlockCache(l1);// fromL1BlockCache(l1); - CacheEngine l2Engine = CacheEngines.fromBlockCache(l2); - CacheTopology topology = new TieredExclusiveTopology(name, l1Engine, l2Engine); - return new TopologyBackedCacheAccessService(topology, policy); + return fromTieredExclusiveCacheEngines(name, CacheEngines.fromBlockCache(l1), + CacheEngines.fromBlockCache(l2), policy); } /** @@ -146,27 +140,13 @@ public static TopologyBackedCacheAccessService fromTieredExclusiveBlockCaches(St } /** - * Creates a topology-backed cache access service from two legacy block caches using an inclusive - * tiered topology. - *

- * The first supplied block cache is treated as the L1 tier and the second supplied block cache is - * treated as the L2 tier. Both legacy caches are adapted to {@link CacheEngine} instances using - * {@link CacheEngines#fromBlockCache(BlockCache)} and then assembled into a - * {@link TieredInclusiveTopology}. - *

- *

- * This helper is intended for compatibility with legacy inclusive combined-cache configurations - * while moving cache access and diagnostics to the {@link CacheAccessService} abstraction. - * Inclusive topology semantics differ from exclusive topology semantics: a block may exist in - * both tiers, and eviction from one tier does not necessarily imply eviction from the other tier. - *

- * @param name topology name used for diagnostics - * @param l1 first-level block cache - * @param l2 second-level block cache - * @param policy cache placement and admission policy to use with the topology-backed service - * @return topology-backed cache access service backed by a tiered inclusive topology - * @throws NullPointerException if {@code name}, {@code l1}, {@code l2}, or {@code policy} is - * {@code null} + * Creates a topology-backed cache access service by adapting two legacy block caches as an + * inclusive tiered topology. + * @param name topology name + * @param l1 first-level legacy block cache + * @param l2 second-level legacy block cache + * @param policy cache placement and admission policy + * @return topology-backed cache access service */ public static TopologyBackedCacheAccessService fromTieredInclusiveBlockCaches(String name, BlockCache l1, BlockCache l2, CachePlacementAdmissionPolicy policy) { @@ -177,37 +157,24 @@ public static TopologyBackedCacheAccessService fromTieredInclusiveBlockCaches(St if (l1 instanceof FirstLevelBlockCache) { ((FirstLevelBlockCache) l1).unsetVictimCache(); } - CacheEngine l1Engine = CacheEngines.fromBlockCache(l1); - CacheEngine l2Engine = CacheEngines.fromBlockCache(l2); - CacheTopology topology = new TieredInclusiveTopology(name, l1Engine, l2Engine); - return new TopologyBackedCacheAccessService(topology, policy); + return fromTieredInclusiveCacheEngines(name, CacheEngines.fromBlockCache(l1), + CacheEngines.fromBlockCache(l2), policy); } /** - * Creates a topology-backed cache access service for a single legacy {@link BlockCache}. - *

- * The supplied block cache is adapted to a {@link CacheEngine} and placed behind a - * {@link SingleTierTopology}. This makes single-tier caches use the same - * {@link TopologyBackedCacheAccessService} path as combined caches while preserving the existing - * block cache implementation underneath. - *

- * @param name topology name used for diagnostics - * @param blockCache legacy block cache to adapt + * Creates a topology-backed cache access service by adapting a legacy block cache as a single + * cache engine. + * @param name topology name + * @param blockCache legacy block cache * @param policy cache placement and admission policy - * @return topology-backed cache access service backed by a single-tier topology - * @throws NullPointerException if {@code name}, {@code blockCache}, or {@code policy} is - * {@code null} + * @return topology-backed cache access service */ public static TopologyBackedCacheAccessService fromSingleBlockCache(String name, BlockCache blockCache, CachePlacementAdmissionPolicy policy) { - Objects.requireNonNull(name, "name must not be null"); Objects.requireNonNull(blockCache, "blockCache must not be null"); - Objects.requireNonNull(policy, "policy must not be null"); - CacheEngine engine = CacheEngines.fromBlockCache(blockCache); - CacheTopology topology = new SingleTierTopology(name, engine); - return new TopologyBackedCacheAccessService(topology, policy); + return fromSingleCacheEngine(name, CacheEngines.fromBlockCache(blockCache), policy); } /** @@ -289,4 +256,62 @@ public static BlockCache getBlockCache(CacheAccessService cacheAccessService) { return getBlockCache(cacheAccessService, CacheTier.SINGLE); } + + /** + * Creates a topology-backed cache access service using a single cache engine. + * @param name topology name + * @param engine cache engine + * @param policy cache placement and admission policy + * @return topology-backed cache access service + */ + public static TopologyBackedCacheAccessService fromSingleCacheEngine(String name, + CacheEngine engine, CachePlacementAdmissionPolicy policy) { + Objects.requireNonNull(name, "name must not be null"); + Objects.requireNonNull(engine, "engine must not be null"); + Objects.requireNonNull(policy, "policy must not be null"); + + CacheTopology topology = new SingleTierTopology(name, engine); + return new TopologyBackedCacheAccessService(topology, policy); + } + + /** + * Creates a topology-backed cache access service using two independent cache engines in an + * exclusive tiered topology. + * @param name topology name + * @param l1 first-level cache engine + * @param l2 second-level cache engine + * @param policy cache placement and admission policy + * @return topology-backed cache access service + */ + public static TopologyBackedCacheAccessService fromTieredExclusiveCacheEngines(String name, + CacheEngine l1, CacheEngine l2, CachePlacementAdmissionPolicy policy) { + Objects.requireNonNull(name, "name must not be null"); + Objects.requireNonNull(l1, "l1 must not be null"); + Objects.requireNonNull(l2, "l2 must not be null"); + Objects.requireNonNull(policy, "policy must not be null"); + + CacheTopology topology = new TieredExclusiveTopology(name, l1, l2); + return new TopologyBackedCacheAccessService(topology, policy); + } + + /** + * Creates a topology-backed cache access service using two independent cache engines in an + * inclusive tiered topology. + * @param name topology name + * @param l1 first-level cache engine + * @param l2 second-level cache engine + * @param policy cache placement and admission policy + * @return topology-backed cache access service + */ + public static TopologyBackedCacheAccessService fromTieredInclusiveCacheEngines(String name, + CacheEngine l1, CacheEngine l2, CachePlacementAdmissionPolicy policy) { + Objects.requireNonNull(name, "name must not be null"); + Objects.requireNonNull(l1, "l1 must not be null"); + Objects.requireNonNull(l2, "l2 must not be null"); + Objects.requireNonNull(policy, "policy must not be null"); + + CacheTopology topology = new TieredInclusiveTopology(name, l1, l2); + return new TopologyBackedCacheAccessService(topology, policy); + } + } diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheConfig.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheConfig.java index 2aa5a8d69912..4952f0d90439 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheConfig.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheConfig.java @@ -21,14 +21,10 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; -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 static org.junit.jupiter.api.Assertions.fail; import java.io.IOException; -import java.lang.management.ManagementFactory; -import java.lang.management.MemoryUsage; import java.nio.ByteBuffer; import java.util.concurrent.atomic.AtomicInteger; import org.apache.hadoop.conf.Configuration; @@ -43,8 +39,12 @@ import org.apache.hadoop.hbase.io.ByteBuffAllocator; import org.apache.hadoop.hbase.io.hfile.BlockType.BlockCategory; import org.apache.hadoop.hbase.io.hfile.bucket.BucketCache; +import org.apache.hadoop.hbase.io.hfile.cache.BlockCacheBackedCacheEngine; import org.apache.hadoop.hbase.io.hfile.cache.CacheAccessService; import org.apache.hadoop.hbase.io.hfile.cache.CacheAccessServiceTestFactory; +import org.apache.hadoop.hbase.io.hfile.cache.CacheEngine; +import org.apache.hadoop.hbase.io.hfile.cache.CacheTier; +import org.apache.hadoop.hbase.io.hfile.cache.LruCacheEngine; import org.apache.hadoop.hbase.io.hfile.cache.NoOpCacheAccessService; import org.apache.hadoop.hbase.io.hfile.cache.TopologyBackedCacheAccessService; import org.apache.hadoop.hbase.io.hfile.cache.TopologyBackedCacheAccessServices; @@ -279,7 +279,8 @@ public void testCacheConfigDefaultLRUBlockCache() { assertTrue(CacheConfig.DEFAULT_IN_MEMORY == cc.isInMemory()); CacheAccessService service = CacheAccessServiceTestFactory.fromConfiguration(this.conf); basicBlockCacheOps(service, cc, false, true); - assertTrue(CacheAccessServiceTestFactory.blockCache(service) instanceof LruBlockCache); + CacheEngine l1 = CacheAccessServiceTestFactory.getCacheEngine(service, CacheTier.SINGLE); + assertTrue(l1 instanceof LruCacheEngine); } /** @@ -305,136 +306,56 @@ public void testFileBucketCacheConfig() throws IOException { } } + /** + * Verifies BucketCache configuration with a native LRU first-level cache engine. + */ private void doBucketCacheConfigTest() { final int bcSize = 100; this.conf.setInt(HConstants.BUCKET_CACHE_SIZE_KEY, bcSize); + CacheConfig cc = new CacheConfig(this.conf); CacheAccessService service = CacheAccessServiceTestFactory.fromConfiguration(this.conf); + basicBlockCacheOps(service, cc, false, false); assertTrue(CacheAccessServiceTestFactory.isCombinedBlockCacheEquivalent(service)); - // TODO: Assert sizes allocated are right and proportions. - LruBlockCache lbc = - (LruBlockCache) CacheAccessServiceTestFactory.getFirstLevelBlockCache(service); - assertEquals(MemorySizeUtil.getOnHeapCacheSize(this.conf), lbc.getMaxSize()); - BucketCache bc = (BucketCache) CacheAccessServiceTestFactory.getSecondLevelBlockCache(service); - // getMaxSize comes back in bytes but we specified size in MB - assertEquals(bcSize, bc.getMaxSize() / (1024 * 1024)); + + CacheEngine l1 = CacheAccessServiceTestFactory.getCacheEngine(service, CacheTier.L1); + assertInstanceOf(LruCacheEngine.class, l1); + assertEquals(MemorySizeUtil.getOnHeapCacheSize(this.conf), l1.getMaxSize()); + + BlockCache l2 = CacheAccessServiceTestFactory.getSecondLevelBlockCache(service); + assertInstanceOf(BucketCache.class, l2); + assertEquals(bcSize, l2.getMaxSize() / (1024 * 1024)); } - /** - * Verifies the legacy two-tier block cache layout used when bucket cache is enabled but the - * combined-cache mode is disabled. - *

- * In this configuration HBase should deploy an {@link LruBlockCache} as the first-level in-memory - * cache and a {@link BucketCache} as the second-level victim cache. Blocks are inserted into L1 - * first. When L1 evicts blocks under memory pressure, the evicted blocks should be passed to the - * configured L2 victim cache. - *

- *

- * This test intentionally verifies the L1-to-L2 victim-cache relationship without relying on an - * exact final L1 block count or on a specific block key being evicted. {@link LruBlockCache} - * eviction policy does not guarantee which block will be selected for eviction, only that some - * blocks may be evicted when cache pressure exceeds the configured threshold. - *

- *

- * The previous version of this test attempted to force eviction by inserting a single synthetic - * block whose size was {@code acceptableSize() + 1}, and then waited until the L1 block count - * returned to its original value. That approach was flawed for two reasons: - *

- *
    - *
  1. {@link LruBlockCache} rejects any block larger than its maximum cacheable block size before - * eviction can run. Since {@code acceptableSize()} depends on the JVM heap size while the maximum - * cacheable block size is fixed by configuration, {@code acceptableSize() + 1} may be larger than - * the maximum cacheable block size. In that case the block is rejected and no eviction is - * triggered.
  2. - *
  3. If the synthetic block is small enough to be accepted, eviction runs only after the block - * is inserted. The eviction policy is not required to restore the exact previous block count, nor - * is it required to evict the originally inserted block. Waiting for an exact L1 block count can - * therefore hang indefinitely.
  4. - *
- *

- * The test now creates cache pressure using normal cacheable blocks and waits with a timeout - * until L2 receives at least one block from L1 eviction. This directly verifies the intended - * contract: L1 is wired with L2 as its victim cache. - *

- */ @Test public void testBucketCacheConfigL1L2Setup() throws Exception { this.conf.set(HConstants.BUCKET_CACHE_IOENGINE_KEY, "offheap"); - // this.conf.setLong("hbase.lru.max.block.size", 1L << 30); - // from L1 happens, it does not fail because L2 can't take the eviction because block too big. this.conf.setFloat(HConstants.HFILE_BLOCK_CACHE_SIZE_KEY, 0.001f); - MemoryUsage mu = ManagementFactory.getMemoryMXBean().getHeapMemoryUsage(); + long lruExpectedSize = MemorySizeUtil.getOnHeapCacheSize(this.conf); final int bcSize = 100; - long bcExpectedSize = 100 * 1024 * 1024; // MB. + long bcExpectedSize = 100 * 1024 * 1024; assertTrue(lruExpectedSize < bcExpectedSize); this.conf.setInt(HConstants.BUCKET_CACHE_SIZE_KEY, bcSize); + CacheConfig cc = new CacheConfig(this.conf); CacheAccessService service = CacheAccessServiceTestFactory.fromConfiguration(this.conf); + basicBlockCacheOps(service, cc, false, false); assertTrue(CacheAccessServiceTestFactory.isCombinedBlockCacheEquivalent(service)); - // TODO: Assert sizes allocated are right and proportions. - FirstLevelBlockCache lbc = - (FirstLevelBlockCache) CacheAccessServiceTestFactory.getFirstLevelBlockCache(service); - assertEquals(lruExpectedSize, lbc.getMaxSize()); - BlockCache bc = CacheAccessServiceTestFactory.getSecondLevelBlockCache(service); - // getMaxSize comes back in bytes but we specified size in MB - assertEquals(bcExpectedSize, ((BucketCache) bc).getMaxSize()); - /* - * The topology-backed cache path intentionally clears legacy L1 victim-cache wiring when L1 and - * L2 are adapted as independent topology engines. Direct calls to the unwrapped L1 cache should - * therefore not be used to verify L1-to-L2 victim movement. Tier placement, promotion, and - * lookup are now owned by TopologyBackedCacheAccessService. - */ - long initialL1BlockCount = lbc.getBlockCount(); - long initialL2BlockCount = bc.getBlockCount(); - Cacheable c = new DataCacheEntry(); - BlockCacheKey bck = new BlockCacheKey("bck", 0); - - lbc.cacheBlock(bck, c, false); - - assertEquals(initialL1BlockCount + 1, lbc.getBlockCount()); - assertEquals(initialL2BlockCount, bc.getBlockCount()); - assertNotNull(lbc.getBlock(bck, true, false, true)); - assertNull(bc.getBlock(bck, true, false, true)); - } - /** - * Adds cacheable blocks to L1 until L1 eviction moves at least one block into L2. - *

- * The helper does not wait for a particular key to appear in L2. {@link LruBlockCache} eviction - * is policy-driven and does not guarantee that the first inserted block, or any specific later - * block, will be evicted first. The observable contract needed by this test is only that an L1 - * eviction is forwarded to the configured L2 victim cache. - *

- *

- * This helper also avoids using a single oversized block to force eviction. Oversized blocks may - * be rejected by {@link LruBlockCache} before eviction can run. Instead, it inserts regular - * cacheable blocks and relies on cumulative cache pressure. - *

- * @param l1Cache first-level cache - * @param l2Cache second-level victim cache - * @param initialL2BlockCount L2 block count before creating L1 pressure - * @throws Exception if the expected L1-to-L2 movement does not happen before the wait timeout - */ - private void waitForAnyBlockToMoveFromL1ToL2(FirstLevelBlockCache l1Cache, BlockCache l2Cache, - long initialL2BlockCount) throws Exception { - AtomicInteger blockIndex = new AtomicInteger(); + CacheEngine l1 = CacheAccessServiceTestFactory.getCacheEngine(service, CacheTier.L1); + assertTrue(l1 instanceof LruCacheEngine); + assertEquals(lruExpectedSize, l1.getMaxSize()); - /* - * Do not try to force eviction with one block of size acceptableSize() + 1. LruBlockCache - * rejects blocks larger than maxBlockSize before eviction can run. For accepted blocks, - * eviction runs after insertion and does not guarantee which block will be evicted. Therefore - * this test should not wait for a particular block key to appear in L2. The intended contract - * is only that an L1 eviction moves some evicted block into the configured L2 victim cache. - */ - Waiter.waitFor(this.conf, 10000, () -> { - BlockCacheKey evictionKey = new BlockCacheKey("eviction-" + blockIndex.getAndIncrement(), 0); - l1Cache.cacheBlock(evictionKey, new DataCacheEntry(), false); + CacheEngine l2 = CacheAccessServiceTestFactory.getCacheEngine(service, CacheTier.L2); + assertTrue(l2 instanceof BlockCacheBackedCacheEngine); + assertEquals(bcExpectedSize, l2.getMaxSize()); - return l2Cache.getBlockCount() > initialL2BlockCount; - }); + BlockCache blockCache = CacheAccessServiceTestFactory.getSecondLevelBlockCache(service); + assertTrue(blockCache instanceof BucketCache); + assertEquals(bcExpectedSize, blockCache.getMaxSize()); } @Test @@ -510,4 +431,41 @@ void testCacheAccessServiceIsNoOpWhenBlockCacheIsNull() { assertInstanceOf(NoOpCacheAccessService.class, service); assertFalse(service.isCacheEnabled()); } + + /** + * Verifies that a capacity-driven eviction from the native L1 cache engine is propagated to the + * configured L2 cache engine. + * @throws Exception if waiting for an L1 eviction to reach L2 fails + */ + @Test + public void testL1CapacityEvictionMovesBlockToL2() throws Exception { + this.conf.set(HConstants.BUCKET_CACHE_IOENGINE_KEY, "offheap"); + this.conf.setFloat(HConstants.HFILE_BLOCK_CACHE_SIZE_KEY, 0.001f); + this.conf.setInt(HConstants.BUCKET_CACHE_SIZE_KEY, 100); + this.conf.set(BlockCacheFactory.BUCKET_CACHE_BUCKETS_KEY, Integer.toString(2 * 1024 * 1024)); + + CacheAccessService service = CacheAccessServiceTestFactory.fromConfiguration(this.conf); + + CacheEngine l1 = CacheAccessServiceTestFactory.getCacheEngine(service, CacheTier.L1); + assertTrue(l1 instanceof LruCacheEngine); + + CacheEngine l2 = CacheAccessServiceTestFactory.getCacheEngine(service, CacheTier.L2); + assertTrue(l2 instanceof BlockCacheBackedCacheEngine); + + ((LruCacheEngine) l1).setMaxSize(4 * 1024 * 1024); + + long initialL2BlockCount = l2.getBlockCount(); + AtomicInteger blockIndex = new AtomicInteger(); + + Waiter.waitFor(this.conf, 10000, () -> { + BlockCacheKey evictionKey = new BlockCacheKey("eviction-" + blockIndex.getAndIncrement(), 0); + DataCacheEntry entry = new DataCacheEntry(); + LOG.info("entry heapSize={}, L1 currentSize={}, L1 maxSize={}", entry.heapSize(), + l1.getCurrentSize(), l1.getMaxSize()); + l1.cacheBlock(evictionKey, entry, false); + return l2.getBlockCount() > initialL2BlockCount; + }); + + assertTrue(l2.getBlockCount() > initialL2BlockCount); + } } diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheOnWrite.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheOnWrite.java index 75a03d425042..514db64b4bea 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheOnWrite.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestCacheOnWrite.java @@ -166,8 +166,8 @@ private static List getCacheServices() throws IOException { Configuration conf = TEST_UTIL.getConfiguration(); List caches = new ArrayList<>(); // default - caches.add(CacheAccessServiceTestFactory.fromConfiguration(conf)); - + BlockCache defaultBlockCache = BlockCacheFactory.createBlockCache(conf); + caches.add(CacheAccessServices.fromBlockCache(defaultBlockCache)); // set LruBlockCache.LRU_HARD_CAPACITY_LIMIT_FACTOR_CONFIG_NAME to 2.0f due to HBASE-16287 TEST_UTIL.getConfiguration().setFloat(LruBlockCache.LRU_HARD_CAPACITY_LIMIT_FACTOR_CONFIG_NAME, 2.0f); diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestForceCacheImportantBlocks.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestForceCacheImportantBlocks.java index c7a0dd2e4896..d60c4df88944 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestForceCacheImportantBlocks.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestForceCacheImportantBlocks.java @@ -32,6 +32,7 @@ import org.apache.hadoop.hbase.io.compress.Compression.Algorithm; import org.apache.hadoop.hbase.io.hfile.cache.CacheAccessService; import org.apache.hadoop.hbase.io.hfile.cache.CacheAccessServiceTestFactory; +import org.apache.hadoop.hbase.io.hfile.cache.CacheAccessServices; import org.apache.hadoop.hbase.regionserver.BloomType; import org.apache.hadoop.hbase.regionserver.HRegion; import org.apache.hadoop.hbase.testclassification.IOTests; @@ -98,8 +99,9 @@ public void setup() { public void testCacheBlocks() throws IOException { // Set index block size to be the same as normal block size. TEST_UTIL.getConfiguration().setInt(HFileBlockIndex.MAX_CHUNK_SIZE_KEY, BLOCK_SIZE); - CacheAccessService cache = - CacheAccessServiceTestFactory.fromConfiguration(TEST_UTIL.getConfiguration()); + + BlockCache blockCache = BlockCacheFactory.createBlockCache(TEST_UTIL.getConfiguration()); + CacheAccessService cache = CacheAccessServices.fromBlockCache(blockCache); ColumnFamilyDescriptor cfd = ColumnFamilyDescriptorBuilder.newBuilder(Bytes.toBytes(CF)).setMaxVersions(MAX_VERSIONS) .setCompressionType(COMPRESSION_ALGORITHM).setBloomFilterType(BLOOM_TYPE) @@ -117,8 +119,11 @@ public void testCacheBlocks() throws IOException { assertTrue(HFile.DATABLOCK_READ_COUNT.sum() > 0); long missCount = stats.getMissCount(); region.get(new Get(Bytes.toBytes("row" + 0))); - if (this.cfCacheEnabled) assertEquals(missCount, stats.getMissCount()); - else assertTrue(stats.getMissCount() > missCount); + if (this.cfCacheEnabled) { + assertEquals(missCount, stats.getMissCount()); + } else { + assertTrue(stats.getMissCount() > missCount); + } } private void writeTestData(HRegion region) throws IOException { diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestHFile.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestHFile.java index 2c6cf46c5b9e..214cc0665c45 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestHFile.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/TestHFile.java @@ -79,6 +79,10 @@ import org.apache.hadoop.hbase.io.hfile.ReaderContext.ReaderType; import org.apache.hadoop.hbase.io.hfile.cache.CacheAccessService; import org.apache.hadoop.hbase.io.hfile.cache.CacheAccessServiceTestFactory; +import org.apache.hadoop.hbase.io.hfile.cache.CacheEngine; +import org.apache.hadoop.hbase.io.hfile.cache.CacheTier; +import org.apache.hadoop.hbase.io.hfile.cache.LruCacheEngine; +import org.apache.hadoop.hbase.io.hfile.cache.TopologyBackedCacheAccessService; import org.apache.hadoop.hbase.monitoring.ThreadLocalServerSideScanMetrics; import org.apache.hadoop.hbase.nio.ByteBuff; import org.apache.hadoop.hbase.nio.RefCnt; @@ -272,16 +276,24 @@ private void assertBytesReadFromCache(boolean isScanMetricsEnabled, DataBlockEnc Path storeFilePath = writeStoreFile(); // Initialize the block cache and HFile reader - CacheAccessService lru = CacheAccessServiceTestFactory.fromConfiguration(conf); - assertTrue(CacheAccessServiceTestFactory.blockCache(lru) instanceof LruBlockCache); - CacheConfig cacheConfig = new CacheConfig(conf, null, lru, ByteBuffAllocator.HEAP); + CacheAccessService cacheAccessService = CacheAccessServiceTestFactory.fromConfiguration(conf); + + assertTrue(cacheAccessService instanceof TopologyBackedCacheAccessService); + + TopologyBackedCacheAccessService topologyService = + (TopologyBackedCacheAccessService) cacheAccessService; + CacheEngine engine = topologyService.getTopology().getEngine(CacheTier.SINGLE).orElseThrow(); + assertTrue(engine instanceof LruCacheEngine); + + CacheConfig cacheConfig = + new CacheConfig(conf, null, cacheAccessService, ByteBuffAllocator.HEAP); HFileReaderImpl reader = (HFileReaderImpl) HFile.createReader(fs, storeFilePath, cacheConfig, true, conf); // Read the first block in HFile from the block cache. final int offset = 0; BlockCacheKey cacheKey = new BlockCacheKey(storeFilePath.getName(), offset); - HFileBlock block = (HFileBlock) lru.getBlock(cacheKey, false, false, true); + HFileBlock block = (HFileBlock) cacheAccessService.getBlock(cacheKey, false, false, true); assertNull(block); // Assert that first block has not been cached in the block cache and no disk I/O happened to diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServiceTestFactory.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServiceTestFactory.java index dc31fab4607f..8bacfa241571 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServiceTestFactory.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/CacheAccessServiceTestFactory.java @@ -753,4 +753,26 @@ public static BlockCache getBlockCache(CacheAccessService cacheAccessService, Ca return ((BlockCacheBackedCacheEngine) engine).getBlockCache(); } + + /** + * Returns the cache engine for the requested tier. + * @param service cache access service + * @param tier cache tier + * @return cache engine for the requested tier + * @throws IllegalArgumentException if the service is not topology-backed or the requested tier + * does not exist + */ + public static CacheEngine getCacheEngine(CacheAccessService service, CacheTier tier) { + Objects.requireNonNull(service, "service must not be null"); + Objects.requireNonNull(tier, "tier must not be null"); + + if (!(service instanceof TopologyBackedCacheAccessService)) { + throw new IllegalArgumentException("Cache access service is not topology-backed"); + } + + TopologyBackedCacheAccessService topologyService = (TopologyBackedCacheAccessService) service; + + return topologyService.getTopology().getEngine(tier) + .orElseThrow(() -> new IllegalArgumentException("Cache tier is not present: " + tier)); + } } diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/TestAggregateCacheStats.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/TestAggregateCacheStats.java new file mode 100644 index 000000000000..179859238e9f --- /dev/null +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/TestAggregateCacheStats.java @@ -0,0 +1,114 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +import org.apache.hadoop.hbase.io.hfile.BlockType; +import org.apache.hadoop.hbase.io.hfile.CacheStats; +import org.apache.hadoop.hbase.testclassification.IOTests; +import org.apache.hadoop.hbase.testclassification.SmallTests; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; + +/** + * Tests for {@link AggregateCacheStats}. + */ +@Tag(IOTests.TAG) +@Tag(SmallTests.TAG) +public class TestAggregateCacheStats { + + /** + * Verifies that hit, miss, request, and eviction counts are aggregated across all delegates. + */ + @Test + void testAggregatesCacheStatistics() { + CacheStats l1 = new CacheStats("l1"); + CacheStats l2 = new CacheStats("l2"); + + l1.hit(true, true, BlockType.LEAF_INDEX); + l1.miss(true, true, BlockType.DATA); + l1.evict(); + l1.evicted(1L, true); + + l2.hit(true, true, BlockType.DATA); + l2.hit(true, false, BlockType.DATA); + l2.miss(true, false, BlockType.DATA); + l2.evict(); + l2.evicted(2L, false); + + CacheStats stats = new AggregateCacheStats("aggregate", l1, l2); + + assertEquals(3L, stats.getHitCount()); + assertEquals(2L, stats.getMissCount()); + assertEquals(5L, stats.getRequestCount()); + + assertEquals(1L, stats.getLeafIndexHitCount()); + assertEquals(2L, stats.getDataHitCount()); + assertEquals(2L, stats.getDataMissCount()); + + assertEquals(2L, stats.getPrimaryHitCount()); + assertEquals(1L, stats.getPrimaryMissCount()); + + assertEquals(3L, stats.getHitCachingCount()); + assertEquals(2L, stats.getMissCachingCount()); + assertEquals(5L, stats.getRequestCachingCount()); + + assertEquals(2L, stats.getEvictionCount()); + assertEquals(2L, stats.getEvictedCount()); + assertEquals(1L, stats.getPrimaryEvictedCount()); + } + + /** + * Verifies that failed insertion counts are aggregated across all delegates. + */ + @Test + void testAggregatesFailedInserts() { + CacheStats l1 = new CacheStats("l1"); + CacheStats l2 = new CacheStats("l2"); + + l1.failInsert(); + l2.failInsert(); + l2.failInsert(); + + CacheStats stats = new AggregateCacheStats("aggregate", l1, l2); + + assertEquals(3L, stats.getFailedInserts()); + } + + /** + * Verifies that rolling the aggregate statistics rolls all delegate statistics. + */ + @Test + void testRollMetricsPeriod() { + CacheStats l1 = new CacheStats("l1"); + CacheStats l2 = new CacheStats("l2"); + + l1.hit(true, true, BlockType.DATA); + l2.miss(true, true, BlockType.DATA); + + CacheStats stats = new AggregateCacheStats("aggregate", l1, l2); + + stats.rollMetricsPeriod(); + + assertEquals(1L, stats.getSumHitCountsPastNPeriods()); + assertEquals(2L, stats.getSumRequestCountsPastNPeriods()); + assertEquals(1L, stats.getSumHitCachingCountsPastNPeriods()); + assertEquals(2L, stats.getSumRequestCachingCountsPastNPeriods()); + } +} diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/TestCacheEngineConstruction.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/TestCacheEngineConstruction.java new file mode 100644 index 000000000000..376577c4464a --- /dev/null +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/io/hfile/cache/TestCacheEngineConstruction.java @@ -0,0 +1,55 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache; + +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.HBaseConfiguration; +import org.apache.hadoop.hbase.HConstants; +import org.apache.hadoop.hbase.io.hfile.BlockCacheFactory; +import org.apache.hadoop.hbase.testclassification.IOTests; +import org.apache.hadoop.hbase.testclassification.SmallTests; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; + +@Tag(IOTests.TAG) +@Tag(SmallTests.TAG) +class TestCacheEngineConstruction { + + /** + * Verifies that the LRU cache policy creates a native {@link LruCacheEngine}. + */ + @Test + void testCreateFirstLevelCacheEngineWithLruPolicy() { + Configuration conf = HBaseConfiguration.create(); + conf.setFloat(HConstants.HFILE_BLOCK_CACHE_SIZE_KEY, 0.01f); + conf.set(BlockCacheFactory.BLOCKCACHE_POLICY_KEY, "LRU"); + + CacheEngine cacheEngine = BlockCacheFactory.createFirstLevelCacheEngine(conf); + try { + assertNotNull(cacheEngine); + assertTrue(cacheEngine instanceof LruCacheEngine); + } finally { + if (cacheEngine != null) { + cacheEngine.shutdown(); + } + } + } +}