This is an automated email from the ASF dual-hosted git repository.
xingtanzjr pushed a commit to branch rel/0.13
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/0.13 by this push:
new 417c80bfb7 [To rel/0.13][IOTDB-3937][IOTDB-3974] Fix the problem of
IoTDB Reporter in metric framework (#6793)
417c80bfb7 is described below
commit 417c80bfb76dab6b23b122306862ffac052e7680
Author: ZhangHongYin <[email protected]>
AuthorDate: Thu Jul 28 18:26:52 2022 +0800
[To rel/0.13][IOTDB-3937][IOTDB-3974] Fix the problem of IoTDB Reporter in
metric framework (#6793)
---
.../metrics/dropwizard/reporter/IoTDBReporter.java | 43 +++++++++-------------
.../apache/iotdb/metrics/config/MetricConfig.java | 13 +++----
.../{MetricsUtils.java => IoTDBMetricsUtils.java} | 28 ++++++++++++--
.../micrometer/reporter/IoTDBMeterRegistry.java | 37 ++++++-------------
.../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 4 ++
5 files changed, 64 insertions(+), 61 deletions(-)
diff --git
a/metrics/dropwizard-metrics/src/main/java/org/apache/iotdb/metrics/dropwizard/reporter/IoTDBReporter.java
b/metrics/dropwizard-metrics/src/main/java/org/apache/iotdb/metrics/dropwizard/reporter/IoTDBReporter.java
index 06693b7491..5786f45f0d 100644
---
a/metrics/dropwizard-metrics/src/main/java/org/apache/iotdb/metrics/dropwizard/reporter/IoTDBReporter.java
+++
b/metrics/dropwizard-metrics/src/main/java/org/apache/iotdb/metrics/dropwizard/reporter/IoTDBReporter.java
@@ -22,10 +22,10 @@ package org.apache.iotdb.metrics.dropwizard.reporter;
import org.apache.iotdb.metrics.config.MetricConfig;
import org.apache.iotdb.metrics.config.MetricConfigDescriptor;
import org.apache.iotdb.metrics.dropwizard.MetricName;
-import org.apache.iotdb.metrics.utils.MetricsUtils;
+import org.apache.iotdb.metrics.utils.IoTDBMetricsUtils;
import org.apache.iotdb.rpc.IoTDBConnectionException;
import org.apache.iotdb.rpc.StatementExecutionException;
-import org.apache.iotdb.session.Session;
+import org.apache.iotdb.session.pool.SessionPool;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import com.codahale.metrics.Counter;
@@ -40,7 +40,12 @@ import com.codahale.metrics.Timer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.util.*;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.SortedMap;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@@ -51,7 +56,7 @@ public class IoTDBReporter extends ScheduledReporter {
private static final TimeUnit DURATION_UNIT = TimeUnit.MILLISECONDS;
private static final TimeUnit RATE_UNIT = TimeUnit.SECONDS;
private final String prefix;
- private final Session session;
+ private final SessionPool sessionPool;
protected IoTDBReporter(
MetricRegistry registry,
@@ -68,34 +73,21 @@ public class IoTDBReporter extends ScheduledReporter {
executor,
shutdownExecutorOnStop);
this.prefix = prefix;
- this.session =
- new Session(
+ this.sessionPool =
+ new SessionPool(
ioTDBReporterConfig.getHost(),
ioTDBReporterConfig.getPort(),
ioTDBReporterConfig.getUsername(),
ioTDBReporterConfig.getPassword(),
- true);
- }
-
- @Override
- public void start(long period, TimeUnit unit) {
- super.start(period, unit);
- try {
- session.open();
- } catch (IoTDBConnectionException e) {
- logger.error("Failed to add session", e);
- }
+ 3);
+ IoTDBMetricsUtils.checkOrCreateStorageGroup(sessionPool);
}
@Override
public void stop() {
super.stop();
- try {
- if (session != null) {
- session.close();
- }
- } catch (IoTDBConnectionException e) {
- logger.error("Failed to close session.");
+ if (sessionPool != null) {
+ sessionPool.close();
}
}
@@ -241,7 +233,7 @@ public class IoTDBReporter extends ScheduledReporter {
private void updateValue(String name, Map<String, String> labels, Object
value) {
if (value != null) {
- String deviceId = MetricsUtils.generatePath(name, labels);
+ String deviceId = IoTDBMetricsUtils.generatePath(name, labels);
List<String> sensors = Collections.singletonList("value");
List<TSDataType> dataTypes = new ArrayList<>();
@@ -257,9 +249,10 @@ public class IoTDBReporter extends ScheduledReporter {
dataTypes.add(TSDataType.TEXT);
value = value.toString();
}
+ List<Object> values = Collections.singletonList(value);
try {
- session.insertRecord(deviceId, System.currentTimeMillis(), sensors,
dataTypes, value);
+ sessionPool.insertRecord(deviceId, System.currentTimeMillis(),
sensors, dataTypes, values);
} catch (IoTDBConnectionException | StatementExecutionException e) {
logger.warn("Failed to insert record");
}
diff --git
a/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java
b/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java
index 17ce4b5588..53df7ed3dd 100644
---
a/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java
+++
b/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java
@@ -152,6 +152,11 @@ public class MetricConfig {
ioTDBReporterConfig = newMetricConfig.ioTDBReporterConfig;
}
+ public void updateInstance(String instanceHost, Integer instancePort) {
+ this.instanceHost = instanceHost;
+ this.instancePort = instancePort;
+ }
+
public Boolean getEnableMetric() {
return enableMetric;
}
@@ -220,18 +225,10 @@ public class MetricConfig {
return instanceHost;
}
- public void setInstanceHost(String instanceHost) {
- this.instanceHost = instanceHost;
- }
-
public Integer getInstancePort() {
return instancePort;
}
- public void setInstancePort(Integer instancePort) {
- this.instancePort = instancePort;
- }
-
@Override
public boolean equals(Object obj) {
if (!(obj instanceof MetricConfig)) {
diff --git
a/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/MetricsUtils.java
b/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/IoTDBMetricsUtils.java
similarity index 64%
rename from
metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/MetricsUtils.java
rename to
metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/IoTDBMetricsUtils.java
index ec359191c2..87d29d3e99 100644
---
a/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/MetricsUtils.java
+++
b/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/IoTDBMetricsUtils.java
@@ -21,18 +21,27 @@ package org.apache.iotdb.metrics.utils;
import org.apache.iotdb.metrics.config.MetricConfig;
import org.apache.iotdb.metrics.config.MetricConfigDescriptor;
+import org.apache.iotdb.rpc.IoTDBConnectionException;
+import org.apache.iotdb.rpc.StatementExecutionException;
+import org.apache.iotdb.session.pool.SessionDataSetWrapper;
+import org.apache.iotdb.session.pool.SessionPool;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.util.Map;
-public class MetricsUtils {
+public class IoTDBMetricsUtils {
+ private static final Logger logger =
LoggerFactory.getLogger(IoTDBMetricsUtils.class);
private static final MetricConfig metricConfig =
MetricConfigDescriptor.getInstance().getMetricConfig();
+ private static final String STORAGE_GROUP =
+ "root." + metricConfig.getIoTDBReporterConfig().getDatabase();
public static String generatePath(String name, Map<String, String> labels) {
StringBuilder stringBuilder = new StringBuilder();
stringBuilder
- .append("root.")
- .append(metricConfig.getIoTDBReporterConfig().getDatabase())
+ .append(STORAGE_GROUP)
.append(".\"")
.append(metricConfig.getInstanceHost())
.append(":")
@@ -53,4 +62,17 @@ public class MetricsUtils {
}
return stringBuilder.toString();
}
+
+ public static void checkOrCreateStorageGroup(SessionPool session) {
+ try (SessionDataSetWrapper result =
+ session.executeQueryStatement("show storage group " + STORAGE_GROUP)) {
+ if (!result.hasNext()) {
+ session.setStorageGroup(STORAGE_GROUP);
+ }
+ } catch (IoTDBConnectionException e) {
+ logger.error("CheckOrCreateStorageGroup failed because ", e);
+ } catch (StatementExecutionException e) {
+ // do nothing
+ }
+ }
}
diff --git
a/metrics/micrometer-metrics/src/main/java/org/apache/iotdb/metrics/micrometer/reporter/IoTDBMeterRegistry.java
b/metrics/micrometer-metrics/src/main/java/org/apache/iotdb/metrics/micrometer/reporter/IoTDBMeterRegistry.java
index 9ec5a3816f..d780c0800e 100644
---
a/metrics/micrometer-metrics/src/main/java/org/apache/iotdb/metrics/micrometer/reporter/IoTDBMeterRegistry.java
+++
b/metrics/micrometer-metrics/src/main/java/org/apache/iotdb/metrics/micrometer/reporter/IoTDBMeterRegistry.java
@@ -21,10 +21,10 @@ package org.apache.iotdb.metrics.micrometer.reporter;
import org.apache.iotdb.metrics.config.MetricConfig;
import org.apache.iotdb.metrics.config.MetricConfigDescriptor;
-import org.apache.iotdb.metrics.utils.MetricsUtils;
+import org.apache.iotdb.metrics.utils.IoTDBMetricsUtils;
import org.apache.iotdb.rpc.IoTDBConnectionException;
import org.apache.iotdb.rpc.StatementExecutionException;
-import org.apache.iotdb.session.Session;
+import org.apache.iotdb.session.pool.SessionPool;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import io.micrometer.core.instrument.Clock;
@@ -40,45 +40,31 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
-import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit;
public class IoTDBMeterRegistry extends StepMeterRegistry {
private static final Logger logger =
LoggerFactory.getLogger(IoTDBMeterRegistry.class);
private static final MetricConfig.IoTDBReporterConfig ioTDBReporterConfig =
MetricConfigDescriptor.getInstance().getMetricConfig().getIoTDBReporterConfig();
- private final Session session;
+ private final SessionPool sessionPool;
public IoTDBMeterRegistry(StepRegistryConfig config, Clock clock) {
super(config, clock);
- session =
- new Session(
+ this.sessionPool =
+ new SessionPool(
ioTDBReporterConfig.getHost(),
ioTDBReporterConfig.getPort(),
ioTDBReporterConfig.getUsername(),
ioTDBReporterConfig.getPassword(),
- true);
- }
-
- @Override
- public void start(ThreadFactory threadFactory) {
- super.start(threadFactory);
- try {
- session.open();
- } catch (IoTDBConnectionException e) {
- logger.error("Failed to add session", e);
- }
+ 3);
+ IoTDBMetricsUtils.checkOrCreateStorageGroup(sessionPool);
}
@Override
public void stop() {
super.stop();
- try {
- if (session != null) {
- session.close();
- }
- } catch (IoTDBConnectionException e) {
- logger.error("Failed to close session.");
+ if (sessionPool != null) {
+ sessionPool.close();
}
}
@@ -143,12 +129,13 @@ public class IoTDBMeterRegistry extends StepMeterRegistry
{
private void updateValue(String name, Map<String, String> labels, Double
value, Long time) {
if (value != null) {
- String deviceId = MetricsUtils.generatePath(name, labels);
+ String deviceId = IoTDBMetricsUtils.generatePath(name, labels);
List<String> sensors = Collections.singletonList("value");
List<TSDataType> dataTypes =
Collections.singletonList(TSDataType.DOUBLE);
+ List<Object> values = Collections.singletonList(value);
try {
- session.insertRecord(deviceId, time, sensors, dataTypes, value);
+ sessionPool.insertRecord(deviceId, time, sensors, dataTypes, values);
} catch (IoTDBConnectionException | StatementExecutionException e) {
logger.warn("Failed to insert record");
}
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index 8d1f1f9b8e..0aff2ac6e1 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -927,6 +927,10 @@ public class IoTDBDescriptor {
} finally {
// update all data seriesPath
conf.updatePath();
+ // update instance in metric
+ MetricConfigDescriptor.getInstance()
+ .getMetricConfig()
+ .updateInstance(conf.getRpcAddress(), conf.getRpcPort());
}
}