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 e6e7bc0e18ef9d9487c47bdde4a451dd158ddebd
Author: chengjianyun <[email protected]>
AuthorDate: Wed Mar 2 18:38:53 2022 +0800

    [rocksdb] complete metadata transfer from mmanage to mrockdb
---
 .../metadata/AcquireLockTimeoutException.java      |   7 +
 .../iotdb/db/metadata/rocksdb/MRocksDBManager.java |  51 ++--
 .../db/metadata/rocksdb/MetaDataTransfer.java      | 278 +++++++++++++++++----
 .../metadata/rocksdb/RocksDBReadWriteHandler.java  |   2 +-
 .../{MRocksDBTest.java => MRocksDBBenchmark.java}  |  30 +--
 ...TestEngine.java => RocksDBBenchmarkEngine.java} |  16 +-
 ...ksDBTestTask.java => RocksDBBenchmarkTask.java} |   4 +-
 .../db/metadata/rocksdb/RocksDBTestUtils.java      |   5 +-
 8 files changed, 284 insertions(+), 109 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/exception/metadata/AcquireLockTimeoutException.java
 
b/server/src/main/java/org/apache/iotdb/db/exception/metadata/AcquireLockTimeoutException.java
new file mode 100644
index 0000000..81c39d5
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/exception/metadata/AcquireLockTimeoutException.java
@@ -0,0 +1,7 @@
+package org.apache.iotdb.db.exception.metadata;
+
+public class AcquireLockTimeoutException extends MetadataException {
+  public AcquireLockTimeoutException(String msg) {
+    super(msg);
+  }
+}
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 ec2bd73..e408263 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
@@ -22,6 +22,7 @@ package org.apache.iotdb.db.metadata.rocksdb;
 import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBConstant;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.exception.metadata.AcquireLockTimeoutException;
 import org.apache.iotdb.db.exception.metadata.AliasAlreadyExistException;
 import org.apache.iotdb.db.exception.metadata.AlignedTimeseriesException;
 import org.apache.iotdb.db.exception.metadata.DataTypeMismatchException;
@@ -118,17 +119,8 @@ import java.util.stream.Collectors;
 
 import static org.apache.iotdb.db.conf.IoTDBConstant.MULTI_LEVEL_PATH_WILDCARD;
 import static org.apache.iotdb.db.conf.IoTDBConstant.ONE_LEVEL_PATH_WILDCARD;
-import static 
org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.DATA_BLOCK_TYPE_SCHEMA;
-import static 
org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.DEFAULT_ALIGNED_ENTITY_VALUE;
-import static 
org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.DEFAULT_NODE_VALUE;
-import static 
org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.FLAG_IS_ALIGNED;
-import static 
org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.FLAG_IS_SCHEMA;
-import static 
org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.FLAG_SET_TTL;
-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_SG;
-import static 
org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.TABLE_NAME_TAGS;
-import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.ZERO;
+import static org.apache.iotdb.db.metadata.rocksdb.RockDBConstants.*;
+import static org.apache.iotdb.db.metadata.rocksdb.RocksDBUtils.*;
 import static 
