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 83da91c220951bd4558264334f1608c17d315564 Author: lisijia <[email protected]> AuthorDate: Wed Mar 16 17:58:30 2022 +0800 fix bug of the incorrect number by counting nodes --- .../iotdb/db/metadata/rocksdb/MRocksDBManager.java | 57 ++++++++++++++-------- .../metadata/rocksdb/RocksDBReadWriteHandler.java | 18 +++++-- 2 files changed, 49 insertions(+), 26 deletions(-) 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 95f23e4..459e924 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,6 +118,7 @@ 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; @@ -1174,14 +1175,22 @@ public class MRocksDBManager implements IMetaManager { } String innerNameByLevel = RocksDBUtils.getLevelPath(pathPattern.getNodes(), pathPattern.getNodeLength() - 1, level); - for (RocksDBMNodeType type : RocksDBMNodeType.values()) { - String getKeyByInnerNameLevel = type.value + innerNameByLevel; - int queryResult = readWriteHandler.getKeyByPrefix(getKeyByInnerNameLevel).size(); - if (queryResult != 0) { - return queryResult; - } - } - return 0; + AtomicInteger atomicInteger = new AtomicInteger(0); + Function<String, Boolean> function = + s -> { + atomicInteger.incrementAndGet(); + return true; + }; + Arrays.stream(ALL_NODE_TYPE_ARRAY) + .parallel() + .forEach( + x -> { + String getKeyByInnerNameLevel = + x + innerNameByLevel + RockDBConstants.PATH_SEPARATOR + level; + readWriteHandler.getKeyByPrefix(getKeyByInnerNameLevel, function); + }); + + return atomicInteger.get(); } /** @@ -1282,13 +1291,17 @@ public class MRocksDBManager implements IMetaManager { pathPattern.getNodeLength()) + RockDBConstants.PATH_SEPARATOR + pathPattern.getNodeLength(); + Function<String, Boolean> function = + s -> { + result.add(RocksDBUtils.getPathByInnerName(s)); + return true; + }; + Arrays.stream(ALL_NODE_TYPE_ARRAY) .parallel() .forEach( x -> { - for (String string : readWriteHandler.getKeyByPrefix(x + innerNameByLevel)) { - result.add(RocksDBUtils.getPathByInnerName(string)); - } + readWriteHandler.getKeyByPrefix(x + innerNameByLevel, function); }); return result; } @@ -1420,16 +1433,18 @@ public class MRocksDBManager implements IMetaManager { /** Get all storage group paths */ @Override public List<PartialPath> getAllStorageGroupPaths() { - List<PartialPath> allStorageGroupPath = new ArrayList<>(); - Set<String> allStorageGroupInnerName = - readWriteHandler.getKeyByPrefix(String.valueOf(NODE_TYPE_SG)); - for (String str : allStorageGroupInnerName) { - try { - allStorageGroupPath.add(new PartialPath(RocksDBUtils.getPathByInnerName(str))); - } catch (IllegalPathException e) { - throw new RuntimeException(e); - } - } + List<PartialPath> allStorageGroupPath = Collections.synchronizedList(new ArrayList<>()); + Function<String, Boolean> function = + s -> { + try { + allStorageGroupPath.add(new PartialPath(RocksDBUtils.getPathByInnerName(s))); + } catch (IllegalPathException e) { + logger.error(e.getMessage()); + return false; + } + return true; + }; + readWriteHandler.getKeyByPrefix(String.valueOf(NODE_TYPE_SG), function); return allStorageGroupPath; } diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBReadWriteHandler.java b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBReadWriteHandler.java index cac86e6..d4a3f05 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBReadWriteHandler.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBReadWriteHandler.java @@ -63,7 +63,17 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiConsumer; import java.util.function.Function; -import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.*; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.DATA_BLOCK_TYPE_ORIGIN_KEY; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.DATA_BLOCK_TYPE_SCHEMA; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.DATA_VERSION; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.DEFAULT_FLAG; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.NODE_TYPE_ENTITY; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.NODE_TYPE_MEASUREMENT; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.NODE_TYPE_ROOT; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.PATH_SEPARATOR; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.ROOT; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.TABLE_NAME_TAGS; +import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.ZERO; public class RocksDBReadWriteHandler { @@ -414,17 +424,15 @@ public class RocksDBReadWriteHandler { return rocksDB.newIterator(columnFamilyHandle); } - public Set<String> getKeyByPrefix(String innerName) { + public void getKeyByPrefix(String innerName, Function<String, Boolean> function) { RocksIterator iterator = rocksDB.newIterator(); - Set<String> result = new HashSet<>(); for (iterator.seek(innerName.getBytes()); iterator.isValid(); iterator.next()) { String keyStr = new String(iterator.key()); if (!keyStr.startsWith(innerName)) { break; } - result.add(keyStr); + function.apply(keyStr); } - return result; } public Map<byte[], byte[]> getKeyValueByPrefix(String innerName) {
