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