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 06fe60e147bee4666fbc67a56d9be1f5de2def3a
Author: chengjianyun <[email protected]>
AuthorDate: Fri Mar 4 15:43:56 2022 +0800

    [rocksdb] complete data transfer
---
 .../resources/conf/iotdb-engine.properties         |   8 +
 .../assembly/resources/tools/rocksdb-transfer.bat  | 126 +++++++++
 .../assembly/resources/tools/rocksdb-transfer.sh   |  82 ++++++
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |  11 +
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |   5 +
 .../iotdb/db/metadata/MetadataManagerType.java     |  14 +
 .../iotdb/db/metadata/rocksdb/MRocksDBManager.java |   1 -
 .../db/metadata/rocksdb/MetaDataTransfer.java      | 283 +++++++++++++--------
 .../java/org/apache/iotdb/db/service/IoTDB.java    |  10 +-
 9 files changed, 426 insertions(+), 114 deletions(-)

diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties 
b/server/src/assembly/resources/conf/iotdb-engine.properties
index 7150de3..fe6d101 100644
--- a/server/src/assembly/resources/conf/iotdb-engine.properties
+++ b/server/src/assembly/resources/conf/iotdb-engine.properties
@@ -515,6 +515,14 @@ timestamp_precision=ms
 # enable_id_table_log_file=false
 
 ####################
+### Metadata Configuration
+####################
+
+# Which metadata manager to be used, right now MEMORY_MANAGER and 
ROCKSDB_MANAGER are supported. Default MEMORY_MANAGER.
+# Datatype: string
+# meta_data_manager=MEMORY_MANAGER
+
+####################
 ### Metadata Cache Configuration
 ####################
 
