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();
}
}