This is an automated email from the ASF dual-hosted git repository.

qiaojialin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 6c341db  Fix data cannot be found when restarting server in HDFS (#412)
6c341db is described below

commit 6c341db333a1c3b7c65c67caf9e9cbdd2ebd706b
Author: Zesong Sun <[email protected]>
AuthorDate: Mon Sep 23 19:19:01 2019 +0800

    Fix data cannot be found when restarting server in HDFS (#412)
---
 .../engine/storagegroup/StorageGroupProcessor.java |  6 +-
 .../iotdb/db/tools/TsFileResourcePrinter.java      |  3 +-
 .../iotdb/db/writelog/recover/LogReplayer.java     |  4 +-
 .../iotdb/tsfile/fileSystem/FileInputFactory.java  |  2 +-
 .../iotdb/tsfile/fileSystem/FileOutputFactory.java |  5 +-
 .../apache/iotdb/tsfile/fileSystem/HDFSFile.java   | 59 +++++++--------
 .../apache/iotdb/tsfile/fileSystem/HDFSOutput.java | 25 +++++--
 .../iotdb/tsfile/fileSystem/TSFileFactory.java     | 84 +++++++++++++++++-----
 8 files changed, 122 insertions(+), 66 deletions(-)

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 3505959..5377adf 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
@@ -253,8 +253,8 @@ public class StorageGroupProcessor {
       // the process was interrupted before the merged files could be named
       continueFailedRenames(fileFolder, MERGE_SUFFIX);
 
-      Collections
-          .addAll(tsFiles, fileFolder.listFiles(file -> 
file.getName().endsWith(TSFILE_SUFFIX)));
+      Collections.addAll(tsFiles,
+          
TSFileFactory.INSTANCE.listFilesBySuffix(fileFolder.getAbsolutePath(), 
TSFILE_SUFFIX));
     }
     tsFiles.sort(this::compareFileName);
     List<TsFileResource> ret = new ArrayList<>();
@@ -263,7 +263,7 @@ public class StorageGroupProcessor {
   }
 
   private void continueFailedRenames(File fileFolder, String suffix) {
-    File[] files = fileFolder.listFiles(file -> 
file.getName().endsWith(suffix));
+    File[] files = 
TSFileFactory.INSTANCE.listFilesBySuffix(fileFolder.getAbsolutePath(), suffix);
     if (files != null) {
       for (File tempResource : files) {
         File originResource = 
TSFileFactory.INSTANCE.getFile(tempResource.getPath().replace(suffix, ""));
diff --git 
a/server/src/main/java/org/apache/iotdb/db/tools/TsFileResourcePrinter.java 
b/server/src/main/java/org/apache/iotdb/db/tools/TsFileResourcePrinter.java
index f26268b..4dafbdb 100644
--- a/server/src/main/java/org/apache/iotdb/db/tools/TsFileResourcePrinter.java
+++ b/server/src/main/java/org/apache/iotdb/db/tools/TsFileResourcePrinter.java
@@ -27,6 +27,7 @@ import java.util.Comparator;
 import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
 import org.apache.iotdb.db.qp.constant.DatetimeUtils;
 import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
+import org.apache.iotdb.tsfile.fileSystem.TSFileFactory;
 
 /**
  * this tool can analyze the tsfile.resource files from a folder.
@@ -40,7 +41,7 @@ public class TsFileResourcePrinter {
       folder = args[0];
     }
     File folderFile = SystemFileFactory.INSTANCE.getFile(folder);
-    File[] files = folderFile.listFiles(file -> 
file.getName().endsWith(".tsfile.resource"));
+    File[] files = 
TSFileFactory.INSTANCE.listFilesBySuffix(folderFile.getAbsolutePath(), 
".tsfile.resource");
     Arrays.sort(files, Comparator.comparingLong(x -> 
Long.valueOf(x.getName().split("-")[0])));
 
     for (File file : files) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/writelog/recover/LogReplayer.java 
b/server/src/main/java/org/apache/iotdb/db/writelog/recover/LogReplayer.java
index 115a304..fd39356 100644
--- a/server/src/main/java/org/apache/iotdb/db/writelog/recover/LogReplayer.java
+++ b/server/src/main/java/org/apache/iotdb/db/writelog/recover/LogReplayer.java
@@ -24,7 +24,6 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 
-import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
 import org.apache.iotdb.db.engine.memtable.IMemTable;
 import org.apache.iotdb.db.engine.modification.Deletion;
 import org.apache.iotdb.db.engine.modification.ModificationFile;
@@ -40,6 +39,7 @@ import org.apache.iotdb.db.writelog.io.ILogReader;
 import org.apache.iotdb.db.writelog.manager.MultiFileLogNodeManager;
 import org.apache.iotdb.db.writelog.node.WriteLogNode;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.fileSystem.TSFileFactory;
 import org.apache.iotdb.tsfile.read.common.Path;
 import org.apache.iotdb.tsfile.write.schema.Schema;
 
@@ -86,7 +86,7 @@ public class LogReplayer {
    */
   public void replayLogs() throws ProcessorException {
     WriteLogNode logNode = MultiFileLogNodeManager.getInstance().getNode(
-        logNodePrefix + 
SystemFileFactory.INSTANCE.getFile(insertFilePath).getName());
+        logNodePrefix + 
TSFileFactory.INSTANCE.getFile(insertFilePath).getName());
 
     ILogReader logReader = logNode.getLogReader();
     try {
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileInputFactory.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileInputFactory.java
index edb2858..efb5927 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileInputFactory.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileInputFactory.java
@@ -34,7 +34,7 @@ public enum FileInputFactory {
   INSTANCE;
 
   private static FSType fsType = 
TSFileDescriptor.getInstance().getConfig().getTSFileStorageFs();
-  private static final Logger logger = 
LoggerFactory.getLogger(TsFileIOWriter.class);
+  private static final Logger logger = 
LoggerFactory.getLogger(FileInputFactory.class);
 
   public TsFileInput getTsFileInput(String filePath) {
     try {
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileOutputFactory.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileOutputFactory.java
index 0853d20..621a04a 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileOutputFactory.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/FileOutputFactory.java
@@ -21,7 +21,6 @@ package org.apache.iotdb.tsfile.fileSystem;
 
 import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
 import org.apache.iotdb.tsfile.write.writer.DefaultTsFileOutput;
-import org.apache.iotdb.tsfile.write.writer.TsFileIOWriter;
 import org.apache.iotdb.tsfile.write.writer.TsFileOutput;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -34,12 +33,12 @@ public enum FileOutputFactory {
   INSTANCE;
 
   private static FSType fsType = 
TSFileDescriptor.getInstance().getConfig().getTSFileStorageFs();
-  private static final Logger logger = 
LoggerFactory.getLogger(TsFileIOWriter.class);
+  private static final Logger logger = 
LoggerFactory.getLogger(FileOutputFactory.class);
 
   public TsFileOutput getTsFileOutput(String filePath, boolean append) {
     try {
       if (fsType.equals(FSType.HDFS)) {
-        return new HDFSOutput(filePath, append);
+        return new HDFSOutput(filePath, !append);
       } else {
         return new DefaultTsFileOutput(new FileOutputStream(filePath, append));
       }
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSFile.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSFile.java
index d827b7a..6dc411a 100644
--- a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSFile.java
+++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSFile.java
@@ -19,12 +19,6 @@
 
 package org.apache.iotdb.tsfile.fileSystem;
 
-import org.apache.hadoop.conf.Configuration;
-import org.apache.hadoop.fs.*;
-import org.apache.iotdb.tsfile.write.TsFileWriter;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
 import java.io.File;
 import java.io.FileFilter;
 import java.io.FilenameFilter;
@@ -33,12 +27,19 @@ import java.net.MalformedURLException;
 import java.net.URI;
 import java.net.URL;
 import java.util.ArrayList;
+import java.util.List;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 public class HDFSFile extends File {
 
   private Path hdfsPath;
   private FileSystem fs;
-  private static final Logger logger = 
LoggerFactory.getLogger(TsFileWriter.class);
+  private static final Logger logger = LoggerFactory.getLogger(HDFSFile.class);
 
 
   public HDFSFile(String pathname) {
@@ -49,13 +50,13 @@ public class HDFSFile extends File {
 
   public HDFSFile(String parent, String child) {
     super(parent, child);
-    hdfsPath = new Path(parent + child);
+    hdfsPath = new Path(parent + File.separator + child);
     setConfAndGetFS();
   }
 
   public HDFSFile(File parent, String child) {
     super(parent, child);
-    hdfsPath = new Path(parent.getAbsolutePath() + child);
+    hdfsPath = new Path(parent.getAbsolutePath() + File.separator + child);
     setConfAndGetFS();
   }
 
@@ -68,6 +69,8 @@ public class HDFSFile extends File {
   private void setConfAndGetFS() {
     Configuration conf = new Configuration();
     conf.set("fs.hdfs.impl", "org.apache.hadoop.hdfs.DistributedFileSystem");
+    conf.set("dfs.client.block.write.replace-datanode-on-failure.policy", 
"NEVER");
+    conf.set("dfs.client.block.write.replace-datanode-on-failure.enable", 
"true");
     try {
       fs = hdfsPath.getFileSystem(conf);
     } catch (IOException e) {
@@ -107,33 +110,11 @@ public class HDFSFile extends File {
 
   @Override
   public File[] listFiles() {
-    ArrayList<HDFSFile> files = new ArrayList<>();
-    try {
-      RemoteIterator<LocatedFileStatus> iterator = fs.listFiles(hdfsPath, 
true);
-      while (iterator.hasNext()) {
-        LocatedFileStatus fileStatus = iterator.next();
-        Path fullPath = fileStatus.getPath();
-        files.add(new HDFSFile(fullPath.toUri()));
-      }
-      return files.toArray(new HDFSFile[files.size()]);
-    } catch (IOException e) {
-      logger.error("Fail to list files in {}. ", hdfsPath.toUri().toString(), 
e);
-      return null;
-    }
-  }
-
-  @Override
-  public File[] listFiles(FileFilter filter) {
-    ArrayList<HDFSFile> files = new ArrayList<>();
+    List<HDFSFile> files = new ArrayList<>();
     try {
-      PathFilter pathFilter = new GlobFilter(filter.toString()); // TODO test 
this filter in the future
-      RemoteIterator<LocatedFileStatus> iterator = fs.listFiles(hdfsPath, 
true);
-      while (iterator.hasNext()) {
-        LocatedFileStatus fileStatus = iterator.next();
-        Path fullPath = fileStatus.getPath();
-        if (pathFilter.accept(fullPath)) {
-          files.add(new HDFSFile(fullPath.toUri()));
-        }
+      for (FileStatus fileStatus : fs.listStatus(hdfsPath)) {
+        Path filePath = fileStatus.getPath();
+        files.add(new HDFSFile(filePath.toUri().toString()));
       }
       return files.toArray(new HDFSFile[files.size()]);
     } catch (IOException e) {
@@ -165,6 +146,9 @@ public class HDFSFile extends File {
   @Override
   public boolean mkdirs() {
     try {
+      if (exists()) {
+        return false;
+      }
       return fs.mkdirs(hdfsPath);
     } catch (IOException e) {
       logger.error("Fail to create directory {}. ", 
hdfsPath.toUri().toString(), e);
@@ -247,6 +231,11 @@ public class HDFSFile extends File {
   }
 
   @Override
+  public File[] listFiles(FileFilter filter) {
+    throw new UnsupportedOperationException("Unsupported operation.");
+  }
+
+  @Override
   public File getAbsoluteFile() {
     throw new UnsupportedOperationException("Unsupported operation.");
   }
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSOutput.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSOutput.java
index 7244758..56b0177 100644
--- a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSOutput.java
+++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/HDFSOutput.java
@@ -26,6 +26,8 @@ import org.apache.hadoop.fs.FSDataOutputStream;
 import org.apache.hadoop.fs.FileSystem;
 import org.apache.hadoop.fs.Path;
 import org.apache.iotdb.tsfile.write.writer.TsFileOutput;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 
 /**
@@ -37,6 +39,7 @@ public class HDFSOutput implements TsFileOutput {
   private FSDataOutputStream fsDataOutputStream;
   private FileSystem fs;
   private Path path;
+  private static final Logger logger = 
LoggerFactory.getLogger(HDFSOutput.class);
 
   public HDFSOutput(String filePath, boolean overwrite) throws IOException {
     this(filePath, new Configuration(), overwrite);
@@ -44,16 +47,19 @@ public class HDFSOutput implements TsFileOutput {
   }
 
 
-  public HDFSOutput(String filePath, Configuration configuration, boolean 
overwriter)
+  public HDFSOutput(String filePath, Configuration configuration, boolean 
overwrite)
       throws IOException {
-    this(new Path(filePath), configuration, overwriter);
+    this(new Path(filePath), configuration, overwrite);
     path = new Path(filePath);
   }
 
-  public HDFSOutput(Path path, Configuration configuration, boolean overwriter)
+  public HDFSOutput(Path path, Configuration configuration, boolean overwrite)
       throws IOException {
     fs = path.getFileSystem(configuration);
-    fsDataOutputStream = fs.create(path, overwriter);
+    configuration.set("fs.hdfs.impl", 
"org.apache.hadoop.hdfs.DistributedFileSystem");
+    
configuration.set("dfs.client.block.write.replace-datanode-on-failure.policy", 
"NEVER");
+    
configuration.set("dfs.client.block.write.replace-datanode-on-failure.enable", 
"true");
+    fsDataOutputStream = fs.exists(path) ? fs.append(path) : fs.create(path, 
overwrite);
     this.path = path;
   }
 
@@ -73,6 +79,7 @@ public class HDFSOutput implements TsFileOutput {
 
   @Override
   public void close() throws IOException {
+    flush();
     fsDataOutputStream.close();
   }
 
@@ -88,6 +95,14 @@ public class HDFSOutput implements TsFileOutput {
 
   @Override
   public void truncate(long position) throws IOException {
-    fs.truncate(path, position);
+    if (fs.exists(path)) {
+      fsDataOutputStream.close();
+    }
+    if (!fs.truncate(path, position)) {
+      logger.error("Failed to truncate file {}. ", path.toUri().toString());
+    }
+    if (fs.exists(path)) {
+      fsDataOutputStream = fs.append(path);
+    }
   }
 }
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/TSFileFactory.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/TSFileFactory.java
index 5a62802..644c3b6 100644
--- a/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/TSFileFactory.java
+++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/fileSystem/TSFileFactory.java
@@ -19,29 +19,42 @@
 
 package org.apache.iotdb.tsfile.fileSystem;
 
+import java.io.BufferedInputStream;
+import java.io.BufferedOutputStream;
+import java.io.BufferedReader;
+import java.io.BufferedWriter;
+import java.io.File;
+import java.io.FileInputStream;
+import java.io.FileOutputStream;
+import java.io.FileReader;
+import java.io.FileWriter;
+import java.io.IOException;
+import java.io.InputStreamReader;
+import java.io.OutputStreamWriter;
+import java.net.URI;
+import java.util.ArrayList;
+import java.util.List;
 import org.apache.commons.io.FileUtils;
 import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileStatus;
 import org.apache.hadoop.fs.FileSystem;
 import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.fs.PathFilter;
 import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
-import org.apache.iotdb.tsfile.write.TsFileWriter;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import java.io.*;
-import java.net.URI;
-
 public enum TSFileFactory {
 
   INSTANCE;
 
   private static FSType fSType = 
TSFileDescriptor.getInstance().getConfig().getTSFileStorageFs();
-  private static final Logger logger = 
LoggerFactory.getLogger(TsFileWriter.class);
+  private static final Logger logger = 
LoggerFactory.getLogger(TSFileFactory.class);
   private FileSystem fs;
   private Configuration conf = new Configuration();
 
   public File getFile(String pathname) {
-    if (fSType.equals(fSType.HDFS)) {
+    if (fSType.equals(FSType.HDFS)) {
       return new HDFSFile(pathname);
     } else {
       return new File(pathname);
@@ -49,7 +62,7 @@ public enum TSFileFactory {
   }
 
   public File getFile(String parent, String child) {
-    if (fSType.equals(fSType.HDFS)) {
+    if (fSType.equals(FSType.HDFS)) {
       return new HDFSFile(parent, child);
     } else {
       return new File(parent, child);
@@ -57,7 +70,7 @@ public enum TSFileFactory {
   }
 
   public File getFile(File parent, String child) {
-    if (fSType.equals(fSType.HDFS)) {
+    if (fSType.equals(FSType.HDFS)) {
       return new HDFSFile(parent, child);
     } else {
       return new File(parent, child);
@@ -65,7 +78,7 @@ public enum TSFileFactory {
   }
 
   public File getFile(URI uri) {
-    if (fSType.equals(fSType.HDFS)) {
+    if (fSType.equals(FSType.HDFS)) {
       return new HDFSFile(uri);
     } else {
       return new File(uri);
@@ -74,7 +87,7 @@ public enum TSFileFactory {
 
   public BufferedReader getBufferedReader(String filePath) {
     try {
-      if (fSType.equals(fSType.HDFS)) {
+      if (fSType.equals(FSType.HDFS)) {
         Path path = new Path(filePath);
         fs = path.getFileSystem(conf);
         return new BufferedReader(new InputStreamReader(fs.open(path)));
@@ -89,7 +102,7 @@ public enum TSFileFactory {
 
   public BufferedWriter getBufferedWriter(String filePath, boolean append) {
     try {
-      if (fSType.equals(fSType.HDFS)) {
+      if (fSType.equals(FSType.HDFS)) {
         Path path = new Path(filePath);
         fs = path.getFileSystem(conf);
         return new BufferedWriter(new OutputStreamWriter(fs.create(path)));
@@ -104,7 +117,7 @@ public enum TSFileFactory {
 
   public BufferedInputStream getBufferedInputStream(String filePath) {
     try {
-      if (fSType.equals(fSType.HDFS)) {
+      if (fSType.equals(FSType.HDFS)) {
         Path path = new Path(filePath);
         fs = path.getFileSystem(conf);
         return new BufferedInputStream(fs.open(path));
@@ -119,7 +132,7 @@ public enum TSFileFactory {
 
   public BufferedOutputStream getBufferedOutputStream(String filePath) {
     try {
-      if (fSType.equals(fSType.HDFS)) {
+      if (fSType.equals(FSType.HDFS)) {
         Path path = new Path(filePath);
         fs = path.getFileSystem(conf);
         return new BufferedOutputStream(fs.create(path));
@@ -134,16 +147,55 @@ public enum TSFileFactory {
 
   public void moveFile(File srcFile, File destFile) {
     try {
-      if (fSType.equals(fSType.HDFS)) {
+      if (fSType.equals(FSType.HDFS)) {
         boolean rename = srcFile.renameTo(destFile);
         if (!rename) {
-          logger.error("Failed to rename file from {} to {}. ", 
srcFile.getName(), destFile.getName());
+          logger.error("Failed to rename file from {} to {}. ", 
srcFile.getName(),
+              destFile.getName());
         }
       } else {
         FileUtils.moveFile(srcFile, destFile);
       }
     } catch (IOException e) {
-      logger.error("Failed to move file from {} to {}. ", 
srcFile.getAbsolutePath(), destFile.getAbsolutePath(), e);
+      logger.error("Failed to move file from {} to {}. ", 
srcFile.getAbsolutePath(),
+          destFile.getAbsolutePath(), e);
+    }
+  }
+
+  public File[] listFilesBySuffix(String fileFolder, String suffix) {
+    if (fSType.equals(FSType.HDFS)) {
+      PathFilter pathFilter = path -> path.toUri().toString().endsWith(suffix);
+      List<HDFSFile> files = listFiles(fileFolder, pathFilter);
+      return files.toArray(new HDFSFile[files.size()]);
+    } else {
+      return new File(fileFolder).listFiles(file -> 
file.getName().endsWith(suffix));
+    }
+  }
+
+  public File[] listFilesByPrefix(String fileFolder, String prefix) {
+    if (fSType.equals(FSType.HDFS)) {
+      PathFilter pathFilter = path -> 
path.toUri().toString().startsWith(prefix);
+      List<HDFSFile> files = listFiles(fileFolder, pathFilter);
+      return files.toArray(new HDFSFile[files.size()]);
+    } else {
+      return new File(fileFolder).listFiles(file -> 
file.getName().startsWith(prefix));
+    }
+  }
+
+  private List<HDFSFile> listFiles(String fileFolder, PathFilter pathFilter) {
+    List<HDFSFile> files = new ArrayList<>();
+    try {
+      Path path = new Path(fileFolder);
+      fs = path.getFileSystem(conf);
+      for (FileStatus fileStatus: fs.listStatus(path)) {
+        Path filePath = fileStatus.getPath();
+        if (pathFilter.accept(filePath)) {
+          files.add(new HDFSFile(filePath.toUri().toString()));
+        }
+      }
+    } catch (IOException e) {
+      logger.error("Failed to list files in {}. ", fileFolder);
     }
+    return files;
   }
 }
\ No newline at end of file

Reply via email to