This is an automated email from the ASF dual-hosted git repository.
haonan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 462c0c2c8c [IOTDB-3438] Fix empty WAL file (#6312)
462c0c2c8c is described below
commit 462c0c2c8c5d98f3369ed692b653d70a5d7ce7bd
Author: Alan Choo <[email protected]>
AuthorDate: Fri Jun 17 10:18:37 2022 +0800
[IOTDB-3438] Fix empty WAL file (#6312)
Co-authored-by: Haonan <[email protected]>
---
.../iotdb/db/engine/storagegroup/DataRegion.java | 13 ++++++++--
.../iotdb/db/wal/buffer/AbstractWALBuffer.java | 8 +++++-
.../org/apache/iotdb/db/wal/buffer/IWALBuffer.java | 3 +++
.../org/apache/iotdb/db/wal/buffer/WALBuffer.java | 30 ++++++++++++----------
.../org/apache/iotdb/db/wal/io/ILogWriter.java | 3 +--
.../java/org/apache/iotdb/db/wal/io/LogWriter.java | 3 +--
.../java/org/apache/iotdb/db/wal/node/WALNode.java | 19 ++++++++++++--
7 files changed, 57 insertions(+), 22 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java
index 4ce9ab0064..f1a79897ab 100755
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java
@@ -509,8 +509,17 @@ public class DataRegion {
// recover and start timed compaction thread
initCompaction();
- logger.info(
- "The data region {}[{}] is recovered successfully",
logicalStorageGroupName, dataRegionId);
+ if (config.isMppMode()
+ ? StorageEngineV2.getInstance().isAllSgReady()
+ : StorageEngine.getInstance().isAllSgReady()) {
+ logger.info(
+ "The data region {}[{}] is created successfully",
logicalStorageGroupName, dataRegionId);
+ } else {
+ logger.info(
+ "The data region {}[{}] is recovered successfully",
+ logicalStorageGroupName,
+ dataRegionId);
+ }
}
private void initCompaction() {
diff --git
a/server/src/main/java/org/apache/iotdb/db/wal/buffer/AbstractWALBuffer.java
b/server/src/main/java/org/apache/iotdb/db/wal/buffer/AbstractWALBuffer.java
index 8cc27f9dfb..16002896b8 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/buffer/AbstractWALBuffer.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/buffer/AbstractWALBuffer.java
@@ -52,7 +52,7 @@ public abstract class AbstractWALBuffer implements IWALBuffer
{
this.logDirectory = logDirectory;
File logDirFile = SystemFileFactory.INSTANCE.getFile(logDirectory);
if (!logDirFile.exists() && logDirFile.mkdirs()) {
- logger.info("create folder {} for wal buffer-{}.", logDirectory,
identifier);
+ logger.info("Create folder {} for wal node-{}'s buffer.", logDirectory,
identifier);
}
currentSearchIndex = startSearchIndex;
currentWALFileVersion.set(startFileVersion);
@@ -68,6 +68,11 @@ public abstract class AbstractWALBuffer implements
IWALBuffer {
return currentWALFileVersion.get();
}
+ @Override
+ public long getCurrentWALFileSize() {
+ return currentWALFileWriter.size();
+ }
+
/** Notice: only called by syncBufferThread and old log writer will be
closed by this function. */
protected void rollLogWriter(long searchIndex) throws IOException {
currentWALFileWriter.close();
@@ -76,5 +81,6 @@ public abstract class AbstractWALBuffer implements IWALBuffer
{
logDirectory,
WALFileUtils.getLogFileName(currentWALFileVersion.incrementAndGet(),
searchIndex));
currentWALFileWriter = new WALWriter(nextLogFile);
+ logger.debug("Open new wal file {} for wal node-{}'s buffer.",
nextLogFile, identifier);
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/wal/buffer/IWALBuffer.java
b/server/src/main/java/org/apache/iotdb/db/wal/buffer/IWALBuffer.java
index ee57e51b2a..42a7946b34 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/buffer/IWALBuffer.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/buffer/IWALBuffer.java
@@ -37,6 +37,9 @@ public interface IWALBuffer extends AutoCloseable {
/** Get current log version id */
int getCurrentWALFileVersion();
+ /** Get current wal file's size */
+ long getCurrentWALFileSize();
+
@Override
void close();
diff --git a/server/src/main/java/org/apache/iotdb/db/wal/buffer/WALBuffer.java
b/server/src/main/java/org/apache/iotdb/db/wal/buffer/WALBuffer.java
index d0e2098c76..ddc77cf6da 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/buffer/WALBuffer.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/buffer/WALBuffer.java
@@ -232,10 +232,15 @@ public class WALBuffer extends AbstractWALBuffer {
private boolean handleSignalEntry(SignalWALEntry signalWALEntry) {
switch (signalWALEntry.getSignalType()) {
case ROLL_WAL_LOG_WRITER_SIGNAL:
+ logger.debug("Handle roll log writer signal for wal node-{}.",
identifier);
rollWALFileWriterListener = signalWALEntry.getWalFlushListener();
fsyncWorkingBuffer(currentSearchIndex, fsyncListeners,
rollWALFileWriterListener);
return true;
case CLOSE_SIGNAL:
+ logger.debug(
+ "Handle close signal for wal node-{}, there are {} entries
left.",
+ identifier,
+ walEntries.size());
boolean dataExists = batchSize > 0;
if (dataExists) {
fsyncWorkingBuffer(currentSearchIndex, fsyncListeners,
rollWALFileWriterListener);
@@ -417,24 +422,23 @@ public class WALBuffer extends AbstractWALBuffer {
}
// try to roll log writer
- try {
- if (rollWALFileWriterListener != null
- || (forceFlag
- && currentWALFileWriter.size() >=
config.getWalFileSizeThresholdInByte())) {
+ if (rollWALFileWriterListener != null
+ || (forceFlag && currentWALFileWriter.size() >=
config.getWalFileSizeThresholdInByte())) {
+ try {
rollLogWriter(searchIndex);
if (rollWALFileWriterListener != null) {
rollWALFileWriterListener.succeed();
}
+ } catch (IOException e) {
+ logger.error(
+ "Fail to roll wal node-{}'s log writer, change system mode to
read-only.",
+ identifier,
+ e);
+ if (rollWALFileWriterListener != null) {
+ rollWALFileWriterListener.fail(e);
+ }
+ config.setReadOnly(true);
}
- } catch (IOException e) {
- logger.error(
- "Fail to roll wal node-{}'s log writer, change system mode to
read-only.",
- identifier,
- e);
- if (rollWALFileWriterListener != null) {
- rollWALFileWriterListener.fail(e);
- }
- config.setReadOnly(true);
}
}
}
diff --git a/server/src/main/java/org/apache/iotdb/db/wal/io/ILogWriter.java
b/server/src/main/java/org/apache/iotdb/db/wal/io/ILogWriter.java
index 9d26c17fa4..932918d79c 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/io/ILogWriter.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/io/ILogWriter.java
@@ -55,7 +55,6 @@ public interface ILogWriter extends Closeable {
* Returns the current size of this file.
*
* @return size
- * @throws IOException if an I/O error occurs
*/
- long size() throws IOException;
+ long size();
}
diff --git a/server/src/main/java/org/apache/iotdb/db/wal/io/LogWriter.java
b/server/src/main/java/org/apache/iotdb/db/wal/io/LogWriter.java
index 2b4c7c8cb8..9b14388ab2 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/io/LogWriter.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/io/LogWriter.java
@@ -37,7 +37,6 @@ import java.nio.channels.FileChannel;
* and writing {@link Checkpoint} into .checkpoint file.
*/
public abstract class LogWriter implements ILogWriter {
- public static final String FILE_PREFIX = "_";
private static final Logger logger =
LoggerFactory.getLogger(LogWriter.class);
private final File logFile;
@@ -76,7 +75,7 @@ public abstract class LogWriter implements ILogWriter {
}
@Override
- public long size() throws IOException {
+ public long size() {
return size;
}
diff --git a/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
b/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
index 9c3a3241ab..41e3fb8514 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
@@ -221,7 +221,9 @@ public class WALNode implements IWALNode {
firstValidVersionId = checkpointManager.getFirstValidWALVersionId();
if (firstValidVersionId == Integer.MIN_VALUE) {
// roll wal log writer to delete current wal file
- rollWALFile();
+ if (buffer.getCurrentWALFileSize() > 0) {
+ rollWALFile();
+ }
// update firstValidVersionId
firstValidVersionId = checkpointManager.getFirstValidWALVersionId();
if (firstValidVersionId == Integer.MIN_VALUE) {
@@ -229,6 +231,12 @@ public class WALNode implements IWALNode {
}
}
+ logger.debug(
+ "Start deleting outdated wal files for wal node-{}, the first valid
version id is {}, and the safely deleted search index is {}.",
+ identifier,
+ firstValidVersionId,
+ safelyDeletedSearchIndex);
+
// delete outdated files
deleteOutdatedFiles();
@@ -262,10 +270,13 @@ public class WALNode implements IWALNode {
}
private void deleteOutdatedFiles() {
+ int deletedFilesNum = 0;
File[] filesToDelete = logDirectory.listFiles(this::filterFilesToDelete);
if (filesToDelete != null) {
for (File file : filesToDelete) {
- if (!file.delete()) {
+ if (file.delete()) {
+ deletedFilesNum++;
+ } else {
logger.info("Fail to delete outdated wal file {} of wal node-{}.",
file, identifier);
}
// update totalRamCostOfFlushedMemTables
@@ -276,6 +287,10 @@ public class WALNode implements IWALNode {
}
}
}
+ logger.debug(
+ "Successfully delete {} outdated wal files for wal node-{}.",
+ deletedFilesNum,
+ identifier);
}
private boolean filterFilesToDelete(File dir, String name) {