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 a79f252  add plan index and flush/close listeners (#1850)
a79f252 is described below

commit a79f2525d554fffabee4487b50d4c2b1e3e037ac
Author: Jiang Tian <[email protected]>
AuthorDate: Tue Oct 27 16:28:30 2020 +0800

    add plan index and flush/close listeners (#1850)
---
 .../org/apache/iotdb/db/engine/StorageEngine.java  | 37 ++++++++++--
 .../iotdb/db/engine/flush/CloseFileListener.java   | 28 +++++++++
 .../iotdb/db/engine/flush/FlushListener.java       | 45 +++++++++++++++
 .../iotdb/db/engine/memtable/AbstractMemTable.java | 23 ++++++++
 .../apache/iotdb/db/engine/memtable/IMemTable.java |  4 ++
 .../engine/storagegroup/StorageGroupProcessor.java | 41 ++++++++++---
 .../db/engine/storagegroup/TsFileProcessor.java    | 67 ++++++++++++++++------
 .../db/engine/storagegroup/TsFileResource.java     | 44 +++++++++++++-
 .../org/apache/iotdb/db/monitor/StatMonitor.java   |  2 +-
 .../apache/iotdb/db/qp/executor/IPlanExecutor.java |  3 +-
 .../apache/iotdb/db/qp/executor/PlanExecutor.java  |  8 +--
 .../apache/iotdb/db/qp/physical/PhysicalPlan.java  |  9 +++
 .../iotdb/db/qp/physical/crud/DeletePlan.java      |  6 ++
 .../iotdb/db/qp/physical/crud/InsertRowPlan.java   |  6 +-
 .../db/qp/physical/crud/InsertTabletPlan.java      |  5 ++
 .../iotdb/db/qp/physical/sys/AuthorPlan.java       |  6 ++
 .../qp/physical/sys/CreateMultiTimeSeriesPlan.java |  6 ++
 .../db/qp/physical/sys/CreateTimeSeriesPlan.java   |  4 ++
 .../iotdb/db/qp/physical/sys/DataAuthPlan.java     |  4 ++
 .../db/qp/physical/sys/DeleteStorageGroupPlan.java |  6 ++
 .../db/qp/physical/sys/DeleteTimeSeriesPlan.java   |  6 ++
 .../db/qp/physical/sys/LoadConfigurationPlan.java  |  3 +
 .../db/qp/physical/sys/SetStorageGroupPlan.java    |  4 ++
 .../iotdb/db/qp/physical/sys/SetTTLPlan.java       |  6 ++
 .../db/qp/physical/sys/ShowTimeSeriesPlan.java     |  4 ++
 .../apache/iotdb/db/writelog/WALFlushListener.java | 49 ++++++++++++++++
 .../engine/modification/DeletionFileNodeTest.java  | 28 ++++-----
 .../db/engine/modification/DeletionQueryTest.java  | 52 ++++++++---------
 .../storagegroup/StorageGroupProcessorTest.java    |  2 +-
 .../db/sync/receiver/load/FileLoaderTest.java      |  3 +
 30 files changed, 430 insertions(+), 81 deletions(-)

diff --git a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java 
b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
index 0823608..a3c3b33 100644
--- a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
+++ b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.db.engine;
 import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.ConcurrentModificationException;
 import java.util.HashMap;
 import java.util.List;
@@ -41,6 +42,8 @@ import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.conf.ServerConfigConsistent;
 import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
+import org.apache.iotdb.db.engine.flush.CloseFileListener;
+import org.apache.iotdb.db.engine.flush.FlushListener;
 import org.apache.iotdb.db.engine.flush.TsFileFlushPolicy;
 import org.apache.iotdb.db.engine.flush.TsFileFlushPolicy.DirectFlushPolicy;
 import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
@@ -125,6 +128,10 @@ public class StorageEngine implements IService {
   private ScheduledExecutorService ttlCheckThread;
   private TsFileFlushPolicy fileFlushPolicy = new DirectFlushPolicy();
 
+  // add customized listeners here for flush and close events
+  private List<CloseFileListener> customCloseFileListeners = new ArrayList<>();
+  private List<FlushListener> customFlushListeners = new ArrayList<>();
+
   /**
    * Time range for dividing storage group, the time unit is the same with 
IoTDB's
    * TimestampPrecision
@@ -180,6 +187,8 @@ public class StorageEngine implements IService {
           StorageGroupProcessor processor = new 
StorageGroupProcessor(systemDir,
               storageGroup.getFullPath(), fileFlushPolicy);
           processor.setDataTTL(storageGroup.getDataTTL());
+          processor.setCustomCloseFileListeners(customCloseFileListeners);
+          processor.setCustomFlushListeners(customFlushListeners);
           processorMap.put(storageGroup.getPartialPath(), processor);
           logger.info("Storage Group Processor {} is recovered successfully",
               storageGroup.getFullPath());
@@ -311,6 +320,8 @@ public class StorageEngine implements IService {
               processor = new StorageGroupProcessor(systemDir, 
storageGroupPath.getFullPath(),
                   fileFlushPolicy);
               processor.setDataTTL(storageGroupMNode.getDataTTL());
+              processor.setCustomFlushListeners(customFlushListeners);
+              processor.setCustomCloseFileListeners(customCloseFileListeners);
               processorMap.put(storageGroupPath, processor);
             }
           }
@@ -447,14 +458,14 @@ public class StorageEngine implements IService {
     // TODO
   }
 
-  public void delete(PartialPath path, long startTime, long endTime)
+  public void delete(PartialPath path, long startTime, long endTime, long 
planIndex)
           throws StorageEngineException {
     try {
       List<PartialPath> sgPaths = 
IoTDB.metaManager.searchAllRelatedStorageGroups(path);
       for (PartialPath storageGroupPath : sgPaths) {
         StorageGroupProcessor storageGroupProcessor = 
getProcessor(storageGroupPath);
         PartialPath newPath = path.alterPrefixPath(storageGroupPath);
-        storageGroupProcessor.delete(newPath, startTime, endTime);
+        storageGroupProcessor.delete(newPath, startTime, endTime, planIndex);
       }
     } catch (IOException | MetadataException e) {
       throw new StorageEngineException(e.getMessage());
@@ -464,13 +475,13 @@ public class StorageEngine implements IService {
   /**
    * delete data of timeseries "{deviceId}.{measurementId}"
    */
-  public void deleteTimeseries(PartialPath path)
+  public void deleteTimeseries(PartialPath path, long planIndex)
       throws StorageEngineException {
     try {
       for (PartialPath storageGroupPath : 
IoTDB.metaManager.searchAllRelatedStorageGroups(path)) {
         StorageGroupProcessor storageGroupProcessor = 
getProcessor(storageGroupPath);
         PartialPath newPath = path.alterPrefixPath(storageGroupPath);
-        storageGroupProcessor.delete(newPath, Long.MIN_VALUE, Long.MAX_VALUE);
+        storageGroupProcessor.delete(newPath, Long.MIN_VALUE, Long.MAX_VALUE, 
planIndex);
       }
     } catch (IOException | MetadataException e) {
       throw new StorageEngineException(e.getMessage());
@@ -691,4 +702,22 @@ public class StorageEngine implements IService {
   public static void setEnablePartition(boolean enablePartition) {
     StorageEngine.enablePartition = enablePartition;
   }
+
+  /**
+   * Add a listener to listen flush start/end events. Notice that this 
addition only applies to
+   * TsFileProcessors created afterwards.
+   * @param listener
+   */
+  public void registerFlushListener(FlushListener listener) {
+    customFlushListeners.add(listener);
+  }
+
+  /**
+   * Add a listener to listen file close events. Notice that this addition 
only applies to
+   * TsFileProcessors created afterwards.
+   * @param listener
+   */
+  public void registerCloseFileListener(CloseFileListener listener) {
+    customCloseFileListeners.add(listener);
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/flush/CloseFileListener.java 
b/server/src/main/java/org/apache/iotdb/db/engine/flush/CloseFileListener.java
new file mode 100644
index 0000000..cd98dad
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/flush/CloseFileListener.java
@@ -0,0 +1,28 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.engine.flush;
+
+import org.apache.iotdb.db.engine.storagegroup.TsFileProcessor;
+import org.apache.iotdb.db.exception.TsFileProcessorException;
+
+@FunctionalInterface
+public interface CloseFileListener {
+  void onClosed(TsFileProcessor processor) throws TsFileProcessorException;
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/flush/FlushListener.java 
b/server/src/main/java/org/apache/iotdb/db/engine/flush/FlushListener.java
new file mode 100644
index 0000000..ff8eb84
--- /dev/null
+++ b/server/src/main/java/org/apache/iotdb/db/engine/flush/FlushListener.java
@@ -0,0 +1,45 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.engine.flush;
+
+import java.io.IOException;
+import org.apache.iotdb.db.engine.memtable.IMemTable;
+
+public interface FlushListener {
+
+  void onFlushStart(IMemTable memTable) throws IOException;
+
+  void onFlushEnd(IMemTable memTable);
+
+  class EmptyListener implements FlushListener {
+
+    public static final EmptyListener INSTANCE = new EmptyListener();
+
+    @Override
+    public void onFlushStart(IMemTable memTable) {
+      // do nothing
+    }
+
+    @Override
+    public void onFlushEnd(IMemTable memTable) {
+      // do nothing
+    }
+  }
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
index 2e4884b..598cdc5 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
@@ -61,6 +61,10 @@ public abstract class AbstractMemTable implements IMemTable {
 
   private long totalPointsNumThreshold = 0;
 
+  private long maxPlanIndex = Long.MIN_VALUE;
+
+  private long minPlanIndex = Long.MAX_VALUE;
+
   public AbstractMemTable() {
     this.memTableMap = new HashMap<>();
   }
@@ -101,6 +105,7 @@ public abstract class AbstractMemTable implements IMemTable 
{
 
   @Override
   public void insert(InsertRowPlan insertRowPlan) {
+    updatePlanIndexes(insertRowPlan.getIndex());
     for (int i = 0; i < insertRowPlan.getValues().length; i++) {
 
       if (insertRowPlan.getValues()[i] == null) {
@@ -120,6 +125,7 @@ public abstract class AbstractMemTable implements IMemTable 
{
   @Override
   public void insertTablet(InsertTabletPlan insertTabletPlan, int start, int 
end)
       throws WriteProcessException {
+    updatePlanIndexes(insertTabletPlan.getIndex());
     try {
       write(insertTabletPlan, start, end);
       memSize += MemUtils.getRecordSize(insertTabletPlan, start, end);
@@ -140,6 +146,7 @@ public abstract class AbstractMemTable implements IMemTable 
{
 
   @Override
   public void write(InsertTabletPlan insertTabletPlan, int start, int end) {
+    updatePlanIndexes(insertTabletPlan.getIndex());
     for (int i = 0; i < insertTabletPlan.getMeasurements().length; i++) {
       if (insertTabletPlan.getColumns()[i] == null) {
         continue;
@@ -192,6 +199,7 @@ public abstract class AbstractMemTable implements IMemTable 
{
     seriesNumber = 0;
     totalPointsNum = 0;
     totalPointsNumThreshold = 0;
+    maxPlanIndex = 0;
   }
 
   @Override
@@ -274,4 +282,19 @@ public abstract class AbstractMemTable implements 
IMemTable {
       }
     }
   }
+
+  @Override
+  public long getMaxPlanIndex() {
+    return maxPlanIndex;
+  }
+
+  @Override
+  public long getMinPlanIndex() {
+    return minPlanIndex;
+  }
+
+  void updatePlanIndexes(long index) {
+    maxPlanIndex = Math.max(index, maxPlanIndex);
+    minPlanIndex = Math.min(index, minPlanIndex);
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java 
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
index 3cadfd9..c6134a0 100644
--- a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
+++ b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
@@ -119,4 +119,8 @@ public interface IMemTable {
   void setVersion(long version);
 
   void release();
+
+  long getMaxPlanIndex();
+
+  long getMinPlanIndex();
 }
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 c791c3c..f26078d 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
@@ -48,6 +48,8 @@ import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.conf.directories.DirectoryManager;
 import org.apache.iotdb.db.engine.StorageEngine;
 import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
+import org.apache.iotdb.db.engine.flush.CloseFileListener;
+import org.apache.iotdb.db.engine.flush.FlushListener;
 import org.apache.iotdb.db.engine.flush.TsFileFlushPolicy;
 import org.apache.iotdb.db.engine.merge.manage.MergeManager;
 import org.apache.iotdb.db.engine.merge.task.RecoverMergeTask;
@@ -235,6 +237,9 @@ public class StorageGroupProcessor {
    */
   private Map<Long, Long> partitionMaxFileVersions = new HashMap<>();
 
+  private List<CloseFileListener> customCloseFileListeners = 
Collections.emptyList();
+  private List<FlushListener> customFlushListeners = Collections.emptyList();
+
   public StorageGroupProcessor(String systemDir, String storageGroupName,
       TsFileFlushPolicy fileFlushPolicy) throws StorageGroupProcessorException 
{
     this.storageGroupName = storageGroupName;
@@ -962,6 +967,8 @@ public class StorageGroupProcessor {
           versionController, this::closeUnsealedTsFileProcessorCallBack,
           this::unsequenceFlushCallback, false);
     }
+    tsFileProcessor.addCloseFileListeners(customCloseFileListeners);
+    tsFileProcessor.addFlushListeners(customFlushListeners);
 
     tsFileProcessor.setTimeRangeId(timePartitionId);
     return tsFileProcessor;
@@ -1345,12 +1352,12 @@ public class StorageGroupProcessor {
   /**
    * Delete data whose timestamp <= 'timestamp' and belongs to the time series
    * deviceId.measurementId.
-   *
-   * @param path the timeseries path of the to be deleted.
+   *  @param path the timeseries path of the to be deleted.
    * @param startTime the startTime of delete range.
    * @param endTime the endTime of delete range.
+   * @param planIndex
    */
-  public void delete(PartialPath path, long startTime, long endTime) throws 
IOException {
+  public void delete(PartialPath path, long startTime, long endTime, long 
planIndex) throws IOException {
     // TODO: how to avoid partial deletion?
     // FIXME: notice that if we may remove a SGProcessor out of memory, we 
need to close all opened
     //mod files in mergingModification, sequenceFileList, and 
unsequenceFileList
@@ -1389,8 +1396,10 @@ public class StorageGroupProcessor {
         updatedModFiles.add(tsFileManagement.mergingModification);
       }
 
-      deleteDataInFiles(tsFileManagement.getTsFileList(true), deletion, 
devicePaths, updatedModFiles);
-      deleteDataInFiles(tsFileManagement.getTsFileList(false), deletion, 
devicePaths, updatedModFiles);
+      deleteDataInFiles(tsFileManagement.getTsFileList(true), deletion, 
devicePaths,
+          updatedModFiles, planIndex);
+      deleteDataInFiles(tsFileManagement.getTsFileList(false), deletion, 
devicePaths,
+          updatedModFiles, planIndex);
 
     } catch (Exception e) {
       // roll back
@@ -1438,8 +1447,8 @@ public class StorageGroupProcessor {
   }
 
   private void deleteDataInFiles(Collection<TsFileResource> 
tsFileResourceList, Deletion deletion,
-      Set<PartialPath> devicePaths, List<ModificationFile> updatedModFiles)
-          throws IOException, MetadataException {
+      Set<PartialPath> devicePaths, List<ModificationFile> updatedModFiles, 
long planIndex)
+          throws IOException {
     for (TsFileResource tsFileResource : tsFileResourceList) {
       if (canSkipDelete(tsFileResource, devicePaths, deletion.getStartTime(), 
deletion.getEndTime())) {
         continue;
@@ -1453,6 +1462,8 @@ public class StorageGroupProcessor {
       // remember to close mod file
       tsFileResource.getModFile().close();
 
+      tsFileResource.updatePlanIndexes(planIndex);
+
       // delete data in memory of unsealed file
       if (!tsFileResource.isClosed()) {
         TsFileProcessor tsfileProcessor = 
tsFileResource.getUnsealedFileProcessor();
@@ -2104,8 +2115,10 @@ public class StorageGroupProcessor {
     }
     partitionDirectFileVersions.computeIfAbsent(filePartitionId,
         p -> new HashSet<>()).addAll(tsFileResource.getHistoricalVersions());
-    updatePartitionFileVersion(filePartitionId,
-        Collections.max(tsFileResource.getHistoricalVersions()));
+    if (!tsFileResource.getHistoricalVersions().isEmpty()) {
+      updatePartitionFileVersion(filePartitionId,
+          Collections.max(tsFileResource.getHistoricalVersions()));
+    }
     return true;
   }
 
@@ -2372,4 +2385,14 @@ public class StorageGroupProcessor {
 
     boolean satisfy(String storageGroupName, long timePartitionId);
   }
+
+  public void setCustomCloseFileListeners(
+      List<CloseFileListener> customCloseFileListeners) {
+    this.customCloseFileListeners = customCloseFileListeners;
+  }
+
+  public void setCustomFlushListeners(
+      List<FlushListener> customFlushListeners) {
+    this.customFlushListeners = customFlushListeners;
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
index 3e26e81..9374ccb 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
@@ -23,6 +23,7 @@ import static 
org.apache.iotdb.db.conf.adapter.IoTDBConfigDynamicAdapter.MEMTABL
 import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.Collection;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
@@ -36,6 +37,8 @@ import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.conf.adapter.ActiveTimeSeriesCounter;
 import org.apache.iotdb.db.conf.adapter.CompressionRatio;
 import org.apache.iotdb.db.conf.adapter.IoTDBConfigDynamicAdapter;
+import org.apache.iotdb.db.engine.flush.CloseFileListener;
+import org.apache.iotdb.db.engine.flush.FlushListener;
 import org.apache.iotdb.db.engine.flush.FlushManager;
 import org.apache.iotdb.db.engine.flush.MemTableFlushTask;
 import org.apache.iotdb.db.engine.flush.NotifyFlushMemTable;
@@ -44,7 +47,6 @@ import org.apache.iotdb.db.engine.modification.Deletion;
 import org.apache.iotdb.db.engine.modification.Modification;
 import org.apache.iotdb.db.engine.modification.ModificationFile;
 import org.apache.iotdb.db.engine.querycontext.ReadOnlyMemChunk;
-import 
org.apache.iotdb.db.engine.storagegroup.StorageGroupProcessor.CloseTsFileCallBack;
 import 
org.apache.iotdb.db.engine.storagegroup.StorageGroupProcessor.UpdateEndTimeCallBack;
 import org.apache.iotdb.db.engine.version.VersionController;
 import org.apache.iotdb.db.exception.TsFileProcessorException;
@@ -57,6 +59,7 @@ import org.apache.iotdb.db.qp.physical.crud.InsertTabletPlan;
 import org.apache.iotdb.db.query.context.QueryContext;
 import org.apache.iotdb.db.rescon.MemTablePool;
 import org.apache.iotdb.db.utils.QueryUtils;
+import org.apache.iotdb.db.writelog.WALFlushListener;
 import org.apache.iotdb.db.writelog.manager.MultiFileLogNodeManager;
 import org.apache.iotdb.db.writelog.node.WriteLogNode;
 import org.apache.iotdb.rpc.RpcUtils;
@@ -97,10 +100,7 @@ public class TsFileProcessor {
   private IMemTable workMemTable;
 
   private final VersionController versionController;
-  /**
-   * this callback is called after the corresponding TsFile is called 
endFile().
-   */
-  private final CloseTsFileCallBack closeTsFileCallback;
+
   /**
    * this callback is called before the workMemtable is added into the 
flushingMemTables.
    */
@@ -112,36 +112,41 @@ public class TsFileProcessor {
   private static final String FLUSH_QUERY_WRITE_LOCKED = "{}: {} get 
flushQueryLock write lock";
   private static final String FLUSH_QUERY_WRITE_RELEASE = "{}: {} get 
flushQueryLock write lock released";
 
+  private List<CloseFileListener> closeFileListeners = new ArrayList<>();
+  private List<FlushListener> flushListeners = new ArrayList<>();
+
   TsFileProcessor(String storageGroupName, File tsfile,
       VersionController versionController,
-      CloseTsFileCallBack closeTsFileCallback,
+      CloseFileListener closeTsFileCallback,
       UpdateEndTimeCallBack updateLatestFlushTimeCallback, boolean sequence)
       throws IOException {
     this.storageGroupName = storageGroupName;
     this.tsFileResource = new TsFileResource(tsfile, this);
     this.versionController = versionController;
     this.writer = new RestorableTsFileIOWriter(tsfile);
-    this.closeTsFileCallback = closeTsFileCallback;
     this.updateLatestFlushTimeCallback = updateLatestFlushTimeCallback;
     this.sequence = sequence;
     logger.info("create a new tsfile processor {}", tsfile.getAbsolutePath());
     // a file generated by flush has only one historical version, which is 
itself
     this.tsFileResource
         
.setHistoricalVersions(Collections.singleton(versionController.currVersion()));
+    flushListeners.add(new WALFlushListener(this));
+    closeFileListeners.add(closeTsFileCallback);
   }
 
   public TsFileProcessor(String storageGroupName, TsFileResource 
tsFileResource,
-      VersionController versionController, CloseTsFileCallBack 
closeUnsealedTsFileProcessor,
+      VersionController versionController, CloseFileListener 
closeUnsealedTsFileProcessor,
       UpdateEndTimeCallBack updateLatestFlushTimeCallback, boolean sequence,
       RestorableTsFileIOWriter writer) {
     this.storageGroupName = storageGroupName;
     this.tsFileResource = tsFileResource;
     this.versionController = versionController;
     this.writer = writer;
-    this.closeTsFileCallback = closeUnsealedTsFileProcessor;
     this.updateLatestFlushTimeCallback = updateLatestFlushTimeCallback;
     this.sequence = sequence;
     logger.info("reopen a tsfile processor {}", tsFileResource.getTsFile());
+    flushListeners.add(new WALFlushListener(this));
+    closeFileListeners.add(closeUnsealedTsFileProcessor);
   }
 
   /**
@@ -176,6 +181,7 @@ public class TsFileProcessor {
       tsFileResource
           .updateEndTime(insertRowPlan.getDeviceId().getFullPath(), 
insertRowPlan.getTime());
     }
+    tsFileResource.updatePlanIndexes(insertRowPlan.getIndex());
   }
 
   /**
@@ -224,6 +230,7 @@ public class TsFileProcessor {
           .updateEndTime(
               insertTabletPlan.getDeviceId().getFullPath(), 
insertTabletPlan.getTimes()[end - 1]);
     }
+    tsFileResource.updatePlanIndexes(insertTabletPlan.getIndex());
   }
 
   /**
@@ -232,8 +239,7 @@ public class TsFileProcessor {
    * <p>
    * Delete data in both working MemTable and flushing MemTables.
    */
-  public void deleteDataInMemory(Deletion deletion, Set<PartialPath> 
devicePaths)
-          throws MetadataException {
+  public void deleteDataInMemory(Deletion deletion, Set<PartialPath> 
devicePaths) {
     flushQueryLock.writeLock().lock();
     if (logger.isDebugEnabled()) {
       logger
@@ -490,6 +496,11 @@ public class TsFileProcessor {
           storageGroupName, tsFileResource.getTsFile().getName(), 
tobeFlushed.getMemTableMap());
       return;
     }
+
+    for (FlushListener flushListener : flushListeners) {
+      flushListener.onFlushStart(tobeFlushed);
+    }
+
     flushingMemTables.addLast(tobeFlushed);
     if (logger.isDebugEnabled()) {
       logger.debug(
@@ -499,9 +510,7 @@ public class TsFileProcessor {
     }
     long cur = versionController.nextVersion();
     tobeFlushed.setVersion(cur);
-    if (IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
-      getLogNode().notifyStartFlush();
-    }
+
     if (!tobeFlushed.isSignalMemTable()) {
       totalMemTableSize += tobeFlushed.memSize();
     }
@@ -557,6 +566,7 @@ public class TsFileProcessor {
   public void flushOneMemTable() {
     IMemTable memTableToFlush;
     memTableToFlush = flushingMemTables.getFirst();
+
     // signal memtable only may appear when calling asyncClose()
     if (!memTableToFlush.isSignalMemTable()) {
       try {
@@ -578,11 +588,12 @@ public class TsFileProcessor {
         }
         Thread.currentThread().interrupt();
       }
+    }
 
-      if (IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
-        getLogNode().notifyEndFlush();
-      }
+    for (FlushListener flushListener : flushListeners) {
+      flushListener.onFlushEnd(memTableToFlush);
     }
+
     if (logger.isDebugEnabled()) {
       logger.debug("{}: {} try get lock to release a memtable (signal={})", 
storageGroupName,
           tsFileResource.getTsFile().getName(), 
memTableToFlush.isSignalMemTable());
@@ -661,7 +672,9 @@ public class TsFileProcessor {
 
     // remove this processor from Closing list in StorageGroupProcessor,
     // mark the TsFileResource closed, no need writer anymore
-    closeTsFileCallback.call(this);
+    for (CloseFileListener closeFileListener : closeFileListeners) {
+      closeFileListener.onClosed(this);
+    }
 
     if (logger.isInfoEnabled()) {
       long closeEndTime = System.currentTimeMillis();
@@ -684,7 +697,7 @@ public class TsFileProcessor {
     this.managedByFlushManager = managedByFlushManager;
   }
 
-  WriteLogNode getLogNode() {
+  public WriteLogNode getLogNode() {
     if (logNode == null) {
       logNode = MultiFileLogNodeManager.getInstance()
           .getNode(storageGroupName + "-" + 
tsFileResource.getTsFile().getName());
@@ -802,4 +815,20 @@ public class TsFileProcessor {
       throw new TsFileProcessorException(e);
     }
   }
+
+  public void addFlushListener(FlushListener listener) {
+    flushListeners.add(listener);
+  }
+
+  public void addCloseFileListener(CloseFileListener listener) {
+    closeFileListeners.add(listener);
+  }
+
+  public void addFlushListeners(Collection<FlushListener> listeners) {
+    flushListeners.addAll(listeners);
+  }
+
+  public void addCloseFileListeners(Collection<CloseFileListener> listeners) {
+    closeFileListeners.addAll(listeners);
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
index f71f83a..b295e90 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
@@ -155,6 +155,16 @@ public class TsFileResource {
    */
   private TsFileResource originTsFileResource;
 
+  /**
+   * Maximum index of plans executed within this TsFile.
+   */
+  private long maxPlanIndex = Long.MIN_VALUE;
+
+  /**
+   * Minimum index of plans executed within this TsFile.
+   */
+  private long minPlanIndex = Long.MAX_VALUE;
+
   public TsFileResource() {
   }
 
@@ -275,6 +285,9 @@ public class TsFileResource {
         }
       }
 
+      ReadWriteIOUtils.write(maxPlanIndex, outputStream);
+      ReadWriteIOUtils.write(minPlanIndex, outputStream);
+
       if (modFile != null && modFile.exists()) {
         String modFileName = new File(modFile.getFilePath()).getName();
         ReadWriteIOUtils.write(modFileName, outputStream);
@@ -324,6 +337,9 @@ public class TsFileResource {
         historicalVersions = Collections.singleton(version);
       }
 
+      maxPlanIndex = ReadWriteIOUtils.readLong(inputStream);
+      minPlanIndex = ReadWriteIOUtils.readLong(inputStream);
+
       if (inputStream.available() > 0) {
         String modFileName = ReadWriteIOUtils.readString(inputStream);
         File modF = new File(file.getParentFile(), modFileName);
@@ -811,7 +827,7 @@ public class TsFileResource {
 
   public long getMaxVersion() {
     long maxVersion = 0;
-    if (historicalVersions != null) {
+    if (historicalVersions != null && !historicalVersions.isEmpty()) {
       maxVersion = Collections.max(historicalVersions);
     }
     return maxVersion;
@@ -824,4 +840,30 @@ public class TsFileResource {
           .getFile(file.toPath() + TsFileResource.RESOURCE_SUFFIX).toPath());
     }
   }
+
+  public long getMaxPlanIndex() {
+    return maxPlanIndex;
+  }
+
+  public long getMinPlanIndex() {
+    return minPlanIndex;
+  }
+
+  public void updatePlanIndexes(long planIndex) {
+    maxPlanIndex = Math.max(maxPlanIndex, planIndex);
+    minPlanIndex = Math.min(minPlanIndex, planIndex);
+    if (closed) {
+      try {
+        serialize();
+      } catch (IOException e) {
+        logger.error("Cannot serialize TsFileResource {} when updating plan 
index {}-{}", this,
+            maxPlanIndex, planIndex);
+      }
+    }
+  }
+
+  public boolean isPlanIndexOverlap(TsFileResource another) {
+    return another.maxPlanIndex >= this.minPlanIndex &&
+           another.minPlanIndex <= this.maxPlanIndex;
+  }
 }
diff --git a/server/src/main/java/org/apache/iotdb/db/monitor/StatMonitor.java 
b/server/src/main/java/org/apache/iotdb/db/monitor/StatMonitor.java
index a4fff85..6a6f8fb 100644
--- a/server/src/main/java/org/apache/iotdb/db/monitor/StatMonitor.java
+++ b/server/src/main/java/org/apache/iotdb/db/monitor/StatMonitor.java
@@ -367,7 +367,7 @@ public class StatMonitor implements IService {
           for (String statParamName : 
entry.getValue().getStatParamsHashMap().keySet()) {
             if (temporaryStatList.contains(statParamName)) {
               fManager.delete(new PartialPath(entry.getKey(), statParamName), 
Long.MIN_VALUE,
-                  currentTimeMillis - statMonitorRetainIntervalSec * 1000);
+                  currentTimeMillis - statMonitorRetainIntervalSec * 1000, -1);
             }
           }
         }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/executor/IPlanExecutor.java 
b/server/src/main/java/org/apache/iotdb/db/qp/executor/IPlanExecutor.java
index 48aab84..744e57b 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/executor/IPlanExecutor.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/executor/IPlanExecutor.java
@@ -80,8 +80,9 @@ public interface IPlanExecutor {
    * @param path       : delete series seriesPath
    * @param startTime start time in delete command
    * @param endTime end time in delete command
+   * @param planIndex index of the deletion plan
    */
-  void delete(PartialPath path, long startTime, long endTime) throws 
QueryProcessException;
+  void delete(PartialPath path, long startTime, long endTime, long planIndex) 
throws QueryProcessException;
 
   /**
    * execute insert command and return whether the operator is successful.
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java 
b/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
index 35fc394..6cfe2b9 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
@@ -711,7 +711,7 @@ public class PlanExecutor implements IPlanExecutor {
   @Override
   public void delete(DeletePlan deletePlan) throws QueryProcessException {
     for (PartialPath path : deletePlan.getPaths()) {
-      delete(path, deletePlan.getDeleteStartTime(), 
deletePlan.getDeleteEndTime());
+      delete(path, deletePlan.getDeleteStartTime(), 
deletePlan.getDeleteEndTime(), deletePlan.getIndex());
     }
   }
 
@@ -881,9 +881,9 @@ public class PlanExecutor implements IPlanExecutor {
   }
 
   @Override
-  public void delete(PartialPath path, long startTime, long endTime) throws 
QueryProcessException {
+  public void delete(PartialPath path, long startTime, long endTime, long 
planIndex) throws QueryProcessException {
     try {
-      StorageEngine.getInstance().delete(path, startTime, endTime);
+      StorageEngine.getInstance().delete(path, startTime, endTime, planIndex);
     } catch (StorageEngineException e) {
       throw new QueryProcessException(e);
     }
@@ -1044,7 +1044,7 @@ public class PlanExecutor implements IPlanExecutor {
     try {
       List<String> failedNames = new LinkedList<>();
       for (PartialPath path : deletePathList) {
-        StorageEngine.getInstance().deleteTimeseries(path);
+        StorageEngine.getInstance().deleteTimeseries(path, 
deleteTimeSeriesPlan.getIndex());
         String failedTimeseries = mManager.deleteTimeseries(path);
         if (!failedTimeseries.isEmpty()) {
           failedNames.add(failedTimeseries);
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 56b3257..eacebf8 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 @@ public abstract class PhysicalPlan {
   //login username, corresponding to cli/session login user info
   private String loginUserName;
 
+  // a bridge from a cluster raft log to a physical plan
+  protected long index;
+
   /**
    * whether the plan can be split into more than one Plans. Only used in the 
cluster mode.
    */
@@ -283,5 +286,11 @@ public abstract class PhysicalPlan {
     DELETE_STORAGE_GROUP, SHOW_TIMESERIES, DELETE_TIMESERIES, 
LOAD_CONFIGURATION, MULTI_CREATE_TIMESERIES
   }
 
+  public long getIndex() {
+    return index;
+  }
 
+  public void setIndex(long index) {
+    this.index = index;
+  }
 }
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 54936de..d4cc835 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
@@ -128,6 +128,8 @@ public class DeletePlan extends PhysicalPlan {
     for (PartialPath path : paths) {
       putString(stream, path.getFullPath());
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -140,6 +142,8 @@ public class DeletePlan extends PhysicalPlan {
     for (PartialPath path : paths) {
       putString(buffer, path.getFullPath());
     }
+
+    buffer.putLong(index);
   }
 
   @Override
@@ -151,5 +155,7 @@ public class DeletePlan extends PhysicalPlan {
     for (int i = 0; i < pathSize; i++) {
       paths.add(new PartialPath(readString(buffer)));
     }
+
+    this.index = buffer.getLong();
   }
 }
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 1271ecd..06ab4ac 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
@@ -247,6 +247,8 @@ public class InsertRowPlan extends InsertPlan {
 
     // the types are not inferred before the plan is serialized
     stream.write((byte) (isNeedInferType ? 1 : 0));
+
+    stream.writeLong(index);
   }
 
   private void putValues(DataOutputStream outputStream) throws 
QueryProcessException, IOException {
@@ -380,11 +382,12 @@ public class InsertRowPlan extends InsertPlan {
     try {
       putValues(buffer);
     } catch (QueryProcessException e) {
-      e.printStackTrace();
+      logger.error("Failed to serialize values for {}", this, e);
     }
 
     // the types are not inferred before the plan is serialized
     buffer.put((byte) (isNeedInferType ? 1 : 0));
+    buffer.putLong(index);
   }
 
   @Override
@@ -408,6 +411,7 @@ public class InsertRowPlan extends InsertPlan {
     }
 
     isNeedInferType = buffer.get() == 1;
+    this.index = buffer.getLong();
   }
 
   @Override
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 f7db816..4119a4a 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
@@ -151,6 +151,8 @@ public class InsertTabletPlan extends InsertPlan {
       stream.write(valueBuffer.array());
       valueBuffer = null;
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -191,6 +193,8 @@ public class InsertTabletPlan extends InsertPlan {
       buffer.put(valueBuffer.array());
       valueBuffer = null;
     }
+
+    buffer.putLong(index);
   }
 
   private void serializeValues(DataOutputStream outputStream) throws 
IOException {
@@ -331,6 +335,7 @@ public class InsertTabletPlan extends InsertPlan {
     times = QueryDataSetUtils.readTimesFromBuffer(buffer, rows);
 
     columns = QueryDataSetUtils.readValuesFromBuffer(buffer, dataTypes, 
measurementSize, rows);
+    this.index = buffer.getLong();
   }
 
   public void setDataTypes(List<Integer> dataTypes) {
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 b43bee9..f09f584 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
@@ -303,6 +303,8 @@ public class AuthorPlan extends PhysicalPlan {
     } else {
       putString(stream, nodeName.getFullPath());
     }
+
+    stream.writeLong(index);
   }
 
 
@@ -329,6 +331,8 @@ public class AuthorPlan extends PhysicalPlan {
     } else {
       putString(buffer, nodeName.getFullPath());
     }
+
+    buffer.putLong(index);
   }
 
   @Override
@@ -354,6 +358,8 @@ public class AuthorPlan extends PhysicalPlan {
     } else {
       this.nodeName = new PartialPath(nodeNameStr);
     }
+
+    this.index = buffer.getLong();
   }
 
   private int getPlanType(OperatorType operatorType) {
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 feb2278..ff2012e 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
@@ -195,6 +195,8 @@ public class CreateMultiTimeSeriesPlan extends PhysicalPlan 
{
     } else {
       stream.write(0);
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -254,6 +256,8 @@ public class CreateMultiTimeSeriesPlan extends PhysicalPlan 
{
     } else {
       buffer.put((byte) 0);
     }
+
+    buffer.putLong(index);
   }
 
   @Override
@@ -299,5 +303,7 @@ public class CreateMultiTimeSeriesPlan extends PhysicalPlan 
{
         attributes.add(ReadWriteIOUtils.readMap(buffer));
       }
     }
+
+    this.index = buffer.getLong();
   }
 }
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 7cace1e..dde1598 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
@@ -181,6 +181,8 @@ public class CreateTimeSeriesPlan extends PhysicalPlan {
     } else {
       stream.write(0);
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -212,6 +214,8 @@ public class CreateTimeSeriesPlan extends PhysicalPlan {
     if (buffer.get() == 1) {
       attributes = ReadWriteIOUtils.readMap(buffer);
     }
+
+    this.index = buffer.getLong();
   }
 
   @Override
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 16cb729..e40fe9b 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
@@ -58,6 +58,8 @@ public class DataAuthPlan extends PhysicalPlan {
     for (String user : users) {
       putString(stream, user);
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -69,6 +71,8 @@ public class DataAuthPlan extends PhysicalPlan {
     for (String user : users) {
       putString(buffer, user);
     }
+
+    buffer.putLong(index);
   }
 
   @Override
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 3c2a39b..5f8a13f 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
@@ -54,6 +54,8 @@ public class DeleteStorageGroupPlan extends PhysicalPlan {
     for (PartialPath path : this.getPaths()) {
       putString(stream, path.getFullPath());
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -64,6 +66,8 @@ public class DeleteStorageGroupPlan extends PhysicalPlan {
     for (PartialPath path : this.getPaths()) {
       putString(buffer, path.getFullPath());
     }
+
+    buffer.putLong(index);
   }
 
   @Override
@@ -73,6 +77,8 @@ public class DeleteStorageGroupPlan extends PhysicalPlan {
     for (int i = 0; i < pathNum; i++) {
       deletePathList.add(new PartialPath(readString(buffer)));
     }
+
+    this.index = buffer.getLong();
   }
 
 }
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 f6c0cc8..b4c8b0c 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
@@ -54,6 +54,8 @@ public class DeleteTimeSeriesPlan extends PhysicalPlan {
     for (PartialPath path : deletePathList) {
       putString(stream, path.getFullPath());
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -64,6 +66,8 @@ public class DeleteTimeSeriesPlan extends PhysicalPlan {
     for (PartialPath path : deletePathList) {
       putString(buffer, path.getFullPath());
     }
+
+    buffer.putLong(index);
   }
 
   @Override
@@ -73,5 +77,7 @@ public class DeleteTimeSeriesPlan extends PhysicalPlan {
     for (int i = 0; i < pathNumber; i++) {
       deletePathList.add(new PartialPath(readString(buffer)));
     }
+
+    this.index = buffer.getLong();
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/LoadConfigurationPlan.java
 
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/LoadConfigurationPlan.java
index b5ff5a8..f88a9f9 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/LoadConfigurationPlan.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/LoadConfigurationPlan.java
@@ -90,6 +90,8 @@ public class LoadConfigurationPlan extends PhysicalPlan {
         }
       }
     }
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -108,6 +110,7 @@ public class LoadConfigurationPlan extends PhysicalPlan {
         }
       }
     }
+    this.index = buffer.getLong();
   }
 
   @Override
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 b0bc1fc..2f40dde 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,6 +64,8 @@ public class SetStorageGroupPlan extends PhysicalPlan {
     byte[] fullPathBytes = path.getFullPath().getBytes();
     stream.writeInt(fullPathBytes.length);
     stream.write(fullPathBytes);
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -72,6 +74,8 @@ public class SetStorageGroupPlan extends PhysicalPlan {
     byte[] fullPathBytes = new byte[length];
     buffer.get(fullPathBytes);
     path = new PartialPath(new String(fullPathBytes));
+
+    this.index = buffer.getLong();
   }
 
   @Override
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 6dad2bc..5aeb6e1 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
@@ -61,6 +61,8 @@ public class SetTTLPlan extends PhysicalPlan {
     stream.writeByte((byte) type);
     stream.writeLong(dataTTL);
     putString(stream, storageGroup.getFullPath());
+
+    stream.writeLong(index);
   }
 
   @Override
@@ -69,12 +71,16 @@ public class SetTTLPlan extends PhysicalPlan {
     buffer.put((byte) type);
     buffer.putLong(dataTTL);
     putString(buffer, storageGroup.getFullPath());
+
+    buffer.putLong(index);
   }
 
   @Override
   public void deserialize(ByteBuffer buffer) throws IllegalPathException {
     this.dataTTL = buffer.getLong();
     this.storageGroup = new PartialPath(readString(buffer));
+
+    this.index = buffer.getLong();
   }
 
   public PartialPath getStorageGroup() {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ShowTimeSeriesPlan.java
 
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ShowTimeSeriesPlan.java
index c393c71..79ce934 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ShowTimeSeriesPlan.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/qp/physical/sys/ShowTimeSeriesPlan.java
@@ -120,6 +120,8 @@ public class ShowTimeSeriesPlan extends ShowPlan {
     outputStream.writeInt(limit);
     outputStream.writeInt(offset);
     outputStream.writeBoolean(orderByHeat);
+
+    outputStream.writeLong(index);
   }
 
   @Override
@@ -132,5 +134,7 @@ public class ShowTimeSeriesPlan extends ShowPlan {
     limit = buffer.getInt();
     limit = buffer.getInt();
     orderByHeat = buffer.get() == 1;
+
+    this.index = buffer.getLong();
   }
 }
\ No newline at end of file
diff --git 
a/server/src/main/java/org/apache/iotdb/db/writelog/WALFlushListener.java 
b/server/src/main/java/org/apache/iotdb/db/writelog/WALFlushListener.java
new file mode 100644
index 0000000..9e58af8
--- /dev/null
+++ b/server/src/main/java/org/apache/iotdb/db/writelog/WALFlushListener.java
@@ -0,0 +1,49 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.writelog;
+
+import java.io.IOException;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.engine.flush.FlushListener;
+import org.apache.iotdb.db.engine.memtable.IMemTable;
+import org.apache.iotdb.db.engine.storagegroup.TsFileProcessor;
+
+public class WALFlushListener implements FlushListener {
+
+  private TsFileProcessor processor;
+
+  public WALFlushListener(TsFileProcessor processor) {
+    this.processor = processor;
+  }
+
+  @Override
+  public void onFlushStart(IMemTable memTable) throws IOException {
+    if (IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
+      processor.getLogNode().notifyStartFlush();
+    }
+  }
+
+  @Override
+  public void onFlushEnd(IMemTable memTable) {
+    if (!memTable.isSignalMemTable() && 
IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
+      processor.getLogNode().notifyEndFlush();
+    }
+  }
+}
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionFileNodeTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionFileNodeTest.java
index 165fe6d..8fd2162 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionFileNodeTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionFileNodeTest.java
@@ -107,10 +107,10 @@ public class DeletionFileNodeTest {
       insertToStorageEngine(record);
     }
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50, -1);
 
     SingleSeriesExpression expression = new SingleSeriesExpression(new 
PartialPath(processorName + TsFileConstant.PATH_SEPARATOR +
         measurements[5]), null);
@@ -143,9 +143,9 @@ public class DeletionFileNodeTest {
     }
     StorageEngine.getInstance().syncCloseAllProcessor();
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30, -1);
 
     Modification[] realModifications = new Modification[]{
         new Deletion(new PartialPath(processorName + 
TsFileConstant.PATH_SEPARATOR + measurements[5]), 201, 50),
@@ -207,10 +207,10 @@ public class DeletionFileNodeTest {
       insertToStorageEngine(record);
     }
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50, -1);
 
     SingleSeriesExpression expression = new SingleSeriesExpression(new 
PartialPath(processorName + TsFileConstant.PATH_SEPARATOR +
         measurements[5]), null);
@@ -256,9 +256,9 @@ public class DeletionFileNodeTest {
     }
     StorageEngine.getInstance().syncCloseAllProcessor();
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30, -1);
 
     Modification[] realModifications = new Modification[]{
         new Deletion(new PartialPath(processorName + 
TsFileConstant.PATH_SEPARATOR + measurements[5]), 301, 50),
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionQueryTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionQueryTest.java
index a5b2530..aa783b2 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionQueryTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/modification/DeletionQueryTest.java
@@ -99,10 +99,10 @@ public class DeletionQueryTest {
       insertToStorageEngine(record);
     }
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50, -1);
 
     List<PartialPath> pathList = new ArrayList<>();
     pathList.add(new PartialPath(processorName + TsFileConstant.PATH_SEPARATOR 
+ measurements[3]));
@@ -138,9 +138,9 @@ public class DeletionQueryTest {
     }
     StorageEngine.getInstance().syncCloseAllProcessor();
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30, -1);
 
     List<PartialPath> pathList = new ArrayList<>();
     pathList.add(new PartialPath(processorName + TsFileConstant.PATH_SEPARATOR 
+ measurements[3]));
@@ -187,10 +187,10 @@ public class DeletionQueryTest {
       insertToStorageEngine(record);
     }
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50, -1);
 
     List<PartialPath> pathList = new ArrayList<>();
     pathList.add(new PartialPath(processorName + TsFileConstant.PATH_SEPARATOR 
+ measurements[3]));
@@ -237,9 +237,9 @@ public class DeletionQueryTest {
     }
     StorageEngine.getInstance().syncCloseAllProcessor();
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 40, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 30, -1);
 
     List<PartialPath> pathList = new ArrayList<>();
     pathList.add(new PartialPath(processorName + TsFileConstant.PATH_SEPARATOR 
+ measurements[3]));
@@ -275,10 +275,10 @@ public class DeletionQueryTest {
       insertToStorageEngine(record);
     }
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50, -1);
 
     StorageEngine.getInstance().syncCloseAllProcessor();
 
@@ -290,10 +290,10 @@ public class DeletionQueryTest {
       insertToStorageEngine(record);
     }
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 250);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 250);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 230);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 230, 250);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 250, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 250, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 230, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 230, 250, -1);
 
     StorageEngine.getInstance().syncCloseAllProcessor();
 
@@ -305,10 +305,10 @@ public class DeletionQueryTest {
       insertToStorageEngine(record);
     }
 
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30);
-    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[3]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[4]), 0, 50, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 0, 30, -1);
+    StorageEngine.getInstance().delete(new PartialPath(processorName, 
measurements[5]), 30, 50, -1);
 
     StorageEngine.getInstance().syncCloseAllProcessor();
 
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
index 983be2f..c14e667 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
@@ -116,7 +116,7 @@ public class StorageGroupProcessorTest {
       insertToStorageGroupProcessor(record);
     }
 
-    processor.delete(new PartialPath(deviceId, measurementId), 0, 15L);
+    processor.delete(new PartialPath(deviceId, measurementId), 0, 15L, -1);
 
     List<TsFileResource> tsfileResourcesForQuery = new ArrayList<>();
     for (TsFileProcessor tsfileProcessor : 
processor.getWorkUnsequenceTsFileProcessor()) {
diff --git 
a/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java
 
b/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java
index c8f7339..0a9a9a4 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/sync/receiver/load/FileLoaderTest.java
@@ -25,6 +25,7 @@ import static org.junit.Assert.assertTrue;
 import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
@@ -128,6 +129,7 @@ public class FileLoaderTest {
         TsFileResource tsFileResource = new TsFileResource(syncFile);
         tsFileResource.putStartTime(String.valueOf(i), (long) j * 10);
         tsFileResource.putEndTime(String.valueOf(i), (long) j * 10 + 5);
+        tsFileResource.setHistoricalVersions(Collections.singleton((long) j));
         tsFileResource.serialize();
       }
     }
@@ -225,6 +227,7 @@ public class FileLoaderTest {
         TsFileResource tsFileResource = new TsFileResource(syncFile);
         tsFileResource.putStartTime(String.valueOf(i), (long) j * 10);
         tsFileResource.putEndTime(String.valueOf(i), (long) j * 10 + 5);
+        tsFileResource.setHistoricalVersions(Collections.singleton((long) j));
         tsFileResource.serialize();
       }
     }

Reply via email to