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

jiangtian pushed a commit to branch improve_wal
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/improve_wal by this push:
     new e932398  fix log loss
e932398 is described below

commit e932398da65a50ae53d816dae4d4adb95ccddd3c
Author: jt <[email protected]>
AuthorDate: Fri Oct 16 14:25:47 2020 +0800

    fix log loss
---
 .../org/apache/iotdb/db/engine/flush/MemTableFlushTask.java  |  1 -
 .../apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java   |  8 ++++++++
 .../src/main/java/org/apache/iotdb/db/utils/CommonUtils.java |  2 +-
 .../apache/iotdb/db/utils/datastructure/BinaryTVList.java    |  4 ++++
 .../apache/iotdb/db/utils/datastructure/BooleanTVList.java   |  4 ++++
 .../apache/iotdb/db/utils/datastructure/DoubleTVList.java    |  4 ++++
 .../org/apache/iotdb/db/utils/datastructure/FloatTVList.java |  4 ++++
 .../org/apache/iotdb/db/utils/datastructure/IntTVList.java   |  4 ++++
 .../org/apache/iotdb/db/utils/datastructure/LongTVList.java  |  4 ++++
 .../iotdb/db/writelog/io/DifferentialBatchLogReader.java     |  3 +--
 .../iotdb/db/writelog/node/DifferentialWriteLogNode.java     | 12 +++++++++---
 11 files changed, 43 insertions(+), 7 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java 
b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
index 3b80b13..b40ea1c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
@@ -90,7 +90,6 @@ public class MemTableFlushTask {
         TVList tvList = series.getSortedTVList();
         long seriesSortTime = System.currentTimeMillis() - startTime;
         sortTime += seriesSortTime;
-        logger.debug("{}.{} sort costs {}", deviceId, measurementId, 
seriesSortTime);
         encodingTaskQueue.add(new Pair<>(tvList, desc));
         // register active time series to the ActiveTimeSeriesCounter
         if 
(IoTDBDescriptor.getInstance().getConfig().isEnableParameterAdapter()) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java
 
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java
index 6ebe53f..310a557 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertTabletPlan.java
@@ -485,4 +485,12 @@ public class InsertTabletPlan extends InsertPlan {
     columns[index] = null;
   }
 
+  @Override
+  public String toString() {
+    return "InsertTabletPlan{" +
+        "maxTime=" + maxTime +
+        ", minTime=" + minTime +
+        ", deviceId=" + deviceId +
+        '}';
+  }
 }
