This is an automated email from the ASF dual-hosted git repository. jianyun pushed a commit to branch rocksdb/dev in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit e1fb08fea92f7a3ae51c4717edf3a033e4b7f0bd Author: chengjianyun <[email protected]> AuthorDate: Fri Mar 11 19:11:36 2022 +0800 [RocksDB] code refine --- .../org/apache/iotdb/db/metadata/IMetaManager.java | 3 - .../org/apache/iotdb/db/metadata/MManager.java | 5 - .../iotdb/db/metadata/rocksdb/MRocksDBManager.java | 132 +-------------------- .../iotdb/db/metadata/rocksdb/RMNodeValueType.java | 18 +++ .../iotdb/db/metadata/rocksdb/RockDBConstants.java | 2 +- .../iotdb/db/metadata/rocksdb/RocksDBUtils.java | 78 ++++++------ .../metadata/rocksdb/mnode/RStorageGroupMNode.java | 4 +- 7 files changed, 67 insertions(+), 175 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/IMetaManager.java b/server/src/main/java/org/apache/iotdb/db/metadata/IMetaManager.java index fcc12da..37a0823 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/IMetaManager.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/IMetaManager.java @@ -247,9 +247,6 @@ public interface IMetaManager { List<MeasurementPath> getAllMeasurementByDevicePath(PartialPath devicePath) throws PathNotExistException; - Map<PartialPath, IMeasurementSchema> getAllMeasurementSchemaByPrefix(PartialPath prefixPath) - throws MetadataException; - IStorageGroupMNode getStorageGroupNodeByStorageGroupPath(PartialPath path) throws MetadataException; diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java index 8efe84e..60fdea1 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java @@ -1510,11 +1510,6 @@ public class MManager implements IMetaManager { return new ArrayList<>(res); } - @Override - public Map<PartialPath, IMeasurementSchema> getAllMeasurementSchemaByPrefix( - PartialPath prefixPath) throws MetadataException { - return mtree.getAllMeasurementSchemaByPrefix(prefixPath); - } // endregion // endregion diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java index 80c9084..ceb8f40 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java @@ -118,7 +118,6 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.BiFunction; -import java.util.function.Function; import java.util.stream.Collectors; import static org.apache.iotdb.db.conf.IoTDBConstant.MULTI_LEVEL_PATH_WILDCARD; @@ -1116,84 +1115,6 @@ public class MRocksDBManager implements IMetaManager { return index; } - private Map<String, byte[]> getKeyNumByPrefix( - PartialPath pathPattern, char nodeType, boolean isPrefixMatch) { - Map<String, byte[]> result = new ConcurrentHashMap<>(); - Set<String> seeds = new HashSet<>(); - - String seedPath; - - int nonWildcardAvailablePosition = - pathPattern.getFullPath().indexOf(ONE_LEVEL_PATH_WILDCARD) - 1; - if (nonWildcardAvailablePosition < 0) { - seedPath = RocksDBUtils.getLevelPath(pathPattern.getNodes(), pathPattern.getNodeLength() - 1); - } else { - seedPath = pathPattern.getFullPath().substring(0, nonWildcardAvailablePosition); - } - - seeds.add(new String(RocksDBUtils.toRocksDBKey(seedPath, nodeType))); - - scanAllKeysRecursively( - seeds, - 0, - s -> { - try { - byte[] value = readWriteHandler.get(null, s.getBytes()); - if (value != null && value.length > 0 && s.charAt(0) == nodeType) { - if (!isPrefixMatch || isMatched(pathPattern, s)) { - result.put(s, value); - return false; - } - } - } catch (RocksDBException e) { - return false; - } - return true; - }, - isPrefixMatch); - return result; - } - - // eg. pathPatter:root.a.b prefixedKey=sroot.2a.2bbb - private boolean isMatched(PartialPath pathPattern, String prefixedKey) { - // path = root.a.bbb - String path = RocksDBUtils.getPathByInnerName(prefixedKey); - if (path.length() <= pathPattern.getFullPath().length()) { - return true; - } else { - String fullPath = pathPattern.getFullPath() + RockDBConstants.PATH_SEPARATOR; - return path.startsWith(fullPath); - } - } - - private void scanAllKeysRecursively( - Set<String> seeds, int level, Function<String, Boolean> op, boolean isPrefixMatch) { - if (seeds == null || seeds.isEmpty()) { - return; - } - Set<String> children = ConcurrentHashMap.newKeySet(); - seeds - .parallelStream() - .forEach( - x -> { - if (op.apply(x)) { - if (isPrefixMatch) { - for (int i = level; i < MAX_PATH_DEPTH; i++) { - // x is not leaf node - String nextLevel = RocksDBUtils.getNextLevelOfPath(x, i); - children.addAll(readWriteHandler.getAllByPrefix(nextLevel)); - } - } else { - String nextLevel = RocksDBUtils.getNextLevelOfPath(x, level); - children.addAll(readWriteHandler.getAllByPrefix(nextLevel)); - } - } - }); - if (!children.isEmpty()) { - scanAllKeysRecursively(children, level + 1, op, isPrefixMatch); - } - } - @Override public int getAllTimeseriesCount(PartialPath pathPattern) throws MetadataException { return getAllTimeseriesCount(pathPattern, false); @@ -1550,21 +1471,6 @@ public class MRocksDBManager implements IMetaManager { return allPath; } - private Set<PartialPath> getMatchedPathWithNodeType( - boolean isPrefixMatch, PartialPath pathPattern, char nodeType) throws MetadataException { - Set<PartialPath> result = new HashSet<>(); - Map<String, byte[]> allMeasurement = getKeyNumByPrefix(pathPattern, nodeType, isPrefixMatch); - for (Entry<String, byte[]> entry : allMeasurement.entrySet()) { - try { - PartialPath path = new PartialPath(RocksDBUtils.getPathByInnerName(entry.getKey())); - result.add(path); - } catch (ClassCastException | IllegalPathException e) { - throw new MetadataException(e); - } - } - return result; - } - /** * Get all device paths and according storage group paths as ShowDevicesResult. * @@ -1767,10 +1673,10 @@ public class MRocksDBManager implements IMetaManager { RocksDBUtils.getLevelPath(fullPath.getNodes(), fullPath.getNodeLength() - 1); Holder<byte[]> holder = new Holder<>(); if (readWriteHandler.keyExistByType(levelKey, RocksDBMNodeType.MEASUREMENT, holder)) { - IMeasurementSchema schema = + MeasurementSchema measurementSchema = (MeasurementSchema) - RocksDBUtils.parseNodeValue(holder.getValue(), DATA_BLOCK_TYPE_SCHEMA); - return schema; + RocksDBUtils.parseNodeValue(holder.getValue(), RMNodeValueType.SCHEMA); + return measurementSchema; } else { throw new PathNotExistException(fullPath.getFullPath()); } @@ -1796,30 +1702,12 @@ public class MRocksDBManager implements IMetaManager { throw new PathNotExistException(e.getMessage()); } MeasurementSchema measurementSchema = - (MeasurementSchema) RocksDBUtils.parseNodeValue(entry.getValue(), FLAG_IS_SCHEMA); + (MeasurementSchema) RocksDBUtils.parseNodeValue(entry.getValue(), RMNodeValueType.SCHEMA); result.add(new MeasurementPath(pathName, measurementSchema)); } return result; } - @Override - public Map<PartialPath, IMeasurementSchema> getAllMeasurementSchemaByPrefix( - PartialPath prefixPath) throws MetadataException { - Map<PartialPath, IMeasurementSchema> result = new HashMap<>(); - Map<String, byte[]> allMeasurement = getKeyNumByPrefix(prefixPath, NODE_TYPE_MEASUREMENT, true); - for (Entry<String, byte[]> entry : allMeasurement.entrySet()) { - try { - MeasurementSchema schema = - (MeasurementSchema) RocksDBUtils.parseNodeValue(entry.getValue(), FLAG_IS_SCHEMA); - PartialPath path = new PartialPath(RocksDBUtils.getPathByInnerName(entry.getKey())); - result.put(path, schema); - } catch (ClassCastException e) { - throw new MetadataException(e); - } - } - return result; - } - /** * E.g., root.sg is storage group given [root, sg], return the MNode of root.sg given [root, sg, * device], return the MNode of root.sg Get storage group node by path. If storage group is not @@ -1842,13 +1730,9 @@ public class MRocksDBManager implements IMetaManager { String levelPath = RocksDBUtils.getLevelPath(nodes, i); Holder<byte[]> holder = new Holder<>(); if (readWriteHandler.keyExistByType(levelPath, RocksDBMNodeType.STORAGE_GROUP, holder)) { - Object ttl = RocksDBUtils.parseNodeValue(holder.getValue(), RockDBConstants.FLAG_SET_TTL); - if (ttl == null) { - ttl = config.getDefaultTTL(); - } node = new RStorageGroupMNode( - MetaUtils.getStorageGroupPathByLevel(path, i).getFullPath(), (Long) ttl); + MetaUtils.getStorageGroupPathByLevel(path, i).getFullPath(), holder.getValue()); break; } } @@ -1873,13 +1757,9 @@ public class MRocksDBManager implements IMetaManager { if (iterator.key()[0] != NODE_TYPE_SG) { break; } - Object ttl = RocksDBUtils.parseNodeValue(iterator.value(), FLAG_SET_TTL); - if (ttl == null) { - ttl = config.getDefaultTTL(); - } result.add( new RStorageGroupMNode( - RocksDBUtils.getPathByInnerName(new String(iterator.key())), (Long) ttl)); + RocksDBUtils.getPathByInnerName(new String(iterator.key())), iterator.value())); } return result; } diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RMNodeValueType.java b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RMNodeValueType.java new file mode 100644 index 0000000..2dee074 --- /dev/null +++ b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RMNodeValueType.java @@ -0,0 +1,18 @@ +package org.apache.iotdb.db.metadata.rocksdb; + +public enum RMNodeValueType { + TTL(RockDBConstants.FLAG_SET_TTL, RockDBConstants.DATA_BLOCK_TYPE_TTL), + SCHEMA(RockDBConstants.FLAG_HAS_SCHEMA, RockDBConstants.DATA_BLOCK_TYPE_SCHEMA), + ALIAS(RockDBConstants.FLAG_HAS_ALIAS, RockDBConstants.DATA_BLOCK_TYPE_ALIAS), + TAGS(RockDBConstants.FLAG_HAS_TAGS, RockDBConstants.DATA_BLOCK_TYPE_TAGS), + ATTRIBUTES(RockDBConstants.FLAG_HAS_ATTRIBUTES, RockDBConstants.DATA_BLOCK_TYPE_ATTRIBUTES), + ORIGIN_KEY(null, RockDBConstants.DATA_BLOCK_TYPE_ORIGIN_KEY); + + byte type; + Byte flag; + + RMNodeValueType(Byte flag, byte type) { + this.type = type; + this.flag = flag; + } +} diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RockDBConstants.java b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RockDBConstants.java index dfc4df2..be49fd6 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RockDBConstants.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RockDBConstants.java @@ -55,7 +55,7 @@ public class RockDBConstants { public static final byte DEFAULT_FLAG = 0x00; public static final byte FLAG_SET_TTL = 0x01; - public static final byte FLAG_IS_SCHEMA = 0x01 << 1; + public static final byte FLAG_HAS_SCHEMA = 0x01 << 1; public static final byte FLAG_HAS_ALIAS = 0x01 << 2; public static final byte FLAG_HAS_TAGS = 0x01 << 3; public static final byte FLAG_HAS_ATTRIBUTES = 0x01 << 4; diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBUtils.java b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBUtils.java index 38906c0..392ba35 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBUtils.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBUtils.java @@ -44,9 +44,8 @@ import java.util.Map; import java.util.stream.Collectors; import static org.apache.iotdb.db.conf.IoTDBConstant.*; -import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.*; -import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.FLAG_IS_ALIGNED; import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.PATH_SEPARATOR; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.*; public class RocksDBUtils { @@ -176,7 +175,7 @@ public class RocksDBUtils { } if (schema != null) { - flag = (byte) (flag | FLAG_IS_SCHEMA); + flag = (byte) (flag | FLAG_HAS_SCHEMA); } ReadWriteIOUtils.write(flag, outputStream); @@ -214,16 +213,18 @@ public class RocksDBUtils { return ReadWriteIOUtils.readBytes(buffer, len); } - public static int indexOfDataBlockType(byte[] data, byte type) { - if ((data[1] & FLAG_SET_TTL) == 0) { + public static int indexOfDataBlockType(byte[] data, RMNodeValueType valueType) { + if (valueType.flag != null && (data[1] & valueType.flag) == 0) { return -1; } int index = -1; boolean typeExist = false; + ByteBuffer byteBuffer = ByteBuffer.wrap(data); - // skip the version flag and node type flag + // skip the data version and filter byte ReadWriteIOUtils.readBytes(byteBuffer, 2); + while (byteBuffer.hasRemaining()) { byte blockType = ReadWriteIOUtils.readByte(byteBuffer); index = byteBuffer.position(); @@ -248,7 +249,7 @@ public class RocksDBUtils { break; } // got the data we need,don't need to read any more - if (type == blockType) { + if (valueType.type == blockType) { typeExist = true; break; } @@ -257,7 +258,7 @@ public class RocksDBUtils { } public static byte[] updateTTL(byte[] origin, long ttl) { - int index = indexOfDataBlockType(origin, DATA_BLOCK_TYPE_TTL); + int index = indexOfDataBlockType(origin, RMNodeValueType.TTL); if (index < 1) { byte[] ttlBlock = new byte[Long.BYTES + 1]; ttlBlock[0] = DATA_BLOCK_TYPE_TTL; @@ -277,44 +278,45 @@ public class RocksDBUtils { * parse value and return a specified type. if no data is required, null is returned. * * @param value value written in default table - * @param type the type of value to obtain + * @param valueType the type of value to obtain */ - public static Object parseNodeValue(byte[] value, byte type) { + public static Object parseNodeValue(byte[] value, RMNodeValueType valueType) { ByteBuffer byteBuffer = ByteBuffer.wrap(value); // skip the version flag and node type flag ReadWriteIOUtils.readByte(byteBuffer); // get block type - byte flag = ReadWriteIOUtils.readByte(byteBuffer); - + byte filter = ReadWriteIOUtils.readByte(byteBuffer); Object obj = null; + if (valueType.flag != null && (filter & valueType.flag) == 0) { + return obj; + } + // this means that the following data contains the information we need - if ((flag & type) > 0) { - while (byteBuffer.hasRemaining()) { - byte blockType = ReadWriteIOUtils.readByte(byteBuffer); - switch (blockType) { - case DATA_BLOCK_TYPE_TTL: - obj = ReadWriteIOUtils.readLong(byteBuffer); - break; - case DATA_BLOCK_TYPE_ALIAS: - obj = ReadWriteIOUtils.readString(byteBuffer); - break; - case DATA_BLOCK_TYPE_ORIGIN_KEY: - obj = readOriginKey(byteBuffer); - break; - case DATA_BLOCK_TYPE_SCHEMA: - obj = MeasurementSchema.deserializeFrom(byteBuffer); - break; - case DATA_BLOCK_TYPE_TAGS: - case DATA_BLOCK_TYPE_ATTRIBUTES: - obj = ReadWriteIOUtils.readMap(byteBuffer); - break; - default: - break; - } - // got the data we need,don't need to read any more - if (type == blockType) { + while (byteBuffer.hasRemaining()) { + byte blockType = ReadWriteIOUtils.readByte(byteBuffer); + switch (blockType) { + case DATA_BLOCK_TYPE_TTL: + obj = ReadWriteIOUtils.readLong(byteBuffer); break; - } + case DATA_BLOCK_TYPE_ALIAS: + obj = ReadWriteIOUtils.readString(byteBuffer); + break; + case DATA_BLOCK_TYPE_ORIGIN_KEY: + obj = readOriginKey(byteBuffer); + break; + case DATA_BLOCK_TYPE_SCHEMA: + obj = MeasurementSchema.deserializeFrom(byteBuffer); + break; + case DATA_BLOCK_TYPE_TAGS: + case DATA_BLOCK_TYPE_ATTRIBUTES: + obj = ReadWriteIOUtils.readMap(byteBuffer); + break; + default: + break; + } + // got the data we need,don't need to read any more + if (valueType.type == blockType) { + break; } } return obj; diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/mnode/RStorageGroupMNode.java b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/mnode/RStorageGroupMNode.java index da22224..198ca29 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/mnode/RStorageGroupMNode.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/mnode/RStorageGroupMNode.java @@ -21,7 +21,7 @@ package org.apache.iotdb.db.metadata.rocksdb.mnode; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.metadata.logfile.MLogWriter; import org.apache.iotdb.db.metadata.mnode.IStorageGroupMNode; -import org.apache.iotdb.db.metadata.rocksdb.RockDBConstants; +import org.apache.iotdb.db.metadata.rocksdb.RMNodeValueType; import org.apache.iotdb.db.metadata.rocksdb.RocksDBUtils; import java.io.IOException; @@ -42,7 +42,7 @@ public class RStorageGroupMNode extends RInternalMNode implements IStorageGroupM public RStorageGroupMNode(String fullPath, byte[] value) { super(fullPath); - Object ttl = RocksDBUtils.parseNodeValue(value, RockDBConstants.FLAG_SET_TTL); + Object ttl = RocksDBUtils.parseNodeValue(value, RMNodeValueType.TTL); if (ttl != null) { ttl = IoTDBDescriptor.getInstance().getConfig().getDefaultTTL(); }
