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",
