This is an automated email from the ASF dual-hosted git repository.

qiaojialin pushed a commit to branch rel/0.12
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/rel/0.12 by this push:
     new 20362e4d5b [To rel/0.12][IOTDB-3306] Use rpc port check to avoid 
starting same IoTDB twice (#6075)
20362e4d5b is described below

commit 20362e4d5bd93cdacb7439a563352dd629e97b0b
Author: Haonan <[email protected]>
AuthorDate: Tue May 31 17:36:33 2022 +0800

    [To rel/0.12][IOTDB-3306] Use rpc port check to avoid starting same IoTDB 
twice (#6075)
---
 .../db/engine/storagegroup/TsFileProcessor.java    |  2 +-
 .../java/org/apache/iotdb/db/service/IoTDB.java    |  9 +++----
 .../db/service/thrift/ThriftServiceThread.java     | 18 +------------
 .../compaction/LevelCompactionMergeTest.java       | 31 +++++++++++++---------
 4 files changed, 24 insertions(+), 36 deletions(-)

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 859ece653c..d473c36406 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
@@ -818,7 +818,7 @@ public class TsFileProcessor {
         MemTableFlushTask flushTask =
             new MemTableFlushTask(memTableToFlush, writer, storageGroupName);
         flushTask.syncFlushMemTable();
-      } catch (Exception e) {
+      } catch (Throwable e) {
         if (writer == null) {
           logger.info(
               "{}: {} is closed during flush, abandon flush task",
diff --git a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java 
b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
index a35f70bd0a..c43c7b2f16 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
@@ -106,6 +106,10 @@ public class IoTDB implements IoTDBMBean {
 
     Runtime.getRuntime().addShutdownHook(new IoTDBShutdownHook());
     setUncaughtExceptionHandler();
+    // in cluster mode, RPC service is not enabled.
+    if (IoTDBDescriptor.getInstance().getConfig().isEnableRpcService()) {
+      registerManager.register(RPCService.getInstance());
+    }
     logger.info("recover the schema...");
     initMManager();
     registerManager.register(JMXService.getInstance());
@@ -121,11 +125,6 @@ public class IoTDB implements IoTDBMBean {
     registerManager.register(UDFClassLoaderManager.getInstance());
     registerManager.register(UDFRegistrationService.getInstance());
 
-    // in cluster mode, RPC service is not enabled.
-    if (IoTDBDescriptor.getInstance().getConfig().isEnableRpcService()) {
-      registerManager.register(RPCService.getInstance());
-    }
-
     if (IoTDBDescriptor.getInstance().getConfig().isEnableMetricService()) {
       registerManager.register(MetricsService.getInstance());
     }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
index 612d187195..c928cc6bf7 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
@@ -108,23 +108,7 @@ public class ThriftServiceThread extends Thread {
 
   @SuppressWarnings("java:S2259")
   public TServerTransport openTransport(String bindAddress, int port) throws 
TTransportException {
-    int maxRetry = 5;
-    long retryIntervalMS = 5000;
-    TTransportException lastExp = null;
-    for (int i = 0; i < maxRetry; i++) {
-      try {
-        return new TServerSocket(new InetSocketAddress(bindAddress, port));
-      } catch (TTransportException e) {
-        lastExp = e;
-        try {
-          Thread.sleep(retryIntervalMS);
-        } catch (InterruptedException interruptedException) {
-          Thread.currentThread().interrupt();
-          break;
-        }
-      }
-    }
-    throw lastExp;
+    return new TServerSocket(new InetSocketAddress(bindAddress, port));
   }
 
   public void setThreadStopLatch(CountDownLatch threadStopLatch) {
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
index 5e2ebfc835..5e14aaf9ae 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
@@ -440,18 +440,23 @@ public class LevelCompactionMergeTest extends 
LevelCompactionTest {
         TSEncoding.PLAIN,
         CompressionType.UNCOMPRESSED,
         Collections.emptyMap());
-    
IoTDBDescriptor.getInstance().getConfig().setMergeWriteThroughputMbPerSec(1024);
-    TsFileResource targetResource =
-        new TsFileResource(
-            new File(resources.get(0).getTsFile().getParentFile(), 
"0-0-1-0.tsfile"));
-    CompactionUtils.merge(
-        targetResource,
-        resources,
-        "root.test",
-        null,
-        new HashSet<>(),
-        true,
-        Collections.EMPTY_LIST,
-        null);
+    int mergeRate = 
IoTDBDescriptor.getInstance().getConfig().getMergeWriteThroughputMbPerSec();
+    try {
+      
IoTDBDescriptor.getInstance().getConfig().setMergeWriteThroughputMbPerSec(1024);
+      TsFileResource targetResource =
+          new TsFileResource(
+              new File(resources.get(0).getTsFile().getParentFile(), 
"0-0-1-0.tsfile"));
+      CompactionUtils.merge(
+          targetResource,
+          resources,
+          "root.test",
+          null,
+          new HashSet<>(),
+          true,
+          Collections.EMPTY_LIST,
+          null);
+    } finally {
+      
IoTDBDescriptor.getInstance().getConfig().setMergeWriteThroughputMbPerSec(mergeRate);
+    }
   }
 }

Reply via email to