This is an automated email from the ASF dual-hosted git repository. xingtanzjr pushed a commit to branch speed_up_restart_recover in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 999d38b319d789b64ff3bd9efb0f6a27f97d4d5b Author: Jinrui.Zhang <[email protected]> AuthorDate: Mon Oct 30 17:08:20 2023 +0800 using CompletableFultre.completedFuture to replace null --- .../src/main/java/org/apache/iotdb/SessionExample.java | 17 +++++++++-------- .../iotdb/db/storageengine/dataregion/DataRegion.java | 7 ++++--- .../db/storageengine/dataregion/flush/FlushManager.java | 3 ++- .../dataregion/memtable/TsFileProcessor.java | 5 +++-- .../db/storageengine/dataregion/wal/WALManager.java | 6 +++--- 5 files changed, 21 insertions(+), 17 deletions(-) diff --git a/example/session/src/main/java/org/apache/iotdb/SessionExample.java b/example/session/src/main/java/org/apache/iotdb/SessionExample.java index 8d1deb764c5..60bf282ed8f 100644 --- a/example/session/src/main/java/org/apache/iotdb/SessionExample.java +++ b/example/session/src/main/java/org/apache/iotdb/SessionExample.java @@ -26,6 +26,7 @@ import org.apache.iotdb.isession.template.Template; import org.apache.iotdb.isession.util.Version; import org.apache.iotdb.rpc.IoTDBConnectionException; import org.apache.iotdb.rpc.StatementExecutionException; +import org.apache.iotdb.rpc.TSStatusCode; import org.apache.iotdb.session.Session; import org.apache.iotdb.session.template.MeasurementNode; import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType; @@ -77,13 +78,13 @@ public class SessionExample { // set session fetchSize session.setFetchSize(10000); - // try { - //// session.createDatabase("root.sg1"); - // } catch (StatementExecutionException e) { - // if (e.getStatusCode() != TSStatusCode.DATABASE_ALREADY_EXISTS.getStatusCode()) { - // throw e; - // } - // } + try { + session.createDatabase("root.sg1"); + } catch (StatementExecutionException e) { + if (e.getStatusCode() != TSStatusCode.DATABASE_ALREADY_EXISTS.getStatusCode()) { + throw e; + } + } // createTemplate(); createTimeseries(); @@ -400,7 +401,7 @@ public class SessionExample { // Method 1 to add tablet data long timestamp = System.currentTimeMillis(); - for (long row = 0; row < 10000000; row++) { + for (long row = 0; row < 100; row++) { int rowIndex = tablet.rowSize++; tablet.addTimestamp(rowIndex, timestamp); for (int s = 0; s < 3; s++) { 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 87d72ab0692..13bea6f0374 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,7 @@ 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; @@ -1371,7 +1372,7 @@ public class DataRegion implements IDataRegionForQuery { if (closingSequenceTsFileProcessor.contains(tsFileProcessor) || closingUnSequenceTsFileProcessor.contains(tsFileProcessor) || tsFileProcessor.alreadyMarkedClosing()) { - return null; + return CompletableFuture.completedFuture(null); } logger.info( "Async close tsfile: {}", @@ -1594,7 +1595,7 @@ public class DataRegion implements IDataRegionForQuery { public void syncCloseAllWorkingTsFileProcessors() { synchronized (closeStorageGroupCondition) { try { - List<Future<?>> futures = asyncCloseAllWorkingTsFileProcessors(); + List<Future<?>> tsFileProcessorsClosingFutures = asyncCloseAllWorkingTsFileProcessors(); long startTime = System.currentTimeMillis(); while (!closingSequenceTsFileProcessor.isEmpty() || !closingUnSequenceTsFileProcessor.isEmpty()) { @@ -1606,7 +1607,7 @@ public class DataRegion implements IDataRegionForQuery { (System.currentTimeMillis() - startTime) / 1000); } } - for (Future<?> f : futures) { + for (Future<?> f : tsFileProcessorsClosingFutures) { if (f != null) { f.get(); } 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 08a6a6dd624..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,6 +32,7 @@ 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; @@ -149,7 +150,7 @@ public class FlushManager implements FlushManagerMBean, IService { } } } - return null; + 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 5215f10e842..487eb8a00fe 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,6 +83,7 @@ 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; @@ -884,7 +885,7 @@ public class TsFileProcessor { } if (shouldClose) { - return null; + 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 @@ -923,7 +924,7 @@ public class TsFileProcessor { FLUSH_QUERY_WRITE_RELEASE, storageGroupName, tsFileResource.getTsFile().getName()); } } - return null; + return CompletableFuture.completedFuture(null); } /** diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/WALManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/WALManager.java index 031fafc4fd3..6ccdb2aa3d3 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/WALManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/WALManager.java @@ -21,6 +21,7 @@ package org.apache.iotdb.db.storageengine.dataregion.wal; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.exception.StartupException; import org.apache.iotdb.commons.service.IService; @@ -273,9 +274,8 @@ public class WALManager implements IService { private void registerScheduleTask(long initDelayMs, long periodMs) { walDeleteThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.WAL_DELETE.getName()); - // ScheduledExecutorUtil.safelyScheduleWithFixedDelay( - // walDeleteThread, this::deleteOutdatedFiles, initDelayMs, periodMs, - // TimeUnit.MILLISECONDS); + ScheduledExecutorUtil.safelyScheduleWithFixedDelay( + walDeleteThread, this::deleteOutdatedFiles, initDelayMs, periodMs, TimeUnit.MILLISECONDS); } @TestOnly
