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(

Reply via email to