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