org.apache.iotdb.tsfile.common.constant.TsFileConstant.PATH_SEPARATOR;
 
 /**
@@ -386,7 +378,6 @@ public class MRocksDBManager implements IMetaManager {
     Holder<byte[]> holder = new Holder<>();
     Lock lock = locksPool.get(levelPath);
     if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
-      Thread.sleep(5);
       lockedLocks.push(lock);
       try {
         CheckKeyResult checkResult = 
readWriteHandler.keyExistByAllTypes(levelPath, holder);
@@ -402,12 +393,12 @@ public class MRocksDBManager implements IMetaManager {
           }
         } else {
           if (start == nodes.length) {
-            throw new PathAlreadyExistException("Measurement node already 
exists");
+            throw new PathAlreadyExistException(levelPath);
           }
 
           if (checkResult.getResult(RocksDBMNodeType.MEASUREMENT)
               || checkResult.getResult(RocksDBMNodeType.ALISA)) {
-            throw new PathAlreadyExistException("Path contains measurement 
node");
+            throw new PathAlreadyExistException(levelPath);
           }
 
           if (start == nodes.length - 1) {
@@ -440,7 +431,7 @@ public class MRocksDBManager implements IMetaManager {
       while (!lockedLocks.isEmpty()) {
         lockedLocks.pop().unlock();
       }
-      throw new MetadataException("acquire lock timeout: " + levelPath);
+      throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
     }
   }
 
@@ -481,7 +472,7 @@ public class MRocksDBManager implements IMetaManager {
           lock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
     } else {
       readWriteHandler.executeBatch(batch);
@@ -559,7 +550,7 @@ public class MRocksDBManager implements IMetaManager {
               throw new PathAlreadyExistException(lockKey);
             }
           } else {
-            throw new MetadataException("acquire lock timeout: " + lockKey);
+            throw new AcquireLockTimeoutException("acquire lock timeout: " + 
lockKey);
           }
         }
         readWriteHandler.executeBatch(batch);
@@ -633,7 +624,7 @@ public class MRocksDBManager implements IMetaManager {
       while (!lockedLocks.isEmpty()) {
         lockedLocks.pop().unlock();
       }
-      throw new MetadataException("acquire lock timeout: " + levelPath);
+      throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
     }
   }
 
@@ -679,7 +670,7 @@ public class MRocksDBManager implements IMetaManager {
             lock.unlock();
           }
         } else {
-          throw new MetadataException("acquire lock timeout, " + 
p.getFullPath());
+          throw new AcquireLockTimeoutException("acquire lock timeout, " + 
p.getFullPath());
         }
 
         // delete parent node if is empty
@@ -713,7 +704,7 @@ public class MRocksDBManager implements IMetaManager {
               curLock.unlock();
             }
           } else {
-            throw new MetadataException("acquire lock timeout, " + 
curNode.getFullPath());
+            throw new AcquireLockTimeoutException("acquire lock timeout, " + 
curNode.getFullPath());
           }
         }
         // TODO: trigger engine update
@@ -819,7 +810,7 @@ public class MRocksDBManager implements IMetaManager {
             lock.unlock();
           }
         } else {
-          throw new MetadataException("acquire lock timeout: " + levelKey);
+          throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelKey);
         }
       }
     } catch (RocksDBException | InterruptedException e) {
@@ -901,7 +892,7 @@ public class MRocksDBManager implements IMetaManager {
           lock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
     } catch (InterruptedException | RocksDBException e) {
       throw new MetadataException(e);
@@ -1817,7 +1808,7 @@ public class MRocksDBManager implements IMetaManager {
                 }
                 batch.put(newAliasKey, 
RocksDBUtils.buildAliasNodeValue(originKey));
               } else {
-                throw new MetadataException("acquire lock timeout: " + 
newAliasLevel);
+                throw new AcquireLockTimeoutException("acquire lock timeout: " 
+ newAliasLevel);
               }
 
               if (StringUtils.isNotEmpty(mNode.getAlias()) && 
!mNode.getAlias().equals(alias)) {
@@ -1836,7 +1827,7 @@ public class MRocksDBManager implements IMetaManager {
                   }
                   batch.delete(oldAliasKey);
                 } else {
-                  throw new MetadataException("acquire lock timeout: " + 
oldAliasLevel);
+                  throw new AcquireLockTimeoutException("acquire lock timeout: 
" + oldAliasLevel);
                 }
               }
               // TODO: need application lock
@@ -1877,7 +1868,7 @@ public class MRocksDBManager implements IMetaManager {
           rawKeyLock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
     } catch (RocksDBException | InterruptedException e) {
       throw new MetadataException(e);
@@ -1918,7 +1909,7 @@ public class MRocksDBManager implements IMetaManager {
           lock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
     } catch (RocksDBException | InterruptedException e) {
       throw new MetadataException(e);
@@ -1967,7 +1958,7 @@ public class MRocksDBManager implements IMetaManager {
           lock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
     } catch (RocksDBException | InterruptedException e) {
       throw new MetadataException(e);
@@ -2022,7 +2013,7 @@ public class MRocksDBManager implements IMetaManager {
           lock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
     } catch (RocksDBException | InterruptedException e) {
       throw new MetadataException(e);
@@ -2070,7 +2061,7 @@ public class MRocksDBManager implements IMetaManager {
           lock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
       if (!readWriteHandler.keyExist(key, holder)) {
         throw new PathNotExistException(path.getFullPath());
@@ -2125,7 +2116,7 @@ public class MRocksDBManager implements IMetaManager {
           lock.unlock();
         }
       } else {
-        throw new MetadataException("acquire lock timeout: " + levelPath);
+        throw new AcquireLockTimeoutException("acquire lock timeout: " + 
levelPath);
       }
     } catch (RocksDBException | InterruptedException e) {
       throw new MetadataException(e);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MetaDataTransfer.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MetaDataTransfer.java
index b99ef5b..84e02dd 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MetaDataTransfer.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MetaDataTransfer.java
@@ -21,15 +21,28 @@ package org.apache.iotdb.db.metadata.rocksdb;
 
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
+import org.apache.iotdb.db.exception.metadata.AcquireLockTimeoutException;
+import org.apache.iotdb.db.exception.metadata.AliasAlreadyExistException;
 import org.apache.iotdb.db.exception.metadata.MetadataException;
+import org.apache.iotdb.db.exception.metadata.PathAlreadyExistException;
+import org.apache.iotdb.db.exception.metadata.StorageGroupAlreadySetException;
 import org.apache.iotdb.db.metadata.MetadataConstant;
 import org.apache.iotdb.db.metadata.logfile.MLogReader;
+import org.apache.iotdb.db.metadata.logfile.MLogWriter;
 import org.apache.iotdb.db.metadata.mnode.IMeasurementMNode;
-import org.apache.iotdb.db.metadata.mnode.MeasurementMNode;
-import org.apache.iotdb.db.metadata.mnode.StorageGroupMNode;
+import org.apache.iotdb.db.metadata.mnode.IStorageGroupMNode;
+import org.apache.iotdb.db.metadata.mtree.MTree;
+import 
org.apache.iotdb.db.metadata.mtree.traverser.collector.MeasurementCollector;
+import org.apache.iotdb.db.metadata.path.PartialPath;
 import org.apache.iotdb.db.qp.physical.PhysicalPlan;
-import org.apache.iotdb.db.qp.physical.sys.MeasurementMNodePlan;
-import org.apache.iotdb.db.qp.physical.sys.StorageGroupMNodePlan;
+import org.apache.iotdb.db.qp.physical.sys.AutoCreateDeviceMNodePlan;
+import org.apache.iotdb.db.qp.physical.sys.ChangeAliasPlan;
+import org.apache.iotdb.db.qp.physical.sys.CreateAlignedTimeSeriesPlan;
+import org.apache.iotdb.db.qp.physical.sys.CreateTimeSeriesPlan;
+import org.apache.iotdb.db.qp.physical.sys.DeleteStorageGroupPlan;
+import org.apache.iotdb.db.qp.physical.sys.DeleteTimeSeriesPlan;
+import org.apache.iotdb.db.qp.physical.sys.SetStorageGroupPlan;
+import org.apache.iotdb.db.qp.physical.sys.SetTTLPlan;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -39,7 +52,6 @@ import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.atomic.AtomicInteger;
-import java.util.concurrent.atomic.AtomicLong;
 
 public class MetaDataTransfer {
 
@@ -47,31 +59,38 @@ public class MetaDataTransfer {
 
   private String mtreeSnapshotPath;
   private MRocksDBManager rocksDBManager;
+  private MLogWriter mLogWriter;
+  private String failedMLogPath =
+      IoTDBDescriptor.getInstance().getConfig().getSchemaDir()
+          + File.separator
+          + MetadataConstant.METADATA_LOG
+          + ".transfer_failed";
 
-  private AtomicInteger storageGroupToCreateCount = new AtomicInteger();
-  private AtomicLong timeSeriesToCreateCount = new AtomicLong();
-
-  private AtomicInteger createStorageGroupCount = new AtomicInteger();
-  private AtomicLong createTimeSeriesCount = new AtomicLong();
+  private AtomicInteger failedPlanCount = new AtomicInteger(0);
+  private List<PhysicalPlan> retryPlans = new ArrayList<>();
 
   MetaDataTransfer() throws MetadataException {
     rocksDBManager = new MRocksDBManager();
   }
 
-  public void init() throws IOException {
-    mtreeSnapshotPath =
-        IoTDBDescriptor.getInstance().getConfig().getSchemaDir()
-            + File.separator
-            + MetadataConstant.MTREE_SNAPSHOT;
+  public static void main(String[] args) {
+    try {
+      MetaDataTransfer transfer = new MetaDataTransfer();
+      transfer.doTransfer();
+    } catch (MetadataException | IOException e) {
+      e.printStackTrace();
+    }
+  }
 
-    File mtreeSnapshot = SystemFileFactory.INSTANCE.getFile(mtreeSnapshotPath);
-    long time = System.currentTimeMillis();
-    if (mtreeSnapshot.exists()) {
-      transferFromSnapshot(mtreeSnapshot);
-      logger.debug(
-          "spend {} ms to deserialize mtree from snapshot", 
System.currentTimeMillis() - time);
+  public void doTransfer() throws IOException {
+    File failedFile = new File(failedMLogPath);
+    if (failedFile.exists()) {
+      failedFile.delete();
     }
 
+    mLogWriter = new MLogWriter(failedMLogPath);
+    mLogWriter.setLogNum(0);
+
     String schemaDir = 
IoTDBDescriptor.getInstance().getConfig().getSchemaDir();
     File schemaFolder = SystemFileFactory.INSTANCE.getFile(schemaDir);
     if (!schemaFolder.exists()) {
@@ -82,6 +101,14 @@ public class MetaDataTransfer {
       }
     }
 
+    mtreeSnapshotPath = schemaDir + File.separator + 
MetadataConstant.MTREE_SNAPSHOT;
+    File mtreeSnapshot = SystemFileFactory.INSTANCE.getFile(mtreeSnapshotPath);
+    long time = System.currentTimeMillis();
+    if (mtreeSnapshot.exists()) {
+      transferFromSnapshot(mtreeSnapshot);
+      logger.info("spend {} ms to transfer data from snapshot", 
System.currentTimeMillis() - time);
+    }
+
     time = System.currentTimeMillis();
     String logFilePath = schemaDir + File.separator + 
MetadataConstant.METADATA_LOG;
     File logFile = SystemFileFactory.INSTANCE.getFile(logFilePath);
@@ -98,12 +125,17 @@ public class MetaDataTransfer {
       logger.info("no mlog.bin file find, skip transfer");
     }
 
-    logger.info("Do transfer success");
+    mLogWriter.close();
+
+    logger.info(
+        "do transfer complete with {} plan failed. Failed plan are persisted 
in mlog.bin.transfer_failed",
+        failedPlanCount.get());
   }
 
   private void transferFromMLog(MLogReader mLogReader) {
     int idx = 0;
     PhysicalPlan plan;
+    List<PhysicalPlan> nonCollisionCollections = new ArrayList<>();
     while (mLogReader.hasNext()) {
       try {
         plan = mLogReader.next();
@@ -116,45 +148,154 @@ public class MetaDataTransfer {
         continue;
       }
       try {
-        rocksDBManager.operation(plan);
+        switch (plan.getOperatorType()) {
+          case CREATE_TIMESERIES:
+          case CREATE_ALIGNED_TIMESERIES:
+          case AUTO_CREATE_DEVICE_MNODE:
+            nonCollisionCollections.add(plan);
+            if (nonCollisionCollections.size() > 100000) {
+              executeOperation(nonCollisionCollections, true);
+            }
+            break;
+          case DELETE_TIMESERIES:
+          case SET_STORAGE_GROUP:
+          case DELETE_STORAGE_GROUP:
+          case TTL:
+          case CHANGE_ALIAS:
+            executeOperation(nonCollisionCollections, true);
+            rocksDBManager.operation(plan);
+            break;
+          case CHANGE_TAG_OFFSET:
+          case CREATE_TEMPLATE:
+          case DROP_TEMPLATE:
+          case APPEND_TEMPLATE:
+          case PRUNE_TEMPLATE:
+          case SET_TEMPLATE:
+          case ACTIVATE_TEMPLATE:
+          case UNSET_TEMPLATE:
+          case CREATE_CONTINUOUS_QUERY:
+          case DROP_CONTINUOUS_QUERY:
+            logger.error("unsupported operations {}", plan.toString());
+            break;
+          default:
+            logger.error("Unrecognizable command {}", plan.getOperatorType());
+        }
       } catch (MetadataException | IOException e) {
         logger.error("Can not operate cmd {} for err:", 
plan.getOperatorType(), e);
+        if (!(e instanceof StorageGroupAlreadySetException)
+            && !(e instanceof PathAlreadyExistException)
+            && !(e instanceof AliasAlreadyExistException)) {
+          persistFailedLog(plan);
+        }
+      }
+    }
+    executeOperation(nonCollisionCollections, true);
+    if (retryPlans.size() > 0) {
+      executeOperation(retryPlans, false);
+    }
+  }
+
+  private void executeOperation(List<PhysicalPlan> plans, boolean 
needsToRetry) {
+    plans
+        .parallelStream()
+        .forEach(
+            x -> {
+              try {
+                rocksDBManager.operation(x);
+              } catch (IOException e) {
+                logger.error("failed to operate plan: {}", x.toString(), e);
+                retryPlans.add(x);
+              } catch (MetadataException e) {
+                logger.error("failed to operate plan: {}", x.toString(), e);
+                if (e instanceof AcquireLockTimeoutException && needsToRetry) {
+                  retryPlans.add(x);
+                } else {
+                  persistFailedLog(x);
+                }
+              } catch (Exception e) {
+                if (needsToRetry) {
+                  retryPlans.add(x);
+                } else {
+                  persistFailedLog(x);
+                }
+              }
+            });
+    logger.info("parallel executed {} operations", plans.size());
+    plans.clear();
+  }
+
+  private void persistFailedLog(PhysicalPlan plan) {
+    logger.info("persist won't retry and failed plan: {}", plan.toString());
+    failedPlanCount.incrementAndGet();
+    try {
+      switch (plan.getOperatorType()) {
+        case CREATE_TIMESERIES:
+          mLogWriter.createTimeseries((CreateTimeSeriesPlan) plan);
+          break;
+        case CREATE_ALIGNED_TIMESERIES:
+          mLogWriter.createAlignedTimeseries((CreateAlignedTimeSeriesPlan) 
plan);
+          break;
+        case AUTO_CREATE_DEVICE_MNODE:
+          mLogWriter.autoCreateDeviceMNode((AutoCreateDeviceMNodePlan) plan);
+          break;
+        case DELETE_TIMESERIES:
+          mLogWriter.deleteTimeseries((DeleteTimeSeriesPlan) plan);
+          break;
+        case SET_STORAGE_GROUP:
+          SetStorageGroupPlan setStorageGroupPlan = (SetStorageGroupPlan) plan;
+          mLogWriter.setStorageGroup(setStorageGroupPlan.getPath());
+          break;
+        case DELETE_STORAGE_GROUP:
+          DeleteStorageGroupPlan deletePlan = (DeleteStorageGroupPlan) plan;
+          for (PartialPath path : deletePlan.getPaths()) {
+            mLogWriter.deleteStorageGroup(path);
+          }
+          break;
+        case TTL:
+          SetTTLPlan ttlPlan = (SetTTLPlan) plan;
+          mLogWriter.setTTL(ttlPlan.getStorageGroup(), ttlPlan.getDataTTL());
+          break;
+        case CHANGE_ALIAS:
+          ChangeAliasPlan changeAliasPlan = (ChangeAliasPlan) plan;
+          mLogWriter.changeAlias(changeAliasPlan.getPath(), 
changeAliasPlan.getAlias());
+          break;
+        case CHANGE_TAG_OFFSET:
+        case CREATE_TEMPLATE:
+        case DROP_TEMPLATE:
+        case APPEND_TEMPLATE:
+        case PRUNE_TEMPLATE:
+        case SET_TEMPLATE:
+        case ACTIVATE_TEMPLATE:
+        case UNSET_TEMPLATE:
+        case CREATE_CONTINUOUS_QUERY:
+        case DROP_CONTINUOUS_QUERY:
+          throw new 
UnsupportedOperationException(plan.getOperatorType().toString());
+        default:
+          logger.error("Unrecognizable command {}", plan.getOperatorType());
       }
+    } catch (IOException e) {
+      logger.error(
+          "fatal error, exception when persist failed plan, metadata transfer 
should be failed", e);
     }
   }
 
   public void transferFromSnapshot(File mtreeSnapshot) {
     try (MLogReader mLogReader = new MLogReader(mtreeSnapshot)) {
       doTransferFromSnapshot(mLogReader);
-    } catch (IOException e) {
+    } catch (IOException | MetadataException e) {
       logger.warn("Failed to deserialize from {}. Use a new MTree.", 
mtreeSnapshot.getPath());
     }
   }
 
   @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity 
warning
-  private void doTransferFromSnapshot(MLogReader mLogReader) {
+  private void doTransferFromSnapshot(MLogReader mLogReader) throws 
IOException, MetadataException {
     long start = System.currentTimeMillis();
-    List<IMeasurementMNode> measurementList = new ArrayList<>();
-    List<StorageGroupMNode> sgList = new ArrayList<>();
-    while (mLogReader.hasNext()) {
-      PhysicalPlan plan = null;
-      try {
-        plan = mLogReader.next();
-        if (plan == null) {
-          continue;
-        }
-        if (plan instanceof StorageGroupMNodePlan) {
-          sgList.add(StorageGroupMNode.deserializeFrom((StorageGroupMNodePlan) 
plan));
-        } else if (plan instanceof MeasurementMNodePlan) {
-          
measurementList.add(MeasurementMNode.deserializeFrom((MeasurementMNodePlan) 
plan));
-        }
-      } catch (Exception e) {
-        logger.error(
-            "Can not operate cmd {} for err:", plan == null ? "" : 
plan.getOperatorType(), e);
-      }
-    }
+    MTree mTree = new MTree();
+    mTree.init();
+    List<IStorageGroupMNode> storageGroupNodes = 
mTree.getAllStorageGroupNodes();
 
-    sgList
+    AtomicInteger errorCount = new AtomicInteger(0);
+    storageGroupNodes
         .parallelStream()
         .forEach(
             sgNode -> {
@@ -164,22 +305,61 @@ public class MetaDataTransfer {
                   rocksDBManager.setTTL(sgNode.getPartialPath(), 
sgNode.getDataTTL());
                 }
               } catch (MetadataException | IOException e) {
-                logger.error("");
+                if (!(e instanceof StorageGroupAlreadySetException)
+                    && !(e instanceof PathAlreadyExistException)
+                    && !(e instanceof AliasAlreadyExistException)) {
+                  errorCount.incrementAndGet();
+                }
+                logger.error(
+                    "create storage group {} failed", 
sgNode.getPartialPath().getFullPath(), e);
               }
             });
 
-    measurementList
+    if (errorCount.get() > 0) {
+      logger.info("Fatal error. create some storage groups fail, terminate 
metadata transfer");
+      return;
+    }
+
+    List<IMeasurementMNode> measurementMNodes = new ArrayList<>();
+
+    MeasurementCollector collector =
+        new MeasurementCollector(
+            mTree.getNodeByPath(new PartialPath("root")), new 
PartialPath("root.**"), -1, -1) {
+          @Override
+          protected void collectMeasurement(IMeasurementMNode node) throws 
MetadataException {
+            measurementMNodes.add(node);
+          }
+        };
+    collector.traverse();
+
+    measurementMNodes
         .parallelStream()
         .forEach(
             mNode -> {
               try {
                 rocksDBManager.createTimeSeries(
                     mNode.getPartialPath(), mNode.getSchema(), 
mNode.getAlias(), null, null);
+              } catch (AcquireLockTimeoutException e) {
+                try {
+                  rocksDBManager.createTimeSeries(
+                      mNode.getPartialPath(), mNode.getSchema(), 
mNode.getAlias(), null, null);
+                } catch (MetadataException metadataException) {
+                  logger.error(
+                      "create timeseries {} failed in retry",
+                      mNode.getPartialPath().getFullPath(),
+                      e);
+                  errorCount.incrementAndGet();
+                }
               } catch (MetadataException e) {
-                logger.error("");
+                logger.error(
+                    "create timeseries {} failed", 
mNode.getPartialPath().getFullPath(), e);
+                errorCount.incrementAndGet();
               }
             });
 
-    logger.info("snapshot transfer complete after {}ms", 
System.currentTimeMillis() - start);
+    logger.info(
+        "metadata snapshot transfer complete after {}ms with {} errors",
+        System.currentTimeMillis() - start,
+        errorCount.get());
   }
 }
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 64359f9..9ab5f42 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
@@ -92,7 +92,7 @@ public class RocksDBReadWriteHandler {
     Options options = new Options();
     options.setCreateIfMissing(true);
     options.setAllowMmapReads(true);
-    options.setRowCache(new LRUCache(900000));
+    options.setRowCache(new LRUCache(9000000));
     options.setDbWriteBufferSize(16 * 1024 * 1024);
 
     org.rocksdb.Logger rocksDBLogger = new RockDBLogger(options, logger);
diff --git 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBTest.java 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBBenchmark.java
similarity index 86%
rename from 
server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBTest.java
rename to 
server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBBenchmark.java
index b255f31..8f054a5 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBBenchmark.java
@@ -13,20 +13,20 @@ import java.util.Collection;
 import java.util.List;
 import java.util.Set;
 
-public class MRocksDBTest {
+public class MRocksDBBenchmark {
   protected static IoTDBConfig config = 
IoTDBDescriptor.getInstance().getConfig();
   private MRocksDBManager rocksDBManager;
 
-  public MRocksDBTest(MRocksDBManager rocksDBManager) {
+  public MRocksDBBenchmark(MRocksDBManager rocksDBManager) {
     this.rocksDBManager = rocksDBManager;
   }
 
-  public List<RocksDBTestTask.BenchmarkResult> benchmarkResults = new 
ArrayList<>();
+  public List<RocksDBBenchmarkTask.BenchmarkResult> benchmarkResults = new 
ArrayList<>();
 
   public void testStorageGroupCreation(List<SetStorageGroupPlan> 
storageGroups) {
-    RocksDBTestTask<SetStorageGroupPlan> task =
-        new RocksDBTestTask<>(storageGroups, 
RocksDBTestUtils.WRITE_CLIENT_NUM, 100);
-    RocksDBTestTask.BenchmarkResult result =
+    RocksDBBenchmarkTask<SetStorageGroupPlan> task =
+        new RocksDBBenchmarkTask<>(storageGroups, 
RocksDBTestUtils.WRITE_CLIENT_NUM, 100);
+    RocksDBBenchmarkTask.BenchmarkResult result =
         task.runWork(
             setStorageGroupPlan -> {
               try {
@@ -42,12 +42,12 @@ public class MRocksDBTest {
 
   public void testTimeSeriesCreation(List<List<CreateTimeSeriesPlan>> 
timeSeriesSet)
       throws IOException {
-    RocksDBTestTask<List<CreateTimeSeriesPlan>> task =
-        new RocksDBTestTask<>(timeSeriesSet, 
RocksDBTestUtils.WRITE_CLIENT_NUM, 100);
-    RocksDBTestTask.BenchmarkResult result =
+    RocksDBBenchmarkTask<List<CreateTimeSeriesPlan>> task =
+        new RocksDBBenchmarkTask<>(timeSeriesSet, 
RocksDBTestUtils.WRITE_CLIENT_NUM, 100);
+    RocksDBBenchmarkTask.BenchmarkResult result =
         task.runBatchWork(
             createTimeSeriesPlans -> {
-              RocksDBTestTask.TaskResult taskResult = new 
RocksDBTestTask.TaskResult();
+              RocksDBBenchmarkTask.TaskResult taskResult = new 
RocksDBBenchmarkTask.TaskResult();
               createTimeSeriesPlans.stream()
                   .forEach(
                       ts -> {
@@ -92,9 +92,9 @@ public class MRocksDBTest {
   //  }
 
   public void testNodeChildrenQuery(Collection<String> queryTsSet) {
-    RocksDBTestTask<String> task =
-        new RocksDBTestTask<>(queryTsSet, RocksDBTestUtils.WRITE_CLIENT_NUM, 
10000);
-    RocksDBTestTask.BenchmarkResult result =
+    RocksDBBenchmarkTask<String> task =
+        new RocksDBBenchmarkTask<>(queryTsSet, 
RocksDBTestUtils.WRITE_CLIENT_NUM, 10000);
+    RocksDBBenchmarkTask.BenchmarkResult result =
         task.runWork(
             s -> {
               try {
@@ -121,8 +121,8 @@ public class MRocksDBTest {
     List<PartialPath> level4 = rocksDBManager.getNodesListInGivenLevel(null, 
4);
     List<PartialPath> level5 = rocksDBManager.getNodesListInGivenLevel(null, 
5);
     long totalCount = level1.size() + level2.size() + level3.size() + 
level4.size() + level5.size();
-    RocksDBTestTask.BenchmarkResult result =
-        new RocksDBTestTask.BenchmarkResult(
+    RocksDBBenchmarkTask.BenchmarkResult result =
+        new RocksDBBenchmarkTask.BenchmarkResult(
             "levelScan", totalCount, 0, System.currentTimeMillis() - start);
     benchmarkResults.add(result);
   }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestEngine.java
 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBBenchmarkEngine.java
similarity index 89%
rename from 
server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestEngine.java
rename to 
server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBBenchmarkEngine.java
index ba124c1..f11ed07 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestEngine.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBBenchmarkEngine.java
@@ -8,7 +8,6 @@ import org.apache.iotdb.db.exception.metadata.MetadataException;
 import org.apache.iotdb.db.metadata.MManager;
 import org.apache.iotdb.db.metadata.MetadataConstant;
 import org.apache.iotdb.db.metadata.logfile.MLogReader;
-import org.apache.iotdb.db.metadata.logfile.MLogWriter;
 import org.apache.iotdb.db.metadata.path.PartialPath;
 import org.apache.iotdb.db.qp.physical.PhysicalPlan;
 import org.apache.iotdb.db.qp.physical.sys.CreateTimeSeriesPlan;
@@ -30,20 +29,19 @@ import java.util.Set;
 
 import static 
org.apache.iotdb.db.metadata.rocksdb.RocksDBReadWriteHandler.ROCKSDB_PATH;
 
-public class RocksDBTestEngine {
+public class RocksDBBenchmarkEngine {
   private static final Logger logger = LoggerFactory.getLogger(MManager.class);
   private static final IoTDBConfig config = 
IoTDBDescriptor.getInstance().getConfig();
 
   private static final int BIN_CAPACITY = 100 * 1000;
 
-  private MLogWriter logWriter;
   private File logFile;
   public static List<List<CreateTimeSeriesPlan>> timeSeriesSet = new 
ArrayList<>();
   public static Set<String> measurementPathSet = new HashSet<>();
   public static Set<String> innerPathSet = new HashSet<>();
   public static List<SetStorageGroupPlan> storageGroups = new ArrayList<>();
 
-  public RocksDBTestEngine() {
+  public RocksDBBenchmarkEngine() {
     String schemaDir = config.getSchemaDir();
     String logFilePath = schemaDir + File.separator + 
MetadataConstant.METADATA_LOG;
     logFile = SystemFileFactory.INSTANCE.getFile(logFilePath);
@@ -58,10 +56,10 @@ public class RocksDBTestEngine {
           storageGroups, timeSeriesSet, measurementPathSet, innerPathSet);
       /** rocksdb benchmark * */
       MRocksDBManager rocksDBManager = new MRocksDBManager();
