This is an automated email from the ASF dual-hosted git repository. lta pushed a commit to branch reimpl_sync in repository https://gitbox.apache.org/repos/asf/incubator-iotdb.git
commit c7b541f98232e89762169e9fca0c16cda36fa9ea Author: lta <[email protected]> AuthorDate: Tue Sep 3 11:26:15 2019 +0800 add log analyzer unit test --- .../org/apache/iotdb/db/concurrent/ThreadName.java | 1 + .../iotdb/db/sync/receiver/load/FileLoader.java | 39 ++-- .../db/sync/receiver/load/FileLoaderManager.java | 4 +- .../receiver/recover/SyncReceiverLogAnalyzer.java | 7 - .../db/sync/receiver/load/FileLoaderTest.java | 26 ++- .../recover/SyncReceiverLogAnalyzerTest.java | 205 +++++++++++++++++++++ .../receiver/recover/SyncReceiverLoggerTest.java | 1 + 7 files changed, 246 insertions(+), 37 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java b/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java index 64b89d1..57e95bc 100644 --- a/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java +++ b/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java @@ -39,6 +39,7 @@ public enum ThreadName { SYNC_CLIENT("Sync-Client"), SYNC_SERVER("Sync-Server"), SYNC_MONITOR("Sync-Monitor"), + LOAD_TSFILE("Load TsFile"), TIME_COST_STATSTIC("TIME_COST_STATSTIC"); private String name; diff --git a/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoader.java b/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoader.java index 5292536..520b42d 100644 --- a/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoader.java +++ b/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoader.java @@ -16,7 +16,9 @@ package org.apache.iotdb.db.sync.receiver.load; import java.io.File; import java.io.IOException; -import java.util.ArrayDeque; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; import org.apache.commons.io.FileUtils; import org.apache.iotdb.db.engine.StorageEngine; import org.apache.iotdb.db.engine.storagegroup.TsFileResource; @@ -30,13 +32,13 @@ public class FileLoader implements IFileLoader { private static final Logger LOGGER = LoggerFactory.getLogger(FileLoader.class); - private static final int WAIT_TIME = 1000; + private static final int WAIT_TIME = 100; private String syncFolderPath; private String senderName; - private ArrayDeque<LoadTask> queue = new ArrayDeque<>(); + private BlockingQueue<LoadTask> queue = new LinkedBlockingQueue<>(); private LoadLogger loadLog; @@ -65,23 +67,16 @@ public class FileLoader implements IFileLoader { private Runnable loadTaskRunner = () -> { try { while (true) { - if (queue.isEmpty()) { - if (endSync) { - cleanUp(); - break; - } - synchronized (queue) { - if (queue.isEmpty()) { - queue.wait(WAIT_TIME); - } - } + if (queue.isEmpty() && endSync) { + cleanUp(); + break; } - if (!queue.isEmpty()) { - LoadTask task = queue.poll(); + LoadTask loadTask = queue.poll(WAIT_TIME, TimeUnit.MILLISECONDS); + if (loadTask != null) { try { - handleLoadTask(task); + handleLoadTask(loadTask); } catch (IOException e) { - LOGGER.error("Can not load task {}", task, e); + LOGGER.error("Can not load task {}", loadTask, e); } } } @@ -92,18 +87,12 @@ public class FileLoader implements IFileLoader { @Override public void addDeletedFileName(File deletedFile) { - synchronized (queue) { - queue.add(new LoadTask(deletedFile, LoadType.DELETE)); - queue.notify(); - } + queue.add(new LoadTask(deletedFile, LoadType.DELETE)); } @Override public void addTsfile(File tsfile) { - synchronized (queue) { - queue.add(new LoadTask(tsfile, LoadType.ADD)); - queue.notify(); - } + queue.add(new LoadTask(tsfile, LoadType.ADD)); } @Override diff --git a/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderManager.java b/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderManager.java index 8fc341a..55187d3 100644 --- a/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderManager.java +++ b/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderManager.java @@ -22,6 +22,8 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.db.concurrent.ThreadName; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -68,7 +70,7 @@ public class FileLoaderManager { fileLoaderMap = new ConcurrentHashMap<>(); } if (loadTaskRunnerPool == null) { - loadTaskRunnerPool = Executors.newCachedThreadPool(); + loadTaskRunnerPool = IoTDBThreadPoolFactory.newCachedThreadPool(ThreadName.LOAD_TSFILE.getName()); } } diff --git a/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogAnalyzer.java b/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogAnalyzer.java index 53a44dd..f2516c4 100644 --- a/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogAnalyzer.java +++ b/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogAnalyzer.java @@ -36,8 +36,6 @@ public class SyncReceiverLogAnalyzer implements ISyncReceiverLogAnalyzer { private static final Logger LOGGER = LoggerFactory.getLogger(SyncReceiverLogAnalyzer.class); - private static final int WAIT_TIMEOUT = 2000; - private SyncReceiverLogAnalyzer() { } @@ -71,11 +69,6 @@ public class SyncReceiverLogAnalyzer implements ISyncReceiverLogAnalyzer { } if (FileLoaderManager.getInstance().containsFileLoader(senderFolder.getName())) { FileLoaderManager.getInstance().getFileLoader(senderFolder.getName()).endSync(); - try { - Thread.sleep(WAIT_TIMEOUT); // wait for file loader to clean up resource - } catch (InterruptedException e) { - LOGGER.error("Can not wait for recovery to complete.", e); - } } else { scanLogger(FileLoader.createFileLoader(senderFolder), new File(senderFolder, Constans.SYNC_LOG_NAME), diff --git a/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java b/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java index 3267e01..4354484 100644 --- a/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java +++ b/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java @@ -1,3 +1,21 @@ +/** + * 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. + */ package org.apache.iotdb.db.sync.receiver.load; import java.io.File; @@ -125,7 +143,7 @@ public class FileLoaderTest { .containsFileLoader(getReceiverFolderFile().getName())) { Thread.sleep(100); waitTime += 100; - LOGGER.info("Has waited for loading new tsfiles {}s", waitTime); + LOGGER.info("Has waited for loading new tsfiles {}ms", waitTime); } } catch (InterruptedException e) { LOGGER.error("Fail to wait for loading new tsfiles", e); @@ -211,7 +229,7 @@ public class FileLoaderTest { .containsFileLoader(getReceiverFolderFile().getName())) { Thread.sleep(100); waitTime += 100; - LOGGER.info("Has waited for loading new tsfiles {}s", waitTime); + LOGGER.info("Has waited for loading new tsfiles {}ms", waitTime); } } catch (InterruptedException e) { LOGGER.error("Fail to wait for loading new tsfiles", e); @@ -315,7 +333,7 @@ public class FileLoaderTest { .containsFileLoader(getReceiverFolderFile().getName())) { Thread.sleep(100); waitTime += 100; - LOGGER.info("Has waited for loading new tsfiles {}s", waitTime); + LOGGER.info("Has waited for loading new tsfiles {}ms", waitTime); } } catch (InterruptedException e) { LOGGER.error("Fail to wait for loading new tsfiles", e); @@ -371,7 +389,7 @@ public class FileLoaderTest { .containsFileLoader(getReceiverFolderFile().getName())) { Thread.sleep(100); waitTime += 100; - LOGGER.info("Has waited for loading new tsfiles {}s", waitTime); + LOGGER.info("Has waited for loading new tsfiles {}ms", waitTime); } } catch (InterruptedException e) { LOGGER.error("Fail to wait for loading new tsfiles", e); diff --git a/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogAnalyzerTest.java b/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogAnalyzerTest.java new file mode 100644 index 0000000..7a591a0 --- /dev/null +++ b/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogAnalyzerTest.java @@ -0,0 +1,205 @@ +/** + * 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. + */ +package org.apache.iotdb.db.sync.receiver.recover; + +import java.io.BufferedReader; +import java.io.File; +import java.io.FileReader; +import java.io.IOException; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Map; +import java.util.Random; +import java.util.Set; +import org.apache.iotdb.db.conf.IoTDBConstant; +import org.apache.iotdb.db.conf.directories.DirectoryManager; +import org.apache.iotdb.db.engine.StorageEngine; +import org.apache.iotdb.db.engine.storagegroup.StorageGroupProcessor; +import org.apache.iotdb.db.engine.storagegroup.TsFileResource; +import org.apache.iotdb.db.exception.DiskSpaceInsufficientException; +import org.apache.iotdb.db.exception.MetadataErrorException; +import org.apache.iotdb.db.exception.StartupException; +import org.apache.iotdb.db.exception.StorageEngineException; +import org.apache.iotdb.db.metadata.MManager; +import org.apache.iotdb.db.service.IoTDB; +import org.apache.iotdb.db.sync.receiver.load.FileLoader; +import org.apache.iotdb.db.sync.receiver.load.FileLoaderManager; +import org.apache.iotdb.db.sync.receiver.load.FileLoaderTest; +import org.apache.iotdb.db.sync.sender.conf.Constans; +import org.apache.iotdb.db.utils.EnvironmentUtils; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class SyncReceiverLogAnalyzerTest { + + private static final Logger LOGGER = LoggerFactory.getLogger(FileLoaderTest.class); + private static final String SG_NAME = "root.sg"; + private static IoTDB daemon; + private String dataDir; + private FileLoader fileLoader; + private SyncReceiverLogAnalyzer logAnalyze; + private SyncReceiverLogger receiverLogger; + + @Before + public void setUp() + throws IOException, InterruptedException, StartupException, DiskSpaceInsufficientException, MetadataErrorException { + EnvironmentUtils.closeStatMonitor(); + daemon = IoTDB.getInstance(); + daemon.active(); + EnvironmentUtils.envSetUp(); + dataDir = new File(DirectoryManager.getInstance().getNextFolderForSequenceFile()) + .getParentFile().getAbsolutePath(); + logAnalyze = SyncReceiverLogAnalyzer.getInstance(); + initMetadata(); + } + + private void initMetadata() throws MetadataErrorException { + MManager mmanager = MManager.getInstance(); + mmanager.init(); + mmanager.clear(); + mmanager.setStorageLevelToMTree("root.sg0"); + mmanager.setStorageLevelToMTree("root.sg1"); + mmanager.setStorageLevelToMTree("root.sg2"); + } + + @After + public void tearDown() throws InterruptedException, IOException, StorageEngineException { + daemon.stop(); + EnvironmentUtils.cleanEnv(); + } + + @Test + public void recover() throws IOException, StorageEngineException { + receiverLogger = new SyncReceiverLogger( + new File(getReceiverFolderFile(), Constans.SYNC_LOG_NAME)); + fileLoader = FileLoader.createFileLoader(getReceiverFolderFile()); + Map<String, Set<File>> allFileList = new HashMap<>(); + Map<String, Set<File>> correctSequenceLoadedFileMap = new HashMap<>(); + + // add some new tsfiles + Random r = new Random(0); + receiverLogger.startSyncTsFiles(); + Set<String> toBeSyncedFiles = new HashSet<>(); + for (int i = 0; i < 3; i++) { + for (int j = 0; j < 10; j++) { + allFileList.putIfAbsent(SG_NAME + i, new HashSet<>()); + correctSequenceLoadedFileMap.putIfAbsent(SG_NAME + i, new HashSet<>()); + String rand = String.valueOf(r.nextInt(10000)); + String fileName = + getSnapshotFolder() + File.separator + SG_NAME + i + File.separator + rand + ".tsfile"; + File syncFile = new File(fileName); + receiverLogger + .finishSyncTsfile(syncFile); + toBeSyncedFiles.add(syncFile.getAbsolutePath()); + File dataFile = new File( + syncFile.getParentFile().getParentFile().getParentFile().getParentFile() + .getParentFile(), IoTDBConstant.SEQUENCE_FLODER_NAME + + File.separatorChar + syncFile.getParentFile().getName() + File.separatorChar + + syncFile.getName()); + correctSequenceLoadedFileMap.get(SG_NAME + i).add(dataFile); + allFileList.get(SG_NAME + i).add(syncFile); + if (!syncFile.getParentFile().exists()) { + syncFile.getParentFile().mkdirs(); + } + if (!syncFile.exists() && !syncFile.createNewFile()) { + LOGGER.error("Can not create new file {}", syncFile.getPath()); + } + if (!new File(syncFile.getAbsolutePath() + TsFileResource.RESOURCE_SUFFIX).exists() + && !new File(syncFile.getAbsolutePath() + TsFileResource.RESOURCE_SUFFIX) + .createNewFile()) { + LOGGER.error("Can not create new file {}", syncFile.getPath()); + } + TsFileResource tsFileResource = new TsFileResource(syncFile); + tsFileResource.getStartTimeMap().put(String.valueOf(i), (long) j * 10); + tsFileResource.getEndTimeMap().put(String.valueOf(i), (long) j * 10 + 5); + tsFileResource.serialize(); + } + } + + for (int i = 0; i < 3; i++) { + StorageGroupProcessor processor = StorageEngine.getInstance().getProcessor(SG_NAME + i); + assert processor.getSequenceFileList().isEmpty(); + assert processor.getUnSequenceFileList().isEmpty(); + } + + assert getReceiverFolderFile().exists(); + for (Set<File> set : allFileList.values()) { + for (File newTsFile : set) { + if (!newTsFile.getName().endsWith(TsFileResource.RESOURCE_SUFFIX)) { + fileLoader.addTsfile(newTsFile); + } + } + } + + receiverLogger.close(); + assert new File(getReceiverFolderFile(), Constans.LOAD_LOG_NAME).exists(); + assert new File(getReceiverFolderFile(), Constans.SYNC_LOG_NAME).exists(); + assert FileLoaderManager.getInstance().containsFileLoader(getReceiverFolderFile().getName()); + int count = 0, mode = 0; + Set<String> toBeSyncedFilesTest = new HashSet<>(); + try (BufferedReader br = new BufferedReader( + new FileReader(new File(getReceiverFolderFile(), Constans.SYNC_LOG_NAME)))) { + String line; + while ((line = br.readLine()) != null) { + count++; + if (line.equals(SyncReceiverLogger.SYNC_DELETED_FILE_NAME_START)) { + mode = -1; + } else if (line.equals(SyncReceiverLogger.SYNC_TSFILE_START)) { + mode = 1; + } else { + if (mode == 1) { + toBeSyncedFilesTest.add(line); + } + } + } + } + assert toBeSyncedFilesTest.size() == toBeSyncedFiles.size(); + assert toBeSyncedFilesTest.containsAll(toBeSyncedFiles); + + logAnalyze.recover(getReceiverFolderFile().getName()); + + try { + long waitTime = 0; + while (FileLoaderManager.getInstance() + .containsFileLoader(getReceiverFolderFile().getName())) { + Thread.sleep(100); + waitTime += 100; + LOGGER.info("Has waited for loading new tsfiles {}ms", waitTime); + } + } catch (InterruptedException e) { + LOGGER.error("Fail to wait for loading new tsfiles", e); + } + + assert !new File(getReceiverFolderFile(), Constans.LOAD_LOG_NAME).exists(); + assert !new File(getReceiverFolderFile(), Constans.SYNC_LOG_NAME).exists(); + } + + + private File getReceiverFolderFile() { + return new File(dataDir + File.separatorChar + Constans.SYNC_RECEIVER + File.separatorChar + + "127.0.0.1_5555"); + } + + private File getSnapshotFolder() { + return new File(getReceiverFolderFile(), Constans.RECEIVER_DATA_FOLDER_NAME); + } +} \ No newline at end of file diff --git a/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLoggerTest.java b/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLoggerTest.java index c3849ec..c6b248d 100644 --- a/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLoggerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLoggerTest.java @@ -74,6 +74,7 @@ public class SyncReceiverLoggerTest { toBeSyncedFiles .add(new File(getReceiverFolderFile(), "new" + i).getAbsolutePath()); } + receiverLogger.close(); int count = 0; int mode = 0; try (BufferedReader br = new BufferedReader(
