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);