diff --git a/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java 
b/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java
index 542252f..369520b 100644
--- a/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java
+++ b/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java
@@ -147,7 +147,7 @@ public class CommonUtils {
   public static void updatePlanWindow(PhysicalPlan plan, int windowLength,
       Queue<PhysicalPlan> planWindow) {
     if (planWindow.size() >= windowLength) {
-      planWindow.remove();
+      planWindow.poll();
     }
     // remove unnecessary fields as bases to reduce memory footprint
     if (plan instanceof InsertRowPlan) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BinaryTVList.java
 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BinaryTVList.java
index b1adc1e..2e9c37c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BinaryTVList.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BinaryTVList.java
@@ -93,6 +93,10 @@ public class BinaryTVList extends TVList {
   }
 
   public void sort() {
+    if (sorted) {
+      return;
+    }
+
     if (sortedTimestamps == null || sortedTimestamps.length < size) {
       sortedTimestamps = (long[][]) PrimitiveArrayPool
           .getInstance().getDataListsByType(TSDataType.INT64, size);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BooleanTVList.java
 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BooleanTVList.java
index ba7d262..e40ac5d 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BooleanTVList.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/BooleanTVList.java
@@ -92,6 +92,10 @@ public class BooleanTVList extends TVList {
   }
 
   public void sort() {
+    if (sorted) {
+      return;
+    }
+
     if (sortedTimestamps == null || sortedTimestamps.length < size) {
       sortedTimestamps = (long[][]) PrimitiveArrayPool
           .getInstance().getDataListsByType(TSDataType.INT64, size);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/DoubleTVList.java
 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/DoubleTVList.java
index e141b75..a20564d 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/DoubleTVList.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/DoubleTVList.java
@@ -93,6 +93,10 @@ public class DoubleTVList extends TVList {
   }
 
   public void sort() {
+    if (sorted) {
+      return;
+    }
+
     if (sortedTimestamps == null || sortedTimestamps.length < size) {
       sortedTimestamps = (long[][]) PrimitiveArrayPool
           .getInstance().getDataListsByType(TSDataType.INT64, size);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/FloatTVList.java 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/FloatTVList.java
index 178a4f7..bb95f62 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/FloatTVList.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/FloatTVList.java
@@ -93,6 +93,10 @@ public class FloatTVList extends TVList {
   }
 
   public void sort() {
+    if (sorted) {
+      return;
+    }
+
     if (sortedTimestamps == null || sortedTimestamps.length < size) {
       sortedTimestamps = (long[][]) PrimitiveArrayPool
           .getInstance().getDataListsByType(TSDataType.INT64, size);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/IntTVList.java 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/IntTVList.java
index 18c7c0f..7d09c80 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/IntTVList.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/IntTVList.java
@@ -92,6 +92,10 @@ public class IntTVList extends TVList {
   }
 
   public void sort() {
+    if (sorted) {
+      return;
+    }
+
     if (sortedTimestamps == null || sortedTimestamps.length < size) {
       sortedTimestamps = (long[][]) PrimitiveArrayPool
           .getInstance().getDataListsByType(TSDataType.INT64, size);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/LongTVList.java 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/LongTVList.java
index c5c198f..5e6a18f 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/LongTVList.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/LongTVList.java
@@ -92,6 +92,10 @@ public class LongTVList extends TVList {
   }
 
   public void sort() {
+    if (sorted) {
+      return;
+    }
+
     if (sortedTimestamps == null || sortedTimestamps.length < size) {
       sortedTimestamps = (long[][]) PrimitiveArrayPool
           .getInstance().getDataListsByType(TSDataType.INT64, size);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/writelog/io/DifferentialBatchLogReader.java
 
b/server/src/main/java/org/apache/iotdb/db/writelog/io/DifferentialBatchLogReader.java
index 8d765cf..1a87ab3 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/writelog/io/DifferentialBatchLogReader.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/writelog/io/DifferentialBatchLogReader.java
@@ -26,7 +26,6 @@ import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.List;
-import java.util.Queue;
 import org.apache.iotdb.db.exception.metadata.IllegalPathException;
 import org.apache.iotdb.db.qp.physical.PhysicalPlan;
 import org.apache.iotdb.db.qp.physical.PhysicalPlan.Factory;
@@ -52,8 +51,8 @@ public class DifferentialBatchLogReader extends 
BatchLogReader {
     while (buffer.position() != buffer.limit()) {
       try {
         PhysicalPlan physicalPlan = Factory.create(buffer, planWindow);
-        plans.add(physicalPlan);
         CommonUtils.updatePlanWindow(physicalPlan, WINDOW_LENGTH, planWindow);
+        plans.add(physicalPlan);
       } catch (IOException | IllegalPathException e) {
         logger.error("Cannot deserialize PhysicalPlans from ByteBuffer, ignore 
remaining logs", e);
         fileCorrupted = true;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/writelog/node/DifferentialWriteLogNode.java
 
b/server/src/main/java/org/apache/iotdb/db/writelog/node/DifferentialWriteLogNode.java
index 803a378..692c6be 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/writelog/node/DifferentialWriteLogNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/writelog/node/DifferentialWriteLogNode.java
@@ -76,9 +76,15 @@ public class DifferentialWriteLogNode extends 
ExclusiveWriteLogNode {
   }
 
   @Override
-  void nextFileWriter() {
-    super.nextFileWriter();
-    planWindow.clear();
+  public void notifyStartFlush() {
+    lock.writeLock().lock();
+    try {
+      close();
+      nextFileWriter();
+      planWindow.clear();
+    } finally {
+      lock.writeLock().unlock();
+    }
   }
 
   private void serialize(PhysicalPlan plan, Pair<PhysicalPlan, Short> 
similarPlanIndex) {

Reply via email to