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
The following commit(s) were added to refs/heads/reimpl_sync by this push:
new bcd3107 add load tsfile unit test
bcd3107 is described below
commit bcd3107d247edcabe321e26c51aa409a5bda8251
Author: lta <[email protected]>
AuthorDate: Mon Sep 2 21:22:05 2019 +0800
add load tsfile unit test
---
.../org/apache/iotdb/db/engine/StorageEngine.java | 11 +-
.../iotdb/db/engine/merge/task/MergeFileTask.java | 10 +-
.../engine/storagegroup/StorageGroupProcessor.java | 32 +-
.../iotdb/db/sync/receiver/load/FileLoader.java | 52 +--
.../iotdb/db/sync/receiver/load/LoadLogger.java | 3 +
.../sync/receiver/recover/SyncReceiverLogger.java | 3 +
.../sender/recover/ISyncSenderLogAnalyzer.java | 2 +-
.../sync/sender/recover/SyncSenderLogAnalyzer.java | 4 +-
.../db/sync/sender/recover/SyncSenderLogger.java | 9 +-
.../sync/sender/transfer/DataTransferManager.java | 1 -
.../db/sync/receiver/load/FileLoaderTest.java | 406 +++++++++++++++++++++
.../receiver/recover/SyncReceiverLoggerTest.java | 108 ++++++
12 files changed, 587 insertions(+), 54 deletions(-)
diff --git a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
index c96ef2e..e061634 100644
--- a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
+++ b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
@@ -115,8 +115,7 @@ public class StorageEngine implements IService {
return ServiceType.STORAGE_ENGINE_SERVICE;
}
-
- private StorageGroupProcessor getProcessor(String path) throws
StorageEngineException {
+ public StorageGroupProcessor getProcessor(String path) throws
StorageEngineException {
String storageGroupName = "";
try {
storageGroupName =
MManager.getInstance().getStorageGroupNameByPath(path);
@@ -340,12 +339,12 @@ public class StorageEngine implements IService {
public void loadNewTsFile(File newTsFile, TsFileResource resource)
- throws TsFileProcessorException {
-
processorMap.get(newTsFile.getParentFile().getName()).loadNewTsFile(newTsFile,
resource);
+ throws TsFileProcessorException, StorageEngineException {
+ getProcessor(newTsFile.getParentFile().getName()).loadNewTsFile(newTsFile,
resource);
}
- public void deleteTsfile(File deletedTsfile){
-
processorMap.get(deletedTsfile.getParentFile().getName()).deleteTsfile(deletedTsfile);
+ public void deleteTsfile(File deletedTsfile) throws StorageEngineException {
+
getProcessor(deletedTsfile.getParentFile().getName()).deleteTsfile(deletedTsfile);
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeFileTask.java
b/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeFileTask.java
index 68d3141..3e5dc84 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeFileTask.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/merge/task/MergeFileTask.java
@@ -154,8 +154,9 @@ class MergeFileTask {
File nextMergeVersionFile = getNextMergeVersionFile(seqFile.getFile());
FileUtils.moveFile(seqFile.getFile(), nextMergeVersionFile);
- FileUtils.moveFile(new File(seqFile.getFile(),
TsFileResource.RESOURCE_SUFFIX),
- new File(nextMergeVersionFile, TsFileResource.RESOURCE_SUFFIX));
+ FileUtils
+ .moveFile(new File(seqFile.getFile().getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX),
+ new File(nextMergeVersionFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX));
seqFile.setFile(nextMergeVersionFile);
} finally {
seqFile.getMergeQueryLock().writeLock().unlock();
@@ -221,8 +222,9 @@ class MergeFileTask {
File nextMergeVersionFile = getNextMergeVersionFile(seqFile.getFile());
FileUtils.moveFile(fileWriter.getFile(), nextMergeVersionFile);
- FileUtils.moveFile(new File(seqFile.getFile(),
TsFileResource.RESOURCE_SUFFIX),
- new File(nextMergeVersionFile, TsFileResource.RESOURCE_SUFFIX));
+ FileUtils
+ .moveFile(new File(seqFile.getFile().getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX),
+ new File(nextMergeVersionFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX));
seqFile.setFile(nextMergeVersionFile);
} finally {
seqFile.getMergeQueryLock().writeLock().unlock();
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
index 4da8e30..3555e3d 100755
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
@@ -1051,10 +1051,12 @@ public class StorageGroupProcessor {
for (Entry<String, Long> entry :
newTsFileResource.getEndTimeMap().entrySet()) {
String device = entry.getKey();
long endTime = newTsFileResource.getEndTimeMap().get(device);
- if(!latestTimeForEachDevice.containsKey(device) ||
latestTimeForEachDevice.get(device) < endTime){
+ if (!latestTimeForEachDevice.containsKey(device)
+ || latestTimeForEachDevice.get(device) < endTime) {
latestTimeForEachDevice.put(device, endTime);
}
- if(!latestFlushedTimeForEachDevice.containsKey(device) ||
latestFlushedTimeForEachDevice.get(device) < endTime){
+ if (!latestFlushedTimeForEachDevice.containsKey(device)
+ || latestFlushedTimeForEachDevice.get(device) < endTime) {
latestFlushedTimeForEachDevice.put(device, endTime);
}
}
@@ -1094,16 +1096,16 @@ public class StorageGroupProcessor {
if (!targetFile.getParentFile().exists()) {
targetFile.getParentFile().mkdirs();
}
- if (!new File(tsFile, TsFileResource.RESOURCE_SUFFIX).exists() && !new
File(
- targetFile, TsFileResource.RESOURCE_SUFFIX).exists()) {
+ if (!new File(tsFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX).exists() && !new File(
+ targetFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX).exists()) {
throw new TsFileProcessorException(
String
.format("The new .resource file {%s} to be loaded does not
exist.",
tsFile.getAbsolutePath()));
}
- if (!new File(targetFile, TsFileResource.RESOURCE_SUFFIX).exists() && !new
File(
- tsFile, TsFileResource.RESOURCE_SUFFIX)
- .renameTo(new File(targetFile, TsFileResource.RESOURCE_SUFFIX))) {
+ if (!new File(targetFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX).exists()
+ && !new File(tsFile.getAbsolutePath() + TsFileResource.RESOURCE_SUFFIX)
+ .renameTo(new File(targetFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX))) {
throw new TsFileProcessorException(String.format(
"File renaming failed when loading .resource file. Origin: %s,
Target: %s",
new File(tsFile, TsFileResource.RESOURCE_SUFFIX).getAbsolutePath(),
@@ -1122,7 +1124,7 @@ public class StorageGroupProcessor {
}
/**
- * Delete tsfile if it exists which.
+ * Delete tsfile if it exists.
*
* Firstly, remove the TsFileResource from
sequenceFileList/unSequenceFileList.
*
@@ -1145,7 +1147,7 @@ public class StorageGroupProcessor {
break;
}
}
- if(deletedTsFileResource == null) {
+ if (deletedTsFileResource == null) {
Iterator<TsFileResource> unsequenceIterator =
unSequenceFileList.iterator();
while (unsequenceIterator.hasNext()) {
TsFileResource unsequenceResource = unsequenceIterator.next();
@@ -1160,13 +1162,13 @@ public class StorageGroupProcessor {
mergeLock.writeLock().unlock();
writeUnlock();
}
- if(deletedTsFileResource == null){
+ if (deletedTsFileResource == null) {
return;
}
deletedTsFileResource.getMergeQueryLock().writeLock().lock();
try {
deletedTsFileResource.getFile().delete();
- new File(deletedTsFileResource.getFile(),
TsFileResource.RESOURCE_SUFFIX).delete();
+ new File(deletedTsFileResource.getFile().getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX).delete();
} finally {
deletedTsFileResource.getMergeQueryLock().writeLock().unlock();
}
@@ -1176,6 +1178,14 @@ public class StorageGroupProcessor {
return workSequenceTsFileProcessor;
}
+ public List<TsFileResource> getSequenceFileList() {
+ return sequenceFileList;
+ }
+
+ public List<TsFileResource> getUnSequenceFileList() {
+ return unSequenceFileList;
+ }
+
@FunctionalInterface
public interface CloseTsFileCallBack {
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 6bd7d49..5292536 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
@@ -1,19 +1,15 @@
/**
- * 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
+ * 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
+ * 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
+ * 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;
@@ -21,8 +17,10 @@ package org.apache.iotdb.db.sync.receiver.load;
import java.io.File;
import java.io.IOException;
import java.util.ArrayDeque;
+import org.apache.commons.io.FileUtils;
import org.apache.iotdb.db.engine.StorageEngine;
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.db.exception.StorageEngineException;
import org.apache.iotdb.db.exception.TsFileProcessorException;
import org.apache.iotdb.db.sync.sender.conf.Constans;
import org.slf4j.Logger;
@@ -82,7 +80,7 @@ public class FileLoader implements IFileLoader {
LoadTask task = queue.poll();
try {
handleLoadTask(task);
- }catch (IOException e){
+ } catch (IOException e) {
LOGGER.error("Can not load task {}", task, e);
}
}
@@ -130,27 +128,32 @@ public class FileLoader implements IFileLoader {
}
private void loadNewTsfile(File newTsFile) throws IOException {
- if (curType != LoadType.DELETE) {
- loadLog.startLoadDeletedFiles();
- curType = LoadType.DELETE;
+ if (curType != LoadType.ADD) {
+ loadLog.startLoadTsFiles();
+ curType = LoadType.ADD;
}
- TsFileResource tsFileResource = new TsFileResource(
- new File(newTsFile, TsFileResource.RESOURCE_SUFFIX));
+ TsFileResource tsFileResource = new TsFileResource(new
File(newTsFile.getAbsolutePath()));
tsFileResource.deSerialize();
try {
StorageEngine.getInstance().loadNewTsFile(newTsFile, tsFileResource);
- } catch (TsFileProcessorException e) {
+ } catch (TsFileProcessorException | StorageEngineException e) {
LOGGER.error("Can not load new tsfile {}", newTsFile.getAbsolutePath(),
e);
+ throw new IOException(e);
}
loadLog.finishLoadDeletedFile(newTsFile);
}
private void loadDeletedFile(File deletedTsFile) throws IOException {
- if (curType != LoadType.ADD) {
- loadLog.startLoadTsFiles();
- curType = LoadType.ADD;
+ if (curType != LoadType.DELETE) {
+ loadLog.startLoadDeletedFiles();
+ curType = LoadType.DELETE;
+ }
+ try {
+ StorageEngine.getInstance().deleteTsfile(deletedTsFile);
+ } catch (StorageEngineException e) {
+ LOGGER.error("Can not load deleted tsfile {}",
deletedTsFile.getAbsolutePath(), e);
+ throw new IOException(e);
}
- StorageEngine.getInstance().deleteTsfile(deletedTsFile);
loadLog.finishLoadTsfile(deletedTsFile);
}
@@ -161,6 +164,7 @@ public class FileLoader implements IFileLoader {
loadLog.close();
new File(syncFolderPath, Constans.SYNC_LOG_NAME).delete();
new File(syncFolderPath, Constans.LOAD_LOG_NAME).delete();
+ FileUtils.deleteDirectory(new File(syncFolderPath,
Constans.RECEIVER_DATA_FOLDER_NAME));
FileLoaderManager.getInstance().removeFileLoader(senderName);
} catch (IOException e) {
LOGGER.error("Can not clean up sync resource.", e);
diff --git
a/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/LoadLogger.java
b/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/LoadLogger.java
index 881a637..5f54d6b 100644
---
a/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/LoadLogger.java
+++
b/server/src/main/java/org/apache/iotdb/db/sync/receiver/load/LoadLogger.java
@@ -30,6 +30,9 @@ public class LoadLogger implements ILoadLogger {
private BufferedWriter bw;
public LoadLogger(File logFile) throws IOException {
+ if (!logFile.getParentFile().exists()) {
+ logFile.getParentFile().mkdirs();
+ }
bw = new BufferedWriter(new FileWriter(logFile));
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogger.java
b/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogger.java
index b940758..4da6c80 100644
---
a/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogger.java
+++
b/server/src/main/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLogger.java
@@ -30,6 +30,9 @@ public class SyncReceiverLogger implements
ISyncReceiverLogger {
private BufferedWriter bw;
public SyncReceiverLogger(File logFile) throws IOException {
+ if (!logFile.getParentFile().exists()) {
+ logFile.getParentFile().mkdirs();
+ }
bw = new BufferedWriter(new FileWriter(logFile));
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/ISyncSenderLogAnalyzer.java
b/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/ISyncSenderLogAnalyzer.java
index 652839b..2f2177b 100644
---
a/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/ISyncSenderLogAnalyzer.java
+++
b/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/ISyncSenderLogAnalyzer.java
@@ -34,6 +34,6 @@ public interface ISyncSenderLogAnalyzer {
void loadLogger(Set<String> deletedFiles, Set<String> newFiles);
- void updateLastLocalFile(Set<String> currentLocalFiles);
+ void updateLastLocalFile(Set<String> currentLocalFiles) throws IOException;
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogAnalyzer.java
b/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogAnalyzer.java
index 5f97f05..d4c65ef 100644
---
a/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogAnalyzer.java
+++
b/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogAnalyzer.java
@@ -112,7 +112,7 @@ public class SyncSenderLogAnalyzer implements
ISyncSenderLogAnalyzer {
}
@Override
- public void updateLastLocalFile(Set<String> currentLocalFiles) {
+ public void updateLastLocalFile(Set<String> currentLocalFiles) throws
IOException {
try (BufferedWriter bw = new BufferedWriter(new
FileWriter(currentLocalFile))) {
for (String line : currentLocalFiles) {
bw.write(line);
@@ -123,6 +123,6 @@ public class SyncSenderLogAnalyzer implements
ISyncSenderLogAnalyzer {
LOGGER.error("Can not clear sync log {}", syncLogFile.getAbsoluteFile(),
e);
}
lastLocalFile.delete();
- currentLocalFile.renameTo(lastLocalFile);
+ FileUtils.moveFile(currentLocalFile, lastLocalFile);
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogger.java
b/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogger.java
index 8e118d7..0b78d3b 100644
---
a/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogger.java
+++
b/server/src/main/java/org/apache/iotdb/db/sync/sender/recover/SyncSenderLogger.java
@@ -29,12 +29,11 @@ public class SyncSenderLogger implements ISyncSenderLogger {
public static final String SYNC_TSFILE_START = "sync tsfile start";
private BufferedWriter bw;
- public SyncSenderLogger(String filePath) throws IOException {
- this.bw = new BufferedWriter(new FileWriter(filePath));
- }
-
public SyncSenderLogger(File file) throws IOException {
- this(file.getAbsolutePath());
+ if (!file.getParentFile().exists()) {
+ file.getParentFile().mkdirs();
+ }
+ this.bw = new BufferedWriter(new FileWriter(file.getAbsolutePath()));
}
@Override
diff --git
a/server/src/main/java/org/apache/iotdb/db/sync/sender/transfer/DataTransferManager.java
b/server/src/main/java/org/apache/iotdb/db/sync/sender/transfer/DataTransferManager.java
index aecf59e..39ebe20 100644
---
a/server/src/main/java/org/apache/iotdb/db/sync/sender/transfer/DataTransferManager.java
+++
b/server/src/main/java/org/apache/iotdb/db/sync/sender/transfer/DataTransferManager.java
@@ -434,7 +434,6 @@ public class DataTransferManager implements
IDataTransferManager {
public void sync() throws IOException {
try {
syncStatus = true;
- syncLog = new SyncSenderLogger(getSchemaLogFile());
for (String sgName : allSG) {
lastLocalFilesMap.putIfAbsent(sgName, new HashSet<>());
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
new file mode 100644
index 0000000..3267e01
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java
@@ -0,0 +1,406 @@
+package org.apache.iotdb.db.sync.receiver.load;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Map.Entry;
+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.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 FileLoaderTest {
+
+ 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;
+
+ @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();
+ 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 loadNewTsfiles() throws IOException, StorageEngineException {
+ 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);
+ 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);
+ 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);
+ }
+ }
+ }
+ fileLoader.endSync();
+
+ try {
+ long waitTime = 0;
+ while (FileLoaderManager.getInstance()
+ .containsFileLoader(getReceiverFolderFile().getName())) {
+ Thread.sleep(100);
+ waitTime += 100;
+ LOGGER.info("Has waited for loading new tsfiles {}s", waitTime);
+ }
+ } catch (InterruptedException e) {
+ LOGGER.error("Fail to wait for loading new tsfiles", e);
+ }
+
+ assert !new File(getReceiverFolderFile(),
Constans.RECEIVER_DATA_FOLDER_NAME).exists();
+ Map<String, Set<File>> sequenceLoadedFileMap = new HashMap<>();
+ for (int i = 0; i < 3; i++) {
+ StorageGroupProcessor processor =
StorageEngine.getInstance().getProcessor(SG_NAME + i);
+ sequenceLoadedFileMap.putIfAbsent(SG_NAME + i, new HashSet<>());
+ assert processor.getSequenceFileList().size() == 10;
+ for (TsFileResource tsFileResource : processor.getSequenceFileList()) {
+ sequenceLoadedFileMap.get(SG_NAME + i).add(tsFileResource.getFile());
+ }
+ assert processor.getUnSequenceFileList().isEmpty();
+ }
+
+ assert sequenceLoadedFileMap.size() == correctSequenceLoadedFileMap.size();
+ for (Entry<String, Set<File>> entry :
correctSequenceLoadedFileMap.entrySet()) {
+ String sg = entry.getKey();
+ assert entry.getValue().size() == sequenceLoadedFileMap.get(sg).size();
+ assert entry.getValue().containsAll(sequenceLoadedFileMap.get(sg));
+ }
+
+
+ // add some overlap new tsfiles
+ fileLoader = FileLoader.createFileLoader(getReceiverFolderFile());
+ Map<String, Set<File>> correctUnSequenceLoadedFileMap = new HashMap<>();
+ allFileList = new HashMap<>();
+ r = new Random(1);
+ for (int i = 0; i < 3; i++) {
+ for (int j = 0; j < 10; j++) {
+ allFileList.putIfAbsent(SG_NAME + i, new HashSet<>());
+ correctUnSequenceLoadedFileMap.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);
+ File dataFile = new File(
+
syncFile.getParentFile().getParentFile().getParentFile().getParentFile()
+ .getParentFile(), IoTDBConstant.UNSEQUENCE_FLODER_NAME
+ + File.separatorChar + syncFile.getParentFile().getName() +
File.separatorChar
+ + syncFile.getName());
+ correctUnSequenceLoadedFileMap.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 +
3);
+ 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);
+ }
+ }
+ }
+ fileLoader.endSync();
+
+ try {
+ long waitTime = 0;
+ while (FileLoaderManager.getInstance()
+ .containsFileLoader(getReceiverFolderFile().getName())) {
+ Thread.sleep(100);
+ waitTime += 100;
+ LOGGER.info("Has waited for loading new tsfiles {}s", waitTime);
+ }
+ } catch (InterruptedException e) {
+ LOGGER.error("Fail to wait for loading new tsfiles", e);
+ }
+
+ assert !new File(getReceiverFolderFile(),
Constans.RECEIVER_DATA_FOLDER_NAME).exists();
+ sequenceLoadedFileMap = new HashMap<>();
+ for (int i = 0; i < 3; i++) {
+ StorageGroupProcessor processor =
StorageEngine.getInstance().getProcessor(SG_NAME + i);
+ sequenceLoadedFileMap.putIfAbsent(SG_NAME + i, new HashSet<>());
+ assert processor.getSequenceFileList().size() == 10;
+ for (TsFileResource tsFileResource : processor.getSequenceFileList()) {
+ sequenceLoadedFileMap.get(SG_NAME + i).add(tsFileResource.getFile());
+ }
+ assert !processor.getUnSequenceFileList().isEmpty();
+ }
+
+ assert sequenceLoadedFileMap.size() == correctSequenceLoadedFileMap.size();
+ for (Entry<String, Set<File>> entry :
correctSequenceLoadedFileMap.entrySet()) {
+ String sg = entry.getKey();
+ assert entry.getValue().size() == sequenceLoadedFileMap.get(sg).size();
+ assert entry.getValue().containsAll(sequenceLoadedFileMap.get(sg));
+ }
+
+ Map<String, Set<File>> unsequenceLoadedFileMap = new HashMap<>();
+ for (int i = 0; i < 3; i++) {
+ StorageGroupProcessor processor =
StorageEngine.getInstance().getProcessor(SG_NAME + i);
+ unsequenceLoadedFileMap.putIfAbsent(SG_NAME + i, new HashSet<>());
+ assert processor.getUnSequenceFileList().size() == 10;
+ for (TsFileResource tsFileResource : processor.getUnSequenceFileList()) {
+ unsequenceLoadedFileMap.get(SG_NAME + i).add(tsFileResource.getFile());
+ }
+ }
+
+ assert unsequenceLoadedFileMap.size() ==
correctUnSequenceLoadedFileMap.size();
+ for (Entry<String, Set<File>> entry :
correctUnSequenceLoadedFileMap.entrySet()) {
+ String sg = entry.getKey();
+ assert entry.getValue().size() == unsequenceLoadedFileMap.get(sg).size();
+ assert entry.getValue().containsAll(unsequenceLoadedFileMap.get(sg));
+ }
+ }
+
+ @Test
+ public void loadDeletedFileName() throws IOException, StorageEngineException
{
+ fileLoader = FileLoader.createFileLoader(getReceiverFolderFile());
+ Map<String, Set<File>> allFileList = new HashMap<>();
+ Map<String, Set<File>> correctLoadedFileMap = new HashMap<>();
+
+ // add some tsfiles
+ Random r = new Random(0);
+ for (int i = 0; i < 3; i++) {
+ for (int j = 0; j < 25; j++) {
+ allFileList.putIfAbsent(SG_NAME + i, new HashSet<>());
+ correctLoadedFileMap.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);
+ File dataFile = new File(
+
syncFile.getParentFile().getParentFile().getParentFile().getParentFile()
+ .getParentFile(), IoTDBConstant.SEQUENCE_FLODER_NAME
+ + File.separatorChar + syncFile.getParentFile().getName() +
File.separatorChar
+ + syncFile.getName());
+ correctLoadedFileMap.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.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);
+ }
+ }
+ }
+ fileLoader.endSync();
+
+ try {
+ long waitTime = 0;
+ while (FileLoaderManager.getInstance()
+ .containsFileLoader(getReceiverFolderFile().getName())) {
+ Thread.sleep(100);
+ waitTime += 100;
+ LOGGER.info("Has waited for loading new tsfiles {}s", waitTime);
+ }
+ } catch (InterruptedException e) {
+ LOGGER.error("Fail to wait for loading new tsfiles", e);
+ }
+
+ assert !new File(getReceiverFolderFile(),
Constans.RECEIVER_DATA_FOLDER_NAME).exists();
+ Map<String, Set<File>> loadedFileMap = new HashMap<>();
+ for (int i = 0; i < 3; i++) {
+ StorageGroupProcessor processor =
StorageEngine.getInstance().getProcessor(SG_NAME + i);
+ loadedFileMap.putIfAbsent(SG_NAME + i, new HashSet<>());
+ assert processor.getSequenceFileList().size() == 25;
+ for (TsFileResource tsFileResource : processor.getSequenceFileList()) {
+ loadedFileMap.get(SG_NAME + i).add(tsFileResource.getFile());
+ }
+ assert processor.getUnSequenceFileList().isEmpty();
+ }
+
+ assert loadedFileMap.size() == correctLoadedFileMap.size();
+ for (Entry<String, Set<File>> entry : correctLoadedFileMap.entrySet()) {
+ String sg = entry.getKey();
+ assert entry.getValue().size() == loadedFileMap.get(sg).size();
+ assert entry.getValue().containsAll(loadedFileMap.get(sg));
+ }
+
+ // delete some tsfiles
+ fileLoader = FileLoader.createFileLoader(getReceiverFolderFile());
+ for(Entry<String, Set<File>> entry:allFileList.entrySet()){
+ String sg = entry.getKey();
+ Set<File> files = entry.getValue();
+ int cnt = 0;
+ for(File snapFile:files){
+ if(!snapFile.getName().endsWith(TsFileResource.RESOURCE_SUFFIX)){
+ File dataFile = new File(
+
snapFile.getParentFile().getParentFile().getParentFile().getParentFile()
+ .getParentFile(), IoTDBConstant.SEQUENCE_FLODER_NAME
+ + File.separatorChar + snapFile.getParentFile().getName() +
File.separatorChar
+ + snapFile.getName());
+ correctLoadedFileMap.get(sg).remove(dataFile);
+ snapFile.delete();
+ fileLoader.addDeletedFileName(snapFile);
+ new File(snapFile + TsFileResource.RESOURCE_SUFFIX).delete();
+ if(++cnt == 15){
+ break;
+ }
+ }
+ }
+ }
+ fileLoader.endSync();
+
+ try {
+ long waitTime = 0;
+ while (FileLoaderManager.getInstance()
+ .containsFileLoader(getReceiverFolderFile().getName())) {
+ Thread.sleep(100);
+ waitTime += 100;
+ LOGGER.info("Has waited for loading new tsfiles {}s", waitTime);
+ }
+ } catch (InterruptedException e) {
+ LOGGER.error("Fail to wait for loading new tsfiles", e);
+ }
+
+ loadedFileMap = new HashMap<>();
+ for (int i = 0; i < 3; i++) {
+ StorageGroupProcessor processor =
StorageEngine.getInstance().getProcessor(SG_NAME + i);
+ loadedFileMap.putIfAbsent(SG_NAME + i, new HashSet<>());
+ for (TsFileResource tsFileResource : processor.getSequenceFileList()) {
+ loadedFileMap.get(SG_NAME + i).add(tsFileResource.getFile());
+ }
+ assert processor.getUnSequenceFileList().isEmpty();
+ }
+
+ assert loadedFileMap.size() == correctLoadedFileMap.size();
+ for (Entry<String, Set<File>> entry : correctLoadedFileMap.entrySet()) {
+ String sg = entry.getKey();
+ assert entry.getValue().size() == loadedFileMap.get(sg).size();
+ assert entry.getValue().containsAll(loadedFileMap.get(sg));
+ }
+ }
+
+ 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
new file mode 100644
index 0000000..c3849ec
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/sync/receiver/recover/SyncReceiverLoggerTest.java
@@ -0,0 +1,108 @@
+/**
+ * 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.HashSet;
+import java.util.Set;
+import org.apache.iotdb.db.conf.directories.DirectoryManager;
+import org.apache.iotdb.db.exception.DiskSpaceInsufficientException;
+import org.apache.iotdb.db.exception.StartupException;
+import org.apache.iotdb.db.exception.StorageEngineException;
+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;
+
+public class SyncReceiverLoggerTest {
+
+ private SyncReceiverLogger receiverLogger;
+ private String dataDir;
+
+ @Before
+ public void setUp()
+ throws IOException, InterruptedException, StartupException,
DiskSpaceInsufficientException {
+ EnvironmentUtils.envSetUp();
+ dataDir = new
File(DirectoryManager.getInstance().getNextFolderForSequenceFile())
+ .getParentFile().getAbsolutePath();
+ }
+
+ @After
+ public void tearDown() throws InterruptedException, IOException,
StorageEngineException {
+ EnvironmentUtils.cleanEnv();
+ }
+
+ @Test
+ public void testSyncReceiverLogger() throws IOException {
+ receiverLogger = new SyncReceiverLogger(
+ new File(getReceiverFolderFile(), Constans.SYNC_LOG_NAME));
+ Set<String> deletedFileNames = new HashSet<>();
+ Set<String> deletedFileNamesTest = new HashSet<>();
+ receiverLogger.startSyncDeletedFilesName();
+ for (int i = 0; i < 200; i++) {
+ receiverLogger
+ .finishSyncDeletedFileName(new File(getReceiverFolderFile(),
"deleted" + i));
+ deletedFileNames
+ .add(new File(getReceiverFolderFile(), "deleted" +
i).getAbsolutePath());
+ }
+ Set<String> toBeSyncedFiles = new HashSet<>();
+ Set<String> toBeSyncedFilesTest = new HashSet<>();
+ receiverLogger.startSyncTsFiles();
+ for (int i = 0; i < 200; i++) {
+ receiverLogger
+ .finishSyncTsfile(new File(getReceiverFolderFile(), "new" + i));
+ toBeSyncedFiles
+ .add(new File(getReceiverFolderFile(), "new" + i).getAbsolutePath());
+ }
+ int count = 0;
+ int mode = 0;
+ 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) {
+ deletedFileNamesTest.add(line);
+ } else if (mode == 1) {
+ toBeSyncedFilesTest.add(line);
+ }
+ }
+ }
+ }
+ assert count == 402;
+ assert deletedFileNames.size() == deletedFileNamesTest.size();
+ assert toBeSyncedFiles.size() == toBeSyncedFilesTest.size();
+ assert deletedFileNames.containsAll(deletedFileNamesTest);
+ assert toBeSyncedFiles.containsAll(toBeSyncedFilesTest);
+ }
+
+ private File getReceiverFolderFile() {
+ return new File(dataDir + File.separatorChar + Constans.SYNC_RECEIVER +
File.separatorChar
+ + "127.0.0.1_5555");
+ }
+}
\ No newline at end of file