This is an automated email from the ASF dual-hosted git repository.
chaow pushed a commit to branch rel/0.12
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/0.12 by this push:
new f3de42b [To rel/0.12][IOTDB-2027] Rollback invalid entry after wal
writing failure (#4418)
f3de42b is described below
commit f3de42b93007adcf5592de9d2033279d0a4b6a81
Author: Alan Choo <[email protected]>
AuthorDate: Fri Nov 19 10:57:35 2021 +0800
[To rel/0.12][IOTDB-2027] Rollback invalid entry after wal writing failure
(#4418)
---
.../org/apache/iotdb/db/mqtt/PublishHandler.java | 16 +++----
.../apache/iotdb/db/qp/physical/PhysicalPlan.java | 22 ++++++++-
.../db/qp/physical/crud/CreateTemplatePlan.java | 2 +-
.../iotdb/db/qp/physical/crud/DeletePlan.java | 2 +-
.../db/qp/physical/crud/InsertMultiTabletPlan.java | 2 +-
.../iotdb/db/qp/physical/crud/InsertRowPlan.java | 4 +-
.../physical/crud/InsertRowsOfOneDevicePlan.java | 2 +-
.../iotdb/db/qp/physical/crud/InsertRowsPlan.java | 2 +-
.../db/qp/physical/crud/InsertTabletPlan.java | 2 +-
.../db/qp/physical/crud/SetDeviceTemplatePlan.java | 2 +-
.../iotdb/db/qp/physical/sys/AuthorPlan.java | 2 +-
.../qp/physical/sys/AutoCreateDeviceMNodePlan.java | 2 +-
.../iotdb/db/qp/physical/sys/ChangeAliasPlan.java | 2 +-
.../db/qp/physical/sys/ChangeTagOffsetPlan.java | 2 +-
.../iotdb/db/qp/physical/sys/CreateIndexPlan.java | 2 +-
.../qp/physical/sys/CreateMultiTimeSeriesPlan.java | 2 +-
.../db/qp/physical/sys/CreateTimeSeriesPlan.java | 2 +-
.../iotdb/db/qp/physical/sys/DataAuthPlan.java | 2 +-
.../db/qp/physical/sys/DeleteStorageGroupPlan.java | 2 +-
.../db/qp/physical/sys/DeleteTimeSeriesPlan.java | 2 +-
.../iotdb/db/qp/physical/sys/DropIndexPlan.java | 2 +-
.../apache/iotdb/db/qp/physical/sys/FlushPlan.java | 2 +-
.../apache/iotdb/db/qp/physical/sys/MNodePlan.java | 2 +-
.../db/qp/physical/sys/MeasurementMNodePlan.java | 2 +-
.../db/qp/physical/sys/SetStorageGroupPlan.java | 2 +-
.../db/qp/physical/sys/SetSystemModePlan.java | 2 +-
.../iotdb/db/qp/physical/sys/SetTTLPlan.java | 2 +-
.../physical/sys/SetUsingDeviceTemplatePlan.java | 2 +-
.../db/qp/physical/sys/StorageGroupMNodePlan.java | 2 +-
.../db/writelog/node/ExclusiveWriteLogNode.java | 2 -
.../iotdb/db/qp/physical/PhysicalPlanTest.java | 53 ++++++++++++----------
31 files changed, 85 insertions(+), 64 deletions(-)
diff --git a/server/src/main/java/org/apache/iotdb/db/mqtt/PublishHandler.java
b/server/src/main/java/org/apache/iotdb/db/mqtt/PublishHandler.java
index 1e4b213..e39c05f 100644
--- a/server/src/main/java/org/apache/iotdb/db/mqtt/PublishHandler.java
+++ b/server/src/main/java/org/apache/iotdb/db/mqtt/PublishHandler.java
@@ -27,7 +27,6 @@ import org.apache.iotdb.db.qp.executor.IPlanExecutor;
import org.apache.iotdb.db.qp.executor.PlanExecutor;
import org.apache.iotdb.db.qp.physical.PhysicalPlan;
import org.apache.iotdb.db.qp.physical.crud.InsertRowPlan;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import io.moquette.interception.AbstractInterceptHandler;
import io.moquette.interception.messages.InterceptPublishMessage;
@@ -93,16 +92,15 @@ public class PublishHandler extends
AbstractInterceptHandler {
continue;
}
- InsertRowPlan plan = new InsertRowPlan();
- plan.setTime(event.getTimestamp());
- plan.setMeasurements(event.getMeasurements().toArray(new String[0]));
- plan.setValues(event.getValues().toArray(new Object[0]));
- plan.setDataTypes(new TSDataType[event.getValues().size()]);
- plan.setNeedInferType(true);
-
boolean status = false;
try {
- plan.setDeviceId(new PartialPath(event.getDevice()));
+ PartialPath path = new PartialPath(event.getDevice());
+ InsertRowPlan plan =
+ new InsertRowPlan(
+ path,
+ event.getTimestamp(),
+ event.getMeasurements().toArray(new String[0]),
+ event.getValues().toArray(new String[0]));
status = executeNonQuery(plan);
} catch (Exception e) {
LOG.warn(
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java
index f87cc05..3b7790d 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/PhysicalPlan.java
@@ -58,6 +58,9 @@ import org.apache.iotdb.db.qp.physical.sys.ShowTimeSeriesPlan;
import org.apache.iotdb.db.qp.physical.sys.StorageGroupMNodePlan;
import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import java.io.DataOutputStream;
import java.io.IOException;
import java.nio.ByteBuffer;
@@ -67,6 +70,7 @@ import java.util.List;
/** This class is a abstract class for all type of PhysicalPlan. */
public abstract class PhysicalPlan {
+ private static final Logger logger =
LoggerFactory.getLogger(PhysicalPlan.class);
private static final String SERIALIZATION_UNIMPLEMENTED = "serialization
unimplemented";
@@ -147,11 +151,27 @@ public abstract class PhysicalPlan {
/**
* Serialize the plan into the given buffer. This is provided for WAL, so
fields that can be
- * recovered will not be serialized.
+ * recovered will not be serialized. If error occurs when serializing this
plan, the buffer will
+ * be reset.
*
* @param buffer
*/
public void serialize(ByteBuffer buffer) {
+ buffer.mark();
+ try {
+ serializeImpl(buffer);
+ } catch (UnsupportedOperationException e) {
+ // ignore and throw
+ throw e;
+ } catch (Exception e) {
+ logger.error(
+ "Rollback buffer entry because error occurs when serializing this
physical plan.", e);
+ buffer.reset();
+ throw e;
+ }
+ }
+
+ protected void serializeImpl(ByteBuffer buffer) {
throw new UnsupportedOperationException(SERIALIZATION_UNIMPLEMENTED);
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/CreateTemplatePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/CreateTemplatePlan.java
index 82c45b4..bfe8430 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/CreateTemplatePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/CreateTemplatePlan.java
@@ -111,7 +111,7 @@ public class CreateTemplatePlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.CREATE_TEMPLATE.ordinal());
ReadWriteIOUtils.write(name, buffer);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/DeletePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/DeletePlan.java
index f3c1b43..4b608f8 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/DeletePlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/DeletePlan.java
@@ -136,7 +136,7 @@ public class DeletePlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.DELETE.ordinal();
buffer.put((byte) type);
buffer.putLong(deleteStartTime);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertMultiTabletPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertMultiTabletPlan.java
index 2ac9504..a744365 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertMultiTabletPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertMultiTabletPlan.java
@@ -256,7 +256,7 @@ public class InsertMultiTabletPlan extends InsertPlan
implements BatchPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.MULTI_BATCH_INSERT.ordinal();
buffer.put((byte) type);
buffer.putInt(insertTabletPlanList.size());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java
index a2f19bc..5faea8b 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowPlan.java
@@ -367,7 +367,7 @@ public class InsertRowPlan extends InsertPlan {
}
private void putValues(ByteBuffer buffer) throws QueryProcessException {
- for (int i = 0; i < values.length; i++) {
+ for (int i = 0; i < measurements.length; i++) {
if (measurements[i] == null) {
continue;
}
@@ -442,7 +442,7 @@ public class InsertRowPlan extends InsertPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.INSERT.ordinal();
buffer.put((byte) type);
subSerialize(buffer);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsOfOneDevicePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsOfOneDevicePlan.java
index 5f6a45e..138b902 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsOfOneDevicePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsOfOneDevicePlan.java
@@ -153,7 +153,7 @@ public class InsertRowsOfOneDevicePlan extends InsertPlan
implements BatchPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.BATCH_INSERT_ONE_DEVICE.ordinal();
buffer.put((byte) type);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsPlan.java
index b0e5550..019ae22 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/InsertRowsPlan.java
@@ -168,7 +168,7 @@ public class InsertRowsPlan extends InsertPlan implements
BatchPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.BATCH_INSERT_ROWS.ordinal();
buffer.put((byte) type);
buffer.putInt(insertRowPlanList.size());
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 fa75882..afb0783 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
@@ -211,7 +211,7 @@ public class InsertTabletPlan extends InsertPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.BATCHINSERT.ordinal();
buffer.put((byte) type);
subSerialize(buffer);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/SetDeviceTemplatePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/SetDeviceTemplatePlan.java
index b76703a..2feeb57 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/SetDeviceTemplatePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/crud/SetDeviceTemplatePlan.java
@@ -65,7 +65,7 @@ public class SetDeviceTemplatePlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.SET_DEVICE_TEMPLATE.ordinal());
ReadWriteIOUtils.write(templateName, buffer);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AuthorPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AuthorPlan.java
index 26c7d8a..694b2b7 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AuthorPlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AuthorPlan.java
@@ -325,7 +325,7 @@ public class AuthorPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = this.getPlanType(super.getOperatorType());
buffer.put((byte) type);
buffer.putInt(authorType.ordinal());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AutoCreateDeviceMNodePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AutoCreateDeviceMNodePlan.java
index ef7412e..eaccf1c 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AutoCreateDeviceMNodePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/AutoCreateDeviceMNodePlan.java
@@ -61,7 +61,7 @@ public class AutoCreateDeviceMNodePlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.AUTO_CREATE_DEVICE_MNODE.ordinal());
putString(buffer, path.getFullPath());
buffer.putLong(index);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeAliasPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeAliasPlan.java
index a6bf1aa..096f8c9 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeAliasPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeAliasPlan.java
@@ -71,7 +71,7 @@ public class ChangeAliasPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.CHANGE_ALIAS.ordinal();
buffer.put((byte) type);
putString(buffer, path.getFullPath());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeTagOffsetPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeTagOffsetPlan.java
index ba80502..56eab5e 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeTagOffsetPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ChangeTagOffsetPlan.java
@@ -71,7 +71,7 @@ public class ChangeTagOffsetPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.CHANGE_TAG_OFFSET.ordinal();
buffer.put((byte) type);
putString(buffer, path.getFullPath());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateIndexPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateIndexPlan.java
index c606a41..009e171 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateIndexPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateIndexPlan.java
@@ -112,7 +112,7 @@ public class CreateIndexPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.CREATE_INDEX.ordinal();
buffer.put((byte) type);
buffer.put((byte) indexType.serialize());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateMultiTimeSeriesPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateMultiTimeSeriesPlan.java
index b98761f..863cbcb 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateMultiTimeSeriesPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateMultiTimeSeriesPlan.java
@@ -221,7 +221,7 @@ public class CreateMultiTimeSeriesPlan extends PhysicalPlan
implements BatchPlan
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.CREATE_MULTI_TIMESERIES.ordinal();
buffer.put((byte) type);
buffer.putInt(paths.size());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateTimeSeriesPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateTimeSeriesPlan.java
index a5d7bd5..28f500d 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateTimeSeriesPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/CreateTimeSeriesPlan.java
@@ -208,7 +208,7 @@ public class CreateTimeSeriesPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.CREATE_TIMESERIES.ordinal());
byte[] bytes = path.getFullPath().getBytes();
buffer.putInt(bytes.length);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DataAuthPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DataAuthPlan.java
index d52b53c..b415c6c 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DataAuthPlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DataAuthPlan.java
@@ -65,7 +65,7 @@ public class DataAuthPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = this.getPlanType(super.getOperatorType());
buffer.put((byte) type);
buffer.putInt(users.size());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteStorageGroupPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteStorageGroupPlan.java
index fecab5c..0918040 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteStorageGroupPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteStorageGroupPlan.java
@@ -60,7 +60,7 @@ public class DeleteStorageGroupPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.DELETE_STORAGE_GROUP.ordinal();
buffer.put((byte) type);
buffer.putInt(this.getPaths().size());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteTimeSeriesPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteTimeSeriesPlan.java
index d0b1677..3796ace 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteTimeSeriesPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DeleteTimeSeriesPlan.java
@@ -69,7 +69,7 @@ public class DeleteTimeSeriesPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.DELETE_TIMESERIES.ordinal();
buffer.put((byte) type);
buffer.putInt(deletePathList.size());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DropIndexPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DropIndexPlan.java
index 0126153..5df5116 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DropIndexPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/DropIndexPlan.java
@@ -79,7 +79,7 @@ public class DropIndexPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.DROP_INDEX.ordinal();
buffer.put((byte) type);
buffer.put((byte) indexType.serialize());
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/FlushPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/FlushPlan.java
index 0be0095..f625940 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/FlushPlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/FlushPlan.java
@@ -148,7 +148,7 @@ public class FlushPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.FLUSH.ordinal();
buffer.put((byte) type);
if (isSeq == null) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MNodePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MNodePlan.java
index d0ab19c..4a1055c 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MNodePlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MNodePlan.java
@@ -70,7 +70,7 @@ public class MNodePlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.MNODE.ordinal());
putString(buffer, name);
buffer.putInt(childSize);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MeasurementMNodePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MeasurementMNodePlan.java
index cee286a..d28977a 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MeasurementMNodePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/MeasurementMNodePlan.java
@@ -55,7 +55,7 @@ public class MeasurementMNodePlan extends MNodePlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.MEASUREMENT_MNODE.ordinal());
putString(buffer, name);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetStorageGroupPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetStorageGroupPlan.java
index 9f20dd2..443b039 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetStorageGroupPlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetStorageGroupPlan.java
@@ -64,7 +64,7 @@ public class SetStorageGroupPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.SET_STORAGE_GROUP.ordinal());
putString(buffer, path.getFullPath());
buffer.putLong(index);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetSystemModePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetSystemModePlan.java
index 208f8d6..0d41a22 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetSystemModePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetSystemModePlan.java
@@ -61,7 +61,7 @@ public class SetSystemModePlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.SET_SYSTEM_MODE.ordinal());
ReadWriteIOUtils.write(isReadOnly, buffer);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetTTLPlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetTTLPlan.java
index 249e6e7..73a0f8a 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetTTLPlan.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetTTLPlan.java
@@ -67,7 +67,7 @@ public class SetTTLPlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
int type = PhysicalPlanType.TTL.ordinal();
buffer.put((byte) type);
buffer.putLong(dataTTL);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetUsingDeviceTemplatePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetUsingDeviceTemplatePlan.java
index 6d20145..2360947 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetUsingDeviceTemplatePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/SetUsingDeviceTemplatePlan.java
@@ -57,7 +57,7 @@ public class SetUsingDeviceTemplatePlan extends PhysicalPlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.SET_USING_DEVICE_TEMPLATE.ordinal());
ReadWriteIOUtils.write(prefixPath.getFullPath(), buffer);
buffer.putLong(index);
diff --git
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/StorageGroupMNodePlan.java
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/StorageGroupMNodePlan.java
index 64f0153..e3100d9 100644
---
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/StorageGroupMNodePlan.java
+++
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/StorageGroupMNodePlan.java
@@ -57,7 +57,7 @@ public class StorageGroupMNodePlan extends MNodePlan {
}
@Override
- public void serialize(ByteBuffer buffer) {
+ public void serializeImpl(ByteBuffer buffer) {
buffer.put((byte) PhysicalPlanType.STORAGE_GROUP_MNODE.ordinal());
putString(buffer, name);
buffer.putLong(dataTTL);
diff --git
a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
index 56fafe2..6f0b94c 100644
---
a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
+++
b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
@@ -120,12 +120,10 @@ public class ExclusiveWriteLogNode implements
WriteLogNode, Comparable<Exclusive
}
private void putLog(PhysicalPlan plan) {
- logBufferWorking.mark();
try {
plan.serialize(logBufferWorking);
} catch (BufferOverflowException e) {
logger.info("WAL BufferOverflow !");
- logBufferWorking.reset();
sync();
plan.serialize(logBufferWorking);
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/qp/physical/PhysicalPlanTest.java
b/server/src/test/java/org/apache/iotdb/db/qp/physical/PhysicalPlanTest.java
index 80f26ae..8d4a43e 100644
--- a/server/src/test/java/org/apache/iotdb/db/qp/physical/PhysicalPlanTest.java
+++ b/server/src/test/java/org/apache/iotdb/db/qp/physical/PhysicalPlanTest.java
@@ -28,30 +28,8 @@ import
org.apache.iotdb.db.exception.runtime.SQLParserException;
import org.apache.iotdb.db.metadata.PartialPath;
import org.apache.iotdb.db.qp.Planner;
import org.apache.iotdb.db.qp.logical.Operator.OperatorType;
-import org.apache.iotdb.db.qp.physical.crud.AggregationPlan;
-import org.apache.iotdb.db.qp.physical.crud.DeletePlan;
-import org.apache.iotdb.db.qp.physical.crud.FillQueryPlan;
-import org.apache.iotdb.db.qp.physical.crud.GroupByTimeFillPlan;
-import org.apache.iotdb.db.qp.physical.crud.GroupByTimePlan;
-import org.apache.iotdb.db.qp.physical.crud.LastQueryPlan;
-import org.apache.iotdb.db.qp.physical.crud.QueryPlan;
-import org.apache.iotdb.db.qp.physical.crud.RawDataQueryPlan;
-import org.apache.iotdb.db.qp.physical.crud.UDTFPlan;
-import org.apache.iotdb.db.qp.physical.sys.AuthorPlan;
-import org.apache.iotdb.db.qp.physical.sys.CreateFunctionPlan;
-import org.apache.iotdb.db.qp.physical.sys.CreateTimeSeriesPlan;
-import org.apache.iotdb.db.qp.physical.sys.CreateTriggerPlan;
-import org.apache.iotdb.db.qp.physical.sys.DataAuthPlan;
-import org.apache.iotdb.db.qp.physical.sys.DropFunctionPlan;
-import org.apache.iotdb.db.qp.physical.sys.DropTriggerPlan;
-import org.apache.iotdb.db.qp.physical.sys.LoadConfigurationPlan;
-import org.apache.iotdb.db.qp.physical.sys.OperateFilePlan;
-import org.apache.iotdb.db.qp.physical.sys.SetStorageGroupPlan;
-import org.apache.iotdb.db.qp.physical.sys.ShowFunctionsPlan;
-import org.apache.iotdb.db.qp.physical.sys.ShowPlan;
-import org.apache.iotdb.db.qp.physical.sys.ShowTriggersPlan;
-import org.apache.iotdb.db.qp.physical.sys.StartTriggerPlan;
-import org.apache.iotdb.db.qp.physical.sys.StopTriggerPlan;
+import org.apache.iotdb.db.qp.physical.crud.*;
+import org.apache.iotdb.db.qp.physical.sys.*;
import org.apache.iotdb.db.query.executor.fill.LinearFill;
import org.apache.iotdb.db.query.executor.fill.PreviousFill;
import org.apache.iotdb.db.query.udf.service.UDFRegistrationService;
@@ -77,6 +55,7 @@ import org.junit.Test;
import java.io.File;
import java.io.IOException;
+import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@@ -1274,4 +1253,30 @@ public class PhysicalPlanTest {
new SingleSeriesExpression(new Path("root.vehicle.d5", "s1"),
ValueFilter.like("string*"));
assertEquals(expect.toString(), queryFilter.toString());
}
+
+ @Test(expected = UnsupportedOperationException.class)
+ public void testSerializationError() {
+ ShowDevicesPlan plan = new ShowDevicesPlan();
+
+ ByteBuffer byteBuffer = ByteBuffer.allocate(10);
+ plan.serialize(byteBuffer);
+ }
+
+ @Test(expected = NullPointerException.class)
+ public void testSerializationRollback() {
+ InsertRowPlan plan = new InsertRowPlan();
+ // only serialize time
+ plan.setTime(0L);
+
+ ByteBuffer byteBuffer = ByteBuffer.allocate(10000);
+ byteBuffer.putInt(0);
+ long position = byteBuffer.position();
+
+ try {
+ plan.serialize(byteBuffer);
+ } catch (NullPointerException e) {
+ Assert.assertEquals(position, byteBuffer.position());
+ throw e;
+ }
+ }
}