diff --git a/server/src/assembly/resources/tools/rocksdb-transfer.bat 
b/server/src/assembly/resources/tools/rocksdb-transfer.bat
new file mode 100644
index 0000000..8fa5433
--- /dev/null
+++ b/server/src/assembly/resources/tools/rocksdb-transfer.bat
@@ -0,0 +1,126 @@
+@REM
+@REM Licensed to the Apache Software Foundation (ASF) under one
+@REM or more contributor license agreements.  See the NOTICE file
+@REM distributed with this work for additional information
+@REM regarding copyright ownership.  The ASF licenses this file
+@REM to you under the Apache License, Version 2.0 (the
+@REM "License"); you may not use this file except in compliance
+@REM with the License.  You may obtain a copy of the License at
+@REM
+@REM     http://www.apache.org/licenses/LICENSE-2.0
+@REM
+@REM Unless required by applicable law or agreed to in writing,
+@REM software distributed under the License is distributed on an
+@REM "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+@REM KIND, either express or implied.  See the License for the
+@REM specific language governing permissions and limitations
+@REM under the License.
+@REM
+
+@echo off
+echo ````````````````````````
+echo Starting IoTDB
+echo ````````````````````````
+
+
+set PATH="%JAVA_HOME%\bin\";%PATH%
+set "FULL_VERSION="
+set "MAJOR_VERSION="
+set "MINOR_VERSION="
+
+
+for /f tokens^=2-5^ delims^=.-_+^" %%j in ('java -fullversion 2^>^&1') do (
+       set "FULL_VERSION=%%j-%%k-%%l-%%m"
+       IF "%%j" == "1" (
+           set "MAJOR_VERSION=%%k"
+           set "MINOR_VERSION=%%l"
+       ) else (
+           set "MAJOR_VERSION=%%j"
+           set "MINOR_VERSION=%%k"
+       )
+)
+
+set JAVA_VERSION=%MAJOR_VERSION%
+
+@REM we do not check jdk that version less than 1.8 because they are too 
stale...
+IF "%JAVA_VERSION%" == "6" (
+       echo IoTDB only supports jdk >= 8, please check your java version.
+       goto finally
+)
+IF "%JAVA_VERSION%" == "7" (
+       echo IoTDB only supports jdk >= 8, please check your java version.
+       goto finally
+)
+
+if "%OS%" == "Windows_NT" setlocal
+
+pushd %~dp0..
+if NOT DEFINED IOTDB_HOME set IOTDB_HOME=%cd%
+popd
+
+set IOTDB_CONF=%IOTDB_HOME%\conf
+set IOTDB_LOGS=%IOTDB_HOME%\logs
+
+@setlocal ENABLEDELAYEDEXPANSION ENABLEEXTENSIONS
+set is_conf_path=false
+for %%i in (%*) do (
+       IF "%%i" == "-c" (
+               set is_conf_path=true
+       ) ELSE IF "!is_conf_path!" == "true" (
+               set is_conf_path=false
+               set IOTDB_CONF=%%i
+       ) ELSE (
+               set CONF_PARAMS=!CONF_PARAMS! %%i
+       )
+)
+
+IF EXIST "%IOTDB_CONF%\iotdb-env.bat" (
+    CALL "%IOTDB_CONF%\iotdb-env.bat" %1
+    ) ELSE (
+    echo "can't find %IOTDB_CONF%\iotdb-env.bat"
+    )
+
+if NOT DEFINED MAIN_CLASS set 
MAIN_CLASS=org.apache.iotdb.db.metadata.rocksdb.MetaDataTransfer
+if NOT DEFINED JAVA_HOME goto :err
+
+@REM 
-----------------------------------------------------------------------------
+@REM JVM Opts we'll use in legacy run or installation
+set JAVA_OPTS=-ea^
+ -Dlogback.configurationFile="%IOTDB_CONF%\logback.xml"^
+ -DIOTDB_HOME="%IOTDB_HOME%"^
+ -DTSFILE_HOME="%IOTDB_HOME%"^
+ -DTSFILE_CONF="%IOTDB_CONF%"^
+ -DIOTDB_CONF="%IOTDB_CONF%"^
+ -Dsun.jnu.encoding=UTF-8^
+ -Dfile.encoding=UTF-8
+
+@REM ***** CLASSPATH library setting *****
+@REM Ensure that any user defined CLASSPATH variables are not used on startup
+set CLASSPATH="%IOTDB_HOME%\lib\*"
+set CLASSPATH=%CLASSPATH%;iotdb.IoTDB
+goto okClasspath
+
+:append
+set CLASSPATH=%CLASSPATH%;%1
+
+goto :eof
+
+@REM 
-----------------------------------------------------------------------------
+:okClasspath
+
+rem echo CLASSPATH: %CLASSPATH%
+
+"%JAVA_HOME%\bin\java" %JAVA_OPTS% %IOTDB_HEAP_OPTS% -cp %CLASSPATH% 
%IOTDB_JMX_OPTS% %MAIN_CLASS% %CONF_PARAMS%
+goto finally
+
+:err
+echo JAVA_HOME environment variable must be set!
+pause
+
+
+@REM 
-----------------------------------------------------------------------------
+:finally
+
+pause
+
+ENDLOCAL
diff --git a/server/src/assembly/resources/tools/rocksdb-transfer.sh 
b/server/src/assembly/resources/tools/rocksdb-transfer.sh
new file mode 100644
index 0000000..cb9040c
--- /dev/null
+++ b/server/src/assembly/resources/tools/rocksdb-transfer.sh
@@ -0,0 +1,82 @@
+#!/bin/bash
+#
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#     http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+#
+
+
+echo ---------------------
+echo Starting IoTDB
+echo ---------------------
+
+if [ -z "${IOTDB_HOME}" ]; then
+  export IOTDB_HOME="`dirname "$0"`/.."
+fi
+
+IOTDB_CONF=${IOTDB_HOME}/conf
+# IOTDB_LOGS=${IOTDB_HOME}/logs
+
+is_conf_path=false
+for arg do
+  shift
+  if [ "$arg" == "-c" ]; then
+    is_conf_path=true
+    continue
+  fi
+  if [ $is_conf_path == true ]; then
+    IOTDB_CONF=$arg
+    is_conf_path=false
+    continue
+  fi
+  set -- "$@" "$arg"
+done
+
+CONF_PARAMS=$*
+
+if [ -f "$IOTDB_CONF/iotdb-env.sh" ]; then
+    if [ "$#" -ge "1" -a "$1" == "printgc" ]; then
+      . "$IOTDB_CONF/iotdb-env.sh" "printgc"
+    else
+        . "$IOTDB_CONF/iotdb-env.sh"
+    fi
+else
+    echo "can't find $IOTDB_CONF/iotdb-env.sh"
+fi
+
+CLASSPATH=""
+for f in ${IOTDB_HOME}/lib/*.jar; do
+  CLASSPATH=${CLASSPATH}":"$f
+done
+classname=org.apache.iotdb.db.metadata.rocksdb.MetaDataTransfer
+
+launch_service()
+{
+       class="$1"
+       iotdb_parms="-Dlogback.configurationFile=${IOTDB_CONF}/logback.xml"
+       iotdb_parms="$iotdb_parms -DIOTDB_HOME=${IOTDB_HOME}"
+       iotdb_parms="$iotdb_parms -DTSFILE_HOME=${IOTDB_HOME}"
+       iotdb_parms="$iotdb_parms -DIOTDB_CONF=${IOTDB_CONF}"
+       iotdb_parms="$iotdb_parms -DTSFILE_CONF=${IOTDB_CONF}"
+       iotdb_parms="$iotdb_parms -Dname=iotdb\.IoTDB"
+       exec "$JAVA" $iotdb_parms $IOTDB_JMX_OPTS -cp "$CLASSPATH" "$class" 
$CONF_PARAMS
+       return $?
+}
+
+# Start up the service
+launch_service "$classname"
+
+exit $?
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java 
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
index b6b23a8..1ee8a1d 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
@@ -25,6 +25,7 @@ import 
org.apache.iotdb.db.engine.compaction.inner.InnerCompactionStrategy;
 import org.apache.iotdb.db.engine.storagegroup.timeindex.TimeIndexLevel;
 import org.apache.iotdb.db.exception.LoadConfigurationException;
 import org.apache.iotdb.db.metadata.MManager;
+import org.apache.iotdb.db.metadata.MetadataManagerType;
 import org.apache.iotdb.db.service.thrift.impl.InfluxDBServiceImpl;
 import org.apache.iotdb.db.service.thrift.impl.TSServiceImpl;
 import org.apache.iotdb.rpc.RpcTransportFactory;
@@ -804,6 +805,8 @@ public class IoTDBConfig {
   /** Encryption provided class parameter */
   private String encryptDecryptProviderParameter;
 