-      MRocksDBTest mRocksDBTest = new MRocksDBTest(rocksDBManager);
-      mRocksDBTest.testStorageGroupCreation(storageGroups);
-      mRocksDBTest.testTimeSeriesCreation(timeSeriesSet);
-      RocksDBTestUtils.printReport(mRocksDBTest.benchmarkResults, "rocksDB");
+      MRocksDBBenchmark mRocksDBBenchmark = new 
MRocksDBBenchmark(rocksDBManager);
+      mRocksDBBenchmark.testStorageGroupCreation(storageGroups);
+      mRocksDBBenchmark.testTimeSeriesCreation(timeSeriesSet);
+      RocksDBTestUtils.printReport(mRocksDBBenchmark.benchmarkResults, 
"rocksDB");
       RocksDBTestUtils.printMemInfo("Benchmark finished");
     } catch (IOException | MetadataException e) {
       logger.error("Error happened when run benchmark", e);
@@ -70,8 +68,6 @@ public class RocksDBTestEngine {
 
   public void prepareBenchmark() throws IOException {
     long time = System.currentTimeMillis();
-    logWriter = new MLogWriter(config.getSchemaDir(), 
MetadataConstant.METADATA_LOG + ".temp");
-    logWriter.setLogNum(0);
     if (!logFile.exists()) {
       throw new FileNotFoundException("we need a mlog.bin to init the 
benchmark test");
     }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestTask.java
 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBBenchmarkTask.java
similarity index 95%
rename from 
server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestTask.java
rename to 
server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBBenchmarkTask.java
index 8c1cfbf..a628783 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestTask.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBBenchmarkTask.java
@@ -9,12 +9,12 @@ import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.Function;
 
-public class RocksDBTestTask<T> {
+public class RocksDBBenchmarkTask<T> {
   private Collection<T> dataSet;
   private int workCount;
   private int timeoutInMin;
 
-  RocksDBTestTask(Collection<T> dataSet, int workCount, int timeoutInMin) {
+  RocksDBBenchmarkTask(Collection<T> dataSet, int workCount, int timeoutInMin) 
{
     this.dataSet = dataSet;
     this.workCount = workCount;
     this.timeoutInMin = timeoutInMin;
diff --git 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestUtils.java
 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestUtils.java
index 65609c4..4c75ae7 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestUtils.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/metadata/rocksdb/RocksDBTestUtils.java
@@ -37,13 +37,14 @@ public class RocksDBTestUtils {
         stageInfo, Runtime.getRuntime().totalMemory() - 
Runtime.getRuntime().freeMemory());
   }
 
-  public static void printReport(List<RocksDBTestTask.BenchmarkResult> 
results, String category) {
+  public static void printReport(
+      List<RocksDBBenchmarkTask.BenchmarkResult> results, String category) {
     System.out.println(
         String.format(
             "\n\n#################################%s benchmark 
statistics#################################",
             category));
     System.out.println(String.format("%25s %15s %10s %15s", "", "success", 
"fail", "cost-in-ms"));
-    for (RocksDBTestTask.BenchmarkResult result : results) {
+    for (RocksDBBenchmarkTask.BenchmarkResult result : results) {
       System.out.println(
           String.format(
               "%25s %15d %10d %15d",

Reply via email to