001/* 002 * Licensed to the Apache Software Foundation (ASF) under one 003 * or more contributor license agreements. See the NOTICE file 004 * distributed with this work for additional information 005 * regarding copyright ownership. The ASF licenses this file 006 * to you under the Apache License, Version 2.0 (the 007 * "License"); you may not use this file except in compliance 008 * with the License. You may obtain a copy of the License at 009 * 010 * http://www.apache.org/licenses/LICENSE-2.0 011 * 012 * Unless required by applicable law or agreed to in writing, software 013 * distributed under the License is distributed on an "AS IS" BASIS, 014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 015 * See the License for the specific language governing permissions and 016 * limitations under the License. 017 */ 018package org.apache.hadoop.hbase.io.hfile.bucket; 019 020import java.io.FileOutputStream; 021import java.io.IOException; 022import java.util.Comparator; 023import java.util.HashMap; 024import java.util.Map; 025import java.util.NavigableSet; 026import java.util.concurrent.ConcurrentHashMap; 027import java.util.concurrent.ConcurrentSkipListSet; 028import java.util.function.Consumer; 029import java.util.function.Function; 030import org.apache.hadoop.hbase.io.ByteBuffAllocator; 031import org.apache.hadoop.hbase.io.ByteBuffAllocator.Recycler; 032import org.apache.hadoop.hbase.io.hfile.BlockCacheKey; 033import org.apache.hadoop.hbase.io.hfile.BlockPriority; 034import org.apache.hadoop.hbase.io.hfile.BlockType; 035import org.apache.hadoop.hbase.io.hfile.CacheableDeserializerIdManager; 036import org.apache.hadoop.hbase.io.hfile.HFileBlock; 037import org.apache.hadoop.hbase.util.Pair; 038import org.apache.yetus.audience.InterfaceAudience; 039 040import org.apache.hbase.thirdparty.com.google.protobuf.ByteString; 041 042import org.apache.hadoop.hbase.shaded.protobuf.generated.BucketCacheProtos; 043 044@InterfaceAudience.Private 045final class BucketProtoUtils { 046 047 final static byte[] PB_MAGIC_V2 = new byte[] { 'V', '2', 'U', 'F' }; 048 049 private BucketProtoUtils() { 050 051 } 052 053 static BucketCacheProtos.BucketCacheEntry toPB(BucketCache cache, 054 BucketCacheProtos.BackingMap.Builder backingMapBuilder) { 055 return BucketCacheProtos.BucketCacheEntry.newBuilder().setCacheCapacity(cache.getMaxSize()) 056 .setIoClass(cache.ioEngine.getClass().getName()) 057 .setMapClass(cache.backingMap.getClass().getName()) 058 .putAllDeserializers(CacheableDeserializerIdManager.save()) 059 .putAllCachedFiles(toCachedPB(cache.fullyCachedFiles)) 060 .setBackingMap(backingMapBuilder.build()) 061 .setChecksum(ByteString 062 .copyFrom(((PersistentIOEngine) cache.ioEngine).calculateChecksum(cache.getAlgorithm()))) 063 .build(); 064 } 065 066 public static void serializeAsPB(BucketCache cache, FileOutputStream fos, long chunkSize) 067 throws IOException { 068 serializeAsPB(cache, fos, chunkSize, entry -> { 069 }); 070 } 071 072 static void serializeAsPB(BucketCache cache, FileOutputStream fos, long chunkSize, 073 Consumer<Map.Entry<BlockCacheKey, BucketEntry>> entryCopiedAction) throws IOException { 074 // Write the new version of magic number. 075 fos.write(PB_MAGIC_V2); 076 077 BucketCacheProtos.BackingMap.Builder builder = BucketCacheProtos.BackingMap.newBuilder(); 078 BucketCacheProtos.BackingMapEntry.Builder entryBuilder = 079 BucketCacheProtos.BackingMapEntry.newBuilder(); 080 081 // Persist the metadata first. 082 toPB(cache, builder).writeDelimitedTo(fos); 083 084 int blockCount = 0; 085 // Persist backing map entries in chunks of size 'chunkSize'. 086 for (Map.Entry<BlockCacheKey, BucketEntry> entry : cache.backingMap.entrySet()) { 087 blockCount++; 088 addEntryToBuilder(entry, entryBuilder, builder); 089 entryCopiedAction.accept(entry); 090 if (blockCount % chunkSize == 0) { 091 builder.build().writeDelimitedTo(fos); 092 builder.clear(); 093 } 094 } 095 // Persist the last chunk. 096 if (builder.getEntryList().size() > 0) { 097 builder.build().writeDelimitedTo(fos); 098 } 099 } 100 101 private static void addEntryToBuilder(Map.Entry<BlockCacheKey, BucketEntry> entry, 102 BucketCacheProtos.BackingMapEntry.Builder entryBuilder, 103 BucketCacheProtos.BackingMap.Builder builder) { 104 entryBuilder.clear(); 105 entryBuilder.setKey(BucketProtoUtils.toPB(entry.getKey())); 106 entryBuilder.setValue(BucketProtoUtils.toPB(entry.getValue())); 107 builder.addEntry(entryBuilder.build()); 108 } 109 110 private static BucketCacheProtos.BlockCacheKey toPB(BlockCacheKey key) { 111 BucketCacheProtos.BlockCacheKey.Builder builder = BucketCacheProtos.BlockCacheKey.newBuilder() 112 .setHfilename(key.getHfileName()).setOffset(key.getOffset()) 113 .setPrimaryReplicaBlock(key.isPrimary()).setBlockType(toPB(key.getBlockType())); 114 if (key.getCfName() != null) { 115 builder.setFamilyName(key.getCfName()); 116 } 117 if (key.getRegionName() != null) { 118 builder.setRegionName(key.getRegionName()); 119 } 120 return builder.build(); 121 } 122 123 private static BucketCacheProtos.BlockType toPB(BlockType blockType) { 124 switch (blockType) { 125 case DATA: 126 return BucketCacheProtos.BlockType.data; 127 case META: 128 return BucketCacheProtos.BlockType.meta; 129 case TRAILER: 130 return BucketCacheProtos.BlockType.trailer; 131 case INDEX_V1: 132 return BucketCacheProtos.BlockType.index_v1; 133 case FILE_INFO: 134 return BucketCacheProtos.BlockType.file_info; 135 case LEAF_INDEX: 136 return BucketCacheProtos.BlockType.leaf_index; 137 case ROOT_INDEX: 138 return BucketCacheProtos.BlockType.root_index; 139 case BLOOM_CHUNK: 140 return BucketCacheProtos.BlockType.bloom_chunk; 141 case ENCODED_DATA: 142 return BucketCacheProtos.BlockType.encoded_data; 143 case GENERAL_BLOOM_META: 144 return BucketCacheProtos.BlockType.general_bloom_meta; 145 case INTERMEDIATE_INDEX: 146 return BucketCacheProtos.BlockType.intermediate_index; 147 case DELETE_FAMILY_BLOOM_META: 148 return BucketCacheProtos.BlockType.delete_family_bloom_meta; 149 default: 150 throw new Error("Unrecognized BlockType."); 151 } 152 } 153 154 private static BucketCacheProtos.BucketEntry toPB(BucketEntry entry) { 155 return BucketCacheProtos.BucketEntry.newBuilder().setOffset(entry.offset()) 156 .setCachedTime(entry.getCachedTime()).setLength(entry.getLength()) 157 .setDiskSizeWithHeader(entry.getOnDiskSizeWithHeader()) 158 .setDeserialiserIndex(entry.deserializerIndex).setAccessCounter(entry.getAccessCounter()) 159 .setPriority(toPB(entry.getPriority())).build(); 160 } 161 162 private static BucketCacheProtos.BlockPriority toPB(BlockPriority p) { 163 switch (p) { 164 case MULTI: 165 return BucketCacheProtos.BlockPriority.multi; 166 case MEMORY: 167 return BucketCacheProtos.BlockPriority.memory; 168 case SINGLE: 169 return BucketCacheProtos.BlockPriority.single; 170 default: 171 throw new Error("Unrecognized BlockPriority."); 172 } 173 } 174 175 static Pair<ConcurrentHashMap<BlockCacheKey, BucketEntry>, NavigableSet<BlockCacheKey>> fromPB( 176 Map<Integer, String> deserializers, BucketCacheProtos.BackingMap backingMap, 177 Function<BucketEntry, Recycler> createRecycler) throws IOException { 178 ConcurrentHashMap<BlockCacheKey, BucketEntry> result = new ConcurrentHashMap<>(); 179 NavigableSet<BlockCacheKey> resultSet = new ConcurrentSkipListSet<>(Comparator 180 .comparing(BlockCacheKey::getHfileName).thenComparingLong(BlockCacheKey::getOffset)); 181 for (BucketCacheProtos.BackingMapEntry entry : backingMap.getEntryList()) { 182 BucketCacheProtos.BlockCacheKey protoKey = entry.getKey(); 183 BlockCacheKey key = new BlockCacheKey(protoKey.getHfilename(), protoKey.getFamilyName(), 184 protoKey.getRegionName(), protoKey.getOffset(), protoKey.getPrimaryReplicaBlock(), 185 fromPb(protoKey.getBlockType()), protoKey.getArchived()); 186 BucketCacheProtos.BucketEntry protoValue = entry.getValue(); 187 // TODO:We use ByteBuffAllocator.HEAP here, because we could not get the ByteBuffAllocator 188 // which created by RpcServer elegantly. 189 BucketEntry value = new BucketEntry(protoValue.getOffset(), protoValue.getLength(), 190 protoValue.getDiskSizeWithHeader(), protoValue.getAccessCounter(), 191 protoValue.getCachedTime(), 192 protoValue.getPriority() == BucketCacheProtos.BlockPriority.memory, createRecycler, 193 ByteBuffAllocator.HEAP); 194 // This is the deserializer that we stored 195 int oldIndex = protoValue.getDeserialiserIndex(); 196 String deserializerClass = deserializers.get(oldIndex); 197 if (deserializerClass == null) { 198 throw new IOException("Found deserializer index without matching entry."); 199 } 200 // Convert it to the identifier for the deserializer that we have in this runtime 201 if (deserializerClass.equals(HFileBlock.BlockDeserializer.class.getName())) { 202 int actualIndex = HFileBlock.BLOCK_DESERIALIZER.getDeserializerIdentifier(); 203 value.deserializerIndex = (byte) actualIndex; 204 } else { 205 // We could make this more plugable, but right now HFileBlock is the only implementation 206 // of Cacheable outside of tests, so this might not ever matter. 207 throw new IOException("Unknown deserializer class found: " + deserializerClass); 208 } 209 result.put(key, value); 210 resultSet.add(key); 211 } 212 return new Pair<>(result, resultSet); 213 } 214 215 private static BlockType fromPb(BucketCacheProtos.BlockType blockType) { 216 switch (blockType) { 217 case data: 218 return BlockType.DATA; 219 case meta: 220 return BlockType.META; 221 case trailer: 222 return BlockType.TRAILER; 223 case index_v1: 224 return BlockType.INDEX_V1; 225 case file_info: 226 return BlockType.FILE_INFO; 227 case leaf_index: 228 return BlockType.LEAF_INDEX; 229 case root_index: 230 return BlockType.ROOT_INDEX; 231 case bloom_chunk: 232 return BlockType.BLOOM_CHUNK; 233 case encoded_data: 234 return BlockType.ENCODED_DATA; 235 case general_bloom_meta: 236 return BlockType.GENERAL_BLOOM_META; 237 case intermediate_index: 238 return BlockType.INTERMEDIATE_INDEX; 239 case delete_family_bloom_meta: 240 return BlockType.DELETE_FAMILY_BLOOM_META; 241 default: 242 throw new Error("Unrecognized BlockType."); 243 } 244 } 245 246 static Map<String, BucketCacheProtos.RegionFileSizeMap> 247 toCachedPB(Map<String, Pair<String, Long>> prefetchedHfileNames) { 248 Map<String, BucketCacheProtos.RegionFileSizeMap> tmpMap = new HashMap<>(); 249 prefetchedHfileNames.forEach((hfileName, regionPrefetchMap) -> { 250 BucketCacheProtos.RegionFileSizeMap tmpRegionFileSize = 251 BucketCacheProtos.RegionFileSizeMap.newBuilder().setRegionName(regionPrefetchMap.getFirst()) 252 .setRegionCachedSize(regionPrefetchMap.getSecond()).build(); 253 tmpMap.put(hfileName, tmpRegionFileSize); 254 }); 255 return tmpMap; 256 } 257 258 static Map<String, Pair<String, Long>> 259 fromPB(Map<String, BucketCacheProtos.RegionFileSizeMap> prefetchHFileNames) { 260 Map<String, Pair<String, Long>> hfileMap = new HashMap<>(); 261 prefetchHFileNames.forEach((hfileName, regionPrefetchMap) -> { 262 hfileMap.put(hfileName, 263 new Pair<>(regionPrefetchMap.getRegionName(), regionPrefetchMap.getRegionCachedSize())); 264 }); 265 return hfileMap; 266 } 267}