+  private MetadataManagerType metadataManagerType = 
MetadataManagerType.MEMORY_MANAGER;
+
   public IoTDBConfig() {
     // empty constructor
   }
@@ -921,6 +924,14 @@ public class IoTDBConfig {
     this.timeIndexLevel = TimeIndexLevel.valueOf(timeIndexLevel);
   }
 
+  public void setMetadataManagerType(String type) {
+    metadataManagerType = MetadataManagerType.of(type);
+  }
+
+  public MetadataManagerType getMetadataManagerType() {
+    return metadataManagerType;
+  }
+
   void updatePath() {
     formulateFolders();
     confirmMultiDirStrategy();
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java 
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index 0d5cb9e..5751a45 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -1413,6 +1413,11 @@ public class IoTDBDescriptor {
             conf.getTimestampPrecision()));
   }
 
+  private void loadMetadataConfig(Properties properties) {
+    conf.setMetadataManagerType(
+        properties.getProperty("meta_data_manager", 
conf.getMetadataManagerType().name()));
+  }
+
   /** Get default encode algorithm by data type */
   public TSEncoding getDefaultEncodingByType(TSDataType dataType) {
     switch (dataType) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/MetadataManagerType.java 
b/server/src/main/java/org/apache/iotdb/db/metadata/MetadataManagerType.java
new file mode 100644
index 0000000..402e2f4
--- /dev/null
+++ b/server/src/main/java/org/apache/iotdb/db/metadata/MetadataManagerType.java
@@ -0,0 +1,14 @@
+package org.apache.iotdb.db.metadata;
+
+public enum MetadataManagerType {
+  MEMORY_MANAGER,
+  ROCKSDB_MANAGER;
+
+  public static MetadataManagerType of(String value) {
+    try {
+      return Enum.valueOf(MetadataManagerType.class, value);
+    } catch (Exception e) {
+      return MEMORY_MANAGER;
+    }
+  }
+}
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 e408263..9ca93ad 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
@@ -1021,7 +1021,6 @@ public class MRocksDBManager implements IMetaManager {
                             }
                             if (RocksDBUtils.suffixMatch(iterator.key(), 
suffixToMatch)) {
                               if (lastIteration) {
-                                System.out.println("matched key: " + new 
String(iterator.key()));
                                 consumer.accept(
                                     RocksDBUtils.getPathByInnerName(new 
String(iterator.key())));
                               } else {
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 84e02dd..c3c8906 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
@@ -29,6 +29,7 @@ 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.IMNode;
 import org.apache.iotdb.db.metadata.mnode.IMeasurementMNode;
 import org.apache.iotdb.db.metadata.mnode.IStorageGroupMNode;
 import org.apache.iotdb.db.metadata.mtree.MTree;
@@ -51,12 +52,21 @@ import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
+import java.util.Queue;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ForkJoinPool;
 import java.util.concurrent.atomic.AtomicInteger;
 
 public class MetaDataTransfer {
 
   private static final Logger logger = 
LoggerFactory.getLogger(MetaDataTransfer.class);
 
+  private static int DEFAULT_TRANSFER_THREAD_POOL_SIZE = 200;
+  private static int DEFAULT_TRANSFER_PLANS_BUFFER_SIZE = 100_000;
+
+  private ForkJoinPool forkJoinPool = new 
ForkJoinPool(DEFAULT_TRANSFER_THREAD_POOL_SIZE);
+
   private String mtreeSnapshotPath;
   private MRocksDBManager rocksDBManager;
   private MLogWriter mLogWriter;
@@ -77,12 +87,12 @@ public class MetaDataTransfer {
     try {
       MetaDataTransfer transfer = new MetaDataTransfer();
       transfer.doTransfer();
-    } catch (MetadataException | IOException e) {
+    } catch (MetadataException | IOException | ExecutionException | 
InterruptedException e) {
       e.printStackTrace();
     }
   }
 
-  public void doTransfer() throws IOException {
+  public void doTransfer() throws IOException, ExecutionException, 
InterruptedException {
     File failedFile = new File(failedMLogPath);
     if (failedFile.exists()) {
       failedFile.delete();
@@ -103,36 +113,34 @@ 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);
+      try {
+        doTransferFromSnapshot();
+      } catch (MetadataException e) {
+        logger.error("Fatal error, terminate data transfer!!!", e);
+      }
     }
 
-    time = System.currentTimeMillis();
     String logFilePath = schemaDir + File.separator + 
MetadataConstant.METADATA_LOG;
     File logFile = SystemFileFactory.INSTANCE.getFile(logFilePath);
     // init the metadata from the operation log
     if (logFile.exists()) {
       try (MLogReader mLogReader = new MLogReader(schemaDir, 
MetadataConstant.METADATA_LOG); ) {
         transferFromMLog(mLogReader);
-        logger.info(
-            "spend {} ms to deserialize mtree from mlog.bin", 
System.currentTimeMillis() - time);
       } catch (Exception e) {
         throw new IOException("Failed to parser mlog.bin for err:" + e);
       }
     } else {
-      logger.info("no mlog.bin file find, skip transfer");
+      logger.info("No mlog.bin file find, skip data transfer");
     }
-
     mLogWriter.close();
 
-    logger.info(
-        "do transfer complete with {} plan failed. Failed plan are persisted 
in mlog.bin.transfer_failed",
-        failedPlanCount.get());
+    logger.info("Transfer metadata from MManager to MRocksDBManager 
complete!");
   }
 
-  private void transferFromMLog(MLogReader mLogReader) {
+  private void transferFromMLog(MLogReader mLogReader)
+      throws IOException, MetadataException, ExecutionException, 
InterruptedException {
+    long time = System.currentTimeMillis();
     int idx = 0;
     PhysicalPlan plan;
     List<PhysicalPlan> nonCollisionCollections = new ArrayList<>();
@@ -142,90 +150,107 @@ public class MetaDataTransfer {
         idx++;
       } catch (Exception e) {
         logger.error("Parse mlog error at lineNumber {} because:", idx, e);
-        break;
+        throw e;
       }
       if (plan == null) {
         continue;
       }
-      try {
-        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);
+
+      switch (plan.getOperatorType()) {
+        case CREATE_TIMESERIES:
+        case CREATE_ALIGNED_TIMESERIES:
+        case AUTO_CREATE_DEVICE_MNODE:
+          nonCollisionCollections.add(plan);
+          if (nonCollisionCollections.size() > 100000) {
+            executeBufferedOperation(nonCollisionCollections);
+          }
+          break;
+        case SET_STORAGE_GROUP:
+        case DELETE_TIMESERIES:
+        case DELETE_STORAGE_GROUP:
+        case TTL:
+        case CHANGE_ALIAS:
+          executeBufferedOperation(nonCollisionCollections);
+          try {
             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);
-        }
+          } catch (IOException e) {
+            rocksDBManager.operation(plan);
+          } catch (MetadataException e) {
+            logger.error("Can not operate cmd {} for err:", 
plan.getOperatorType(), e);
+          }
+          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());
       }
     }
-    executeOperation(nonCollisionCollections, true);
+
+    executeBufferedOperation(nonCollisionCollections);
+
     if (retryPlans.size() > 0) {
-      executeOperation(retryPlans, false);
+      for (PhysicalPlan retryPlan : retryPlans) {
+        try {
+          rocksDBManager.operation(retryPlan);
+        } catch (IOException e) {
+          persistFailedLog(retryPlan);
+        } catch (MetadataException e) {
+          logger.error("Execute plan failed: {}", retryPlan.toString(), e);
+        } catch (Exception e) {
+          persistFailedLog(retryPlan);
+        }
+      }
     }
+    logger.info(
+        "Transfer data from mlog.bin complete after {}ms with {} errors",
+        System.currentTimeMillis() - time,
+        failedPlanCount.get());
   }
 
-  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());
+  private void executeBufferedOperation(List<PhysicalPlan> plans)
+      throws ExecutionException, InterruptedException {
+    if (plans.size() <= 0) {
+      return;
+    }
+    forkJoinPool
+        .submit(
+            () -> {
+              plans
+                  .parallelStream()
+                  .forEach(
+                      x -> {
+                        try {
+                          rocksDBManager.operation(x);
+                        } catch (IOException e) {
+                          retryPlans.add(x);
+                        } catch (MetadataException e) {
+                          if (e instanceof AcquireLockTimeoutException) {
+                            retryPlans.add(x);
+                          } else {
+                            logger.error("Execute plan failed: {}", 
x.toString(), e);
+                          }
+                        } catch (Exception e) {
+                          retryPlans.add(x);
+                        }
+                      });
+            })
+        .get();
+    logger.debug("parallel executed {} operations", plans.size());
     plans.clear();
   }
 
   private void persistFailedLog(PhysicalPlan plan) {
-    logger.info("persist won't retry and failed plan: {}", plan.toString());
+    logger.info("persist failed plan: {}", plan.toString());
     failedPlanCount.incrementAndGet();
     try {
       switch (plan.getOperatorType()) {
@@ -275,20 +300,15 @@ public class MetaDataTransfer {
       }
     } 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 | MetadataException e) {
-      logger.warn("Failed to deserialize from {}. Use a new MTree.", 
mtreeSnapshot.getPath());
+          "Fatal error, exception when persist failed plan, metadata transfer 
should be failed", e);
+      throw new RuntimeException("Terminate transfer as persist log fail.");
     }
   }
 
   @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity 
warning
-  private void doTransferFromSnapshot(MLogReader mLogReader) throws 
IOException, MetadataException {
+  private void doTransferFromSnapshot()
+      throws IOException, MetadataException, ExecutionException, 
InterruptedException {
+    logger.info("Start transfer data from snapshot");
     long start = System.currentTimeMillis();
     MTree mTree = new MTree();
     mTree.init();
@@ -316,15 +336,17 @@ public class MetaDataTransfer {
             });
 
     if (errorCount.get() > 0) {
-      logger.info("Fatal error. create some storage groups fail, terminate 
metadata transfer");
+      logger.error("Fatal error. Create some storage groups fail, terminate 
metadata transfer");
       return;
     }
 
     List<IMeasurementMNode> measurementMNodes = new ArrayList<>();
 
+    IMNode root = mTree.getNodeByPath(new PartialPath(new String[] {"root"}));
+    PartialPath matchAllPath = new PartialPath(new String[] {"root", "**"});
+
     MeasurementCollector collector =
-        new MeasurementCollector(
-            mTree.getNodeByPath(new PartialPath("root")), new 
PartialPath("root.**"), -1, -1) {
+        new MeasurementCollector(root, matchAllPath, -1, -1) {
           @Override
           protected void collectMeasurement(IMeasurementMNode node) throws 
MetadataException {
             measurementMNodes.add(node);
@@ -332,33 +354,70 @@ public class MetaDataTransfer {
         };
     collector.traverse();
 
-    measurementMNodes
-        .parallelStream()
-        .forEach(
-            mNode -> {
-              try {
-                rocksDBManager.createTimeSeries(
-                    mNode.getPartialPath(), mNode.getSchema(), 
mNode.getAlias(), null, null);
-              } catch (AcquireLockTimeoutException e) {
+    Queue<IMeasurementMNode> failCreatedNodes = new ConcurrentLinkedQueue<>();
+    AtomicInteger createdNodeCnt = new AtomicInteger(0);
+    AtomicInteger lastValue = new AtomicInteger(-1);
+    new Thread(
+            () -> {
+              while (lastValue.get() < createdNodeCnt.get()) {
+                try {
+                  lastValue.set(createdNodeCnt.get());
+                  Thread.sleep(10 * 1000);
+                  logger.info("created count: {}", createdNodeCnt.get());
+                } catch (InterruptedException e) {
+                  logger.error("timer thread error", e);
+                }
+              }
+            })
+        .start();
+
+    forkJoinPool
+        .submit(
+            () ->
+                measurementMNodes
+                    .parallelStream()
+                    .forEach(
+                        mNode -> {
+                          try {
+                            rocksDBManager.createTimeSeries(
+                                mNode.getPartialPath(),
+                                mNode.getSchema(),
+                                mNode.getAlias(),
+                                null,
+                                null);
+                            createdNodeCnt.incrementAndGet();
+                          } catch (AcquireLockTimeoutException e) {
+                            failCreatedNodes.add(mNode);
+                          } catch (MetadataException e) {
+                            logger.error(
+                                "create timeseries {} failed",
+                                mNode.getPartialPath().getFullPath(),
+                                e);
+                            errorCount.incrementAndGet();
+                          }
+                        }))
+        .get();
+
+    if (!failCreatedNodes.isEmpty()) {
+      failCreatedNodes.stream()
+          .forEach(
+              mNode -> {
                 try {
                   rocksDBManager.createTimeSeries(
                       mNode.getPartialPath(), mNode.getSchema(), 
mNode.getAlias(), null, null);
-                } catch (MetadataException metadataException) {
+                  createdNodeCnt.incrementAndGet();
+                } catch (Exception e) {
                   logger.error(
                       "create timeseries {} failed in retry",
                       mNode.getPartialPath().getFullPath(),
                       e);
                   errorCount.incrementAndGet();
                 }
-              } catch (MetadataException e) {
-                logger.error(
-                    "create timeseries {} failed", 
mNode.getPartialPath().getFullPath(), e);
-                errorCount.incrementAndGet();
-              }
-            });
+              });
+    }
 
     logger.info(
-        "metadata snapshot transfer complete after {}ms with {} errors",
+        "Transfer data from snapshot complete after {}ms with {} errors",
         System.currentTimeMillis() - start,
         errorCount.get());
   }
diff --git a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java 
b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
index 00d2e30..cb27781 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
@@ -35,6 +35,7 @@ import org.apache.iotdb.db.exception.StartupException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.metadata.IMetaManager;
 import org.apache.iotdb.db.metadata.MManager;
+import org.apache.iotdb.db.metadata.MetadataManagerType;
 import org.apache.iotdb.db.metadata.rocksdb.MRocksDBManager;
 import org.apache.iotdb.db.protocol.influxdb.meta.InfluxDBMetaManager;
 import org.apache.iotdb.db.protocol.rest.RestService;
@@ -78,7 +79,14 @@ public class IoTDB implements IoTDBMBean {
     }
 
     try {
-      metaManager = new MRocksDBManager();
+      if (IoTDBDescriptor.getInstance().getConfig().getMetadataManagerType()
+          == MetadataManagerType.ROCKSDB_MANAGER) {
+        metaManager = new MRocksDBManager();
+        logger.info("Use MRocksDBManager to manage metadata");
+      } else {
+        metaManager = MManager.getInstance();
+        logger.info("Use MManager to manage metadata");
+      }
     } catch (Exception e) {
       logger.error("create meta manager fail", e);
       System.exit(1);

Reply via email to