This is an automated email from the ASF dual-hosted git repository.
xingtanzjr 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 1dd5a200281 Speed up recover stage when DataNode restarting (#11421)
1dd5a200281 is described below
commit 1dd5a200281aff1f42a9d7abdbde9d172445d4f6
Author: Zhang.Jinrui <[email protected]>
AuthorDate: Tue Oct 31 12:02:11 2023 +0800
Speed up recover stage when DataNode restarting (#11421)
---
.../apache/iotdb/db/service/IoTDBShutdownHook.java | 5 ++-
.../iotdb/db/storageengine/StorageEngine.java | 5 +++
.../db/storageengine/dataregion/DataRegion.java | 36 +++++++++++++---------
.../dataregion/flush/FlushManager.java | 7 +++--
.../dataregion/memtable/TsFileProcessor.java | 14 ++++++---
5 files changed, 45 insertions(+), 22 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/IoTDBShutdownHook.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/IoTDBShutdownHook.java
index 36e679fe996..af5d90f5110 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/IoTDBShutdownHook.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/IoTDBShutdownHook.java
@@ -68,7 +68,10 @@ public class IoTDBShutdownHook extends Thread {
WALManager.getInstance().waitAllWALFlushed();
// flush data to Tsfile and remove WAL log files
- if (!IoTDBDescriptor.getInstance().getConfig().isClusterMode()) {
+ if (!IoTDBDescriptor.getInstance()
+ .getConfig()
+ .getDataRegionConsensusProtocolClass()
+ .equals(ConsensusFactory.RATIS_CONSENSUS)) {
StorageEngine.getInstance().syncCloseAllProcessor();
}
WALManager.getInstance().deleteOutdatedFilesInWALNodes();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
index 90470e024c0..dcfbc918f74 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
@@ -296,6 +296,11 @@ public class StorageEngine implements IService {
}
recover();
+ for (DataRegion dataRegion : dataRegionMap.values()) {
+ if (dataRegion != null) {
+ dataRegion.initCompaction();
+ }
+ }
ttlCheckThread =
IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.TTL_CHECK.getName());
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
index 96e32cbe428..513552d6ac1 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
@@ -134,6 +134,9 @@ import java.util.Map;
import java.util.Map.Entry;
import java.util.Set;
import java.util.TreeMap;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -559,9 +562,6 @@ public class DataRegion implements IDataRegionForQuery {
throw new DataRegionException(e);
}
- // recover and start timed compaction thread
- initCompaction();
-
if (StorageEngine.getInstance().isAllSgReady()) {
logger.info("The data region {}[{}] is created successfully",
databaseName, dataRegionId);
} else {
@@ -583,7 +583,7 @@ public class DataRegion implements IDataRegionForQuery {
}
}
- private void initCompaction() {
+ public void initCompaction() {
if (!config.isEnableSeqSpaceCompaction()
&& !config.isEnableUnseqSpaceCompaction()
&& !config.isEnableCrossSpaceCompaction()) {
@@ -1366,21 +1366,21 @@ public class DataRegion implements IDataRegionForQuery {
* @param sequence whether this tsfile processor is sequence or not
* @param tsFileProcessor tsfile processor
*/
- public void asyncCloseOneTsFileProcessor(boolean sequence, TsFileProcessor
tsFileProcessor) {
+ public Future<?> asyncCloseOneTsFileProcessor(boolean sequence,
TsFileProcessor tsFileProcessor) {
// for sequence tsfile, we update the endTimeMap only when the file is
prepared to be closed.
// for unsequence tsfile, we have maintained the endTimeMap when an
insertion comes.
if (closingSequenceTsFileProcessor.contains(tsFileProcessor)
|| closingUnSequenceTsFileProcessor.contains(tsFileProcessor)
|| tsFileProcessor.alreadyMarkedClosing()) {
- return;
+ return CompletableFuture.completedFuture(null);
}
logger.info(
"Async close tsfile: {}",
tsFileProcessor.getTsFileResource().getTsFile().getAbsolutePath());
-
+ Future<?> future;
if (sequence) {
closingSequenceTsFileProcessor.add(tsFileProcessor);
- tsFileProcessor.asyncClose();
+ future = tsFileProcessor.asyncClose();
workSequenceTsFileProcessors.remove(tsFileProcessor.getTimeRangeId());
// if unsequence files don't contain this time range id, we should
remove it's version
@@ -1391,7 +1391,7 @@ public class DataRegion implements IDataRegionForQuery {
logger.info("close a sequence tsfile processor {}", databaseName + "-" +
dataRegionId);
} else {
closingUnSequenceTsFileProcessor.add(tsFileProcessor);
- tsFileProcessor.asyncClose();
+ future = tsFileProcessor.asyncClose();
workUnsequenceTsFileProcessors.remove(tsFileProcessor.getTimeRangeId());
// if sequence files don't contain this time range id, we should remove
it's version
@@ -1400,6 +1400,7 @@ public class DataRegion implements IDataRegionForQuery {
timePartitionIdVersionControllerMap.remove(tsFileProcessor.getTimeRangeId());
}
}
+ return future;
}
/**
@@ -1594,7 +1595,7 @@ public class DataRegion implements IDataRegionForQuery {
public void syncCloseAllWorkingTsFileProcessors() {
synchronized (closeStorageGroupCondition) {
try {
- asyncCloseAllWorkingTsFileProcessors();
+ List<Future<?>> tsFileProcessorsClosingFutures =
asyncCloseAllWorkingTsFileProcessors();
long startTime = System.currentTimeMillis();
while (!closingSequenceTsFileProcessor.isEmpty()
|| !closingUnSequenceTsFileProcessor.isEmpty()) {
@@ -1606,7 +1607,12 @@ public class DataRegion implements IDataRegionForQuery {
(System.currentTimeMillis() - startTime) / 1000);
}
}
- } catch (InterruptedException e) {
+ for (Future<?> f : tsFileProcessorsClosingFutures) {
+ if (f != null) {
+ f.get();
+ }
+ }
+ } catch (InterruptedException | ExecutionException e) {
logger.error(
"CloseFileNodeCondition error occurs while waiting for closing the
storage "
+ "group {}",
@@ -1618,23 +1624,25 @@ public class DataRegion implements IDataRegionForQuery {
}
/** close all working tsfile processors */
- public void asyncCloseAllWorkingTsFileProcessors() {
+ public List<Future<?>> asyncCloseAllWorkingTsFileProcessors() {
writeLock("asyncCloseAllWorkingTsFileProcessors");
+ List<Future<?>> futures = new ArrayList<>();
try {
logger.info("async force close all files in database: {}", databaseName
+ "-" + dataRegionId);
// to avoid concurrent modification problem, we need a new array list
for (TsFileProcessor tsFileProcessor :
new ArrayList<>(workSequenceTsFileProcessors.values())) {
- asyncCloseOneTsFileProcessor(true, tsFileProcessor);
+ futures.add(asyncCloseOneTsFileProcessor(true, tsFileProcessor));
}
// to avoid concurrent modification problem, we need a new array list
for (TsFileProcessor tsFileProcessor :
new ArrayList<>(workUnsequenceTsFileProcessors.values())) {
- asyncCloseOneTsFileProcessor(false, tsFileProcessor);
+ futures.add(asyncCloseOneTsFileProcessor(false, tsFileProcessor));
}
} finally {
writeUnlock();
}
+ return futures;
}
/** force close all working tsfile processors */
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/FlushManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/FlushManager.java
index f00d85676a5..87f0440c3be 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/FlushManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/flush/FlushManager.java
@@ -32,7 +32,9 @@ import
org.apache.iotdb.db.storageengine.dataregion.memtable.TsFileProcessor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentLinkedDeque;
+import java.util.concurrent.Future;
@SuppressWarnings("squid:S6548")
public class FlushManager implements FlushManagerMBean, IService {
@@ -120,7 +122,7 @@ public class FlushManager implements FlushManagerMBean,
IService {
* @param tsFileProcessor tsFileProcessor to be flushed
*/
@SuppressWarnings("squid:S2445")
- public void registerTsFileProcessor(TsFileProcessor tsFileProcessor) {
+ public Future<?> registerTsFileProcessor(TsFileProcessor tsFileProcessor) {
synchronized (tsFileProcessor) {
if (tsFileProcessor.isManagedByFlushManager()) {
LOGGER.debug(
@@ -138,7 +140,7 @@ public class FlushManager implements FlushManagerMBean,
IService {
tsFileProcessorQueue.size());
}
tsFileProcessor.setManagedByFlushManager(true);
- flushPool.submit(new FlushThread());
+ return flushPool.submit(new FlushThread());
} else {
if (LOGGER.isDebugEnabled()) {
LOGGER.debug(
@@ -148,6 +150,7 @@ public class FlushManager implements FlushManagerMBean,
IService {
}
}
}
+ return CompletableFuture.completedFuture(null);
}
private FlushManager() {}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
index 1c01ced6695..7fadcda4267 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
@@ -83,8 +83,10 @@ import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentLinkedDeque;
import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.Future;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
@@ -860,7 +862,7 @@ public class TsFileProcessor {
}
/** async close one tsfile, register and close it by another thread */
- public void asyncClose() {
+ public Future<?> asyncClose() {
flushQueryLock.writeLock().lock();
if (logger.isDebugEnabled()) {
logger.debug(
@@ -889,7 +891,7 @@ public class TsFileProcessor {
}
if (shouldClose) {
- return;
+ return CompletableFuture.completedFuture(null);
}
// when a flush thread serves this TsFileProcessor (because the
processor is submitted by
// registerTsFileProcessor()), the thread will seal the corresponding
TsFile and
@@ -910,9 +912,10 @@ public class TsFileProcessor {
// When invoke closing TsFile after insert data to memTable, we
shouldn't flush until invoke
// flushing memTable in System module.
- addAMemtableIntoFlushingList(tmpMemTable);
+ Future<?> future = addAMemtableIntoFlushingList(tmpMemTable);
logger.info("Memtable {} has been added to flushing list",
tmpMemTable);
shouldClose = true;
+ return future;
} catch (Exception e) {
logger.error(
"{}: {} async close failed, because",
@@ -927,6 +930,7 @@ public class TsFileProcessor {
FLUSH_QUERY_WRITE_RELEASE, storageGroupName,
tsFileResource.getTsFile().getName());
}
}
+ return CompletableFuture.completedFuture(null);
}
/**
@@ -1015,7 +1019,7 @@ public class TsFileProcessor {
* queue, set the current working memtable as null and then register the
tsfileProcessor into the
* flushManager again.
*/
- private void addAMemtableIntoFlushingList(IMemTable tobeFlushed) throws
IOException {
+ private Future<?> addAMemtableIntoFlushingList(IMemTable tobeFlushed) throws
IOException {
Map<String, Long> lastTimeForEachDevice = new HashMap<>();
if (sequence) {
lastTimeForEachDevice = tobeFlushed.getMaxTime();
@@ -1054,7 +1058,7 @@ public class TsFileProcessor {
totalMemTableSize += tobeFlushed.memSize();
}
workMemTable = null;
- FlushManager.getInstance().registerTsFileProcessor(this);
+ return FlushManager.getInstance().registerTsFileProcessor(this);
}
/** put back the memtable to MemTablePool and make metadata in writer
visible */