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}