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) {