From 1554c087f408aa8bbe34d37d001be8174b7a0e69 Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?=E5=88=98=E5=A8=81?=
<51618159+LIU-WEI-git@users.noreply.github.com>
Date: Wed, 30 Mar 2022 19:19:49 +0800
Subject: [PATCH] Refactor IoTDB storage plugin and bump up iotdb-session to
0.12.5 (#8755)
---
CHANGES.md | 2 +
dist-material/release-docs/LICENSE | 8 +-
oap-server-bom/pom.xml | 2 +-
.../storage-iotdb-plugin/pom.xml | 7 +-
.../storage/plugin/iotdb/IoTDBClient.java | 110 ++--------
.../plugin/iotdb/IoTDBStorageProvider.java | 8 +-
.../plugin/iotdb/IoTDBTableMetaInfo.java | 7 +-
.../plugin/iotdb/base/IoTDBInsertRequest.java | 75 ++-----
.../plugin/iotdb/base/IoTDBManagementDAO.java | 3 +-
.../plugin/iotdb/base/IoTDBMetricsDAO.java | 7 +-
.../plugin/iotdb/base/IoTDBNoneStreamDAO.java | 3 +-
.../plugin/iotdb/base/IoTDBRecordDAO.java | 3 +-
.../cache/IoTDBNetworkAddressAliasDAO.java | 9 +-
.../IoTDBUITemplateManagementDAO.java | 52 ++---
.../profile/IoTDBProfileTaskLogQueryDAO.java | 28 +--
.../profile/IoTDBProfileTaskQueryDAO.java | 50 +++--
.../IoTDBProfileThreadSnapshotQueryDAO.java | 69 ++++---
.../iotdb/query/IoTDBAggregationQueryDAO.java | 22 +-
.../iotdb/query/IoTDBAlarmQueryDAO.java | 15 +-
.../iotdb/query/IoTDBBrowserLogQueryDAO.java | 34 +++-
.../query/IoTDBEBPFProfilingDataDAO.java | 11 +-
.../query/IoTDBEBPFProfilingScheduleDAO.java | 18 +-
.../query/IoTDBEBPFProfilingTaskDAO.java | 21 +-
.../iotdb/query/IoTDBEventQueryDAO.java | 43 ++--
.../plugin/iotdb/query/IoTDBLogQueryDAO.java | 36 +++-
.../iotdb/query/IoTDBMetadataQueryDAO.java | 59 +++---
.../iotdb/query/IoTDBMetricsQueryDAO.java | 63 ++++--
.../iotdb/query/IoTDBTopNRecordsQueryDAO.java | 28 ++-
.../iotdb/query/IoTDBTopologyQueryDAO.java | 74 ++++---
.../iotdb/query/IoTDBTraceQueryDAO.java | 43 ++--
.../iotdb/utils/IoTDBDataConverter.java | 189 ++++++++++++++++++
.../plugin/iotdb/utils/IoTDBUtils.java | 70 +++++++
.../cases/alarm/iotdb/docker-compose.yml | 2 +-
.../cases/event/iotdb/docker-compose.yml | 2 +-
.../e2e-v2/cases/log/iotdb/docker-compose.yml | 2 +-
.../profiling/trace/iotdb/docker-compose.yml | 2 +-
.../cases/storage/iotdb/docker-compose.yml | 2 +-
.../e2e-v2/cases/ttl/iotdb/docker-compose.yml | 2 +-
.../known-oap-backend-dependencies.txt | 8 +-
39 files changed, 755 insertions(+), 434 deletions(-)
create mode 100644 oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/utils/IoTDBDataConverter.java
create mode 100644 oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/utils/IoTDBUtils.java
diff --git a/CHANGES.md b/CHANGES.md
index 4fb9b78fed..c8f90770fb 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -148,6 +148,8 @@ NOTICE, this sharding concept is NOT just for splitting data into different data
* Fix event type of export data is incorrect, it was `EventType.TOTAL` always.
* Reduce redundancy ThreadLocal in MAL core. Improve MAL performance.
* Trim tag's key and value in log query.
+* Refactor IoTDB storage plugin, add IoTDBDataConverter and fix ModifyCollectionInEnhancedForLoop bug.
+* Bump up iotdb-session to 0.12.5.
#### UI
diff --git a/dist-material/release-docs/LICENSE b/dist-material/release-docs/LICENSE
index 48af9ec7a7..c69e7f88cb 100755
--- a/dist-material/release-docs/LICENSE
+++ b/dist-material/release-docs/LICENSE
@@ -329,10 +329,10 @@ The text of each license is the standard Apache 2.0 license.
Armeria 1.14.1, http://github.com/line/armeria, Apache 2.0
Brotli4j 1.6.0, https://github.com/hyperxpro/Brotli4j, Apache 2.0
micrometer 1.8.2, https://github.com/micrometer-metrics/micrometer, Apache 2.0
- iotdb-session 0.12.4: https://github.com/apache/iotdb, Apache 2.0
- iotdb-thrift 0.12.4: https://github.com/apache/iotdb, Apache 2.0
- service-rpc 0.12.4: https://github.com/apache/iotdb, Apache 2.0
- tsfile 0.12.4 https://github.com/apache/iotdb Apache 2.0
+ iotdb-session 0.12.5: https://github.com/apache/iotdb, Apache 2.0
+ iotdb-thrift 0.12.5: https://github.com/apache/iotdb, Apache 2.0
+ service-rpc 0.12.5: https://github.com/apache/iotdb, Apache 2.0
+ tsfile 0.12.5 https://github.com/apache/iotdb Apache 2.0
libthrift 0.14.1: https://github.com/apache/thrift Apache 2.0
j2objc 1.3: https://github.com/google/j2objc Apache 2.0
diff --git a/oap-server-bom/pom.xml b/oap-server-bom/pom.xml
index 666123f85f..fee971fea5 100644
--- a/oap-server-bom/pom.xml
+++ b/oap-server-bom/pom.xml
@@ -77,7 +77,7 @@
3.0.0
4.4.13
1.21
- 0.12.4
+ 0.12.5
1.6.0
diff --git a/oap-server/server-storage-plugin/storage-iotdb-plugin/pom.xml b/oap-server/server-storage-plugin/storage-iotdb-plugin/pom.xml
index 40b863a8f1..2fd833b57e 100644
--- a/oap-server/server-storage-plugin/storage-iotdb-plugin/pom.xml
+++ b/oap-server/server-storage-plugin/storage-iotdb-plugin/pom.xml
@@ -17,7 +17,8 @@
~
-->
-
+
server-storage-plugin
org.apache.skywalking
@@ -62,8 +63,8 @@
- org.lz4
- lz4-java
+ org.lz4
+ lz4-java
diff --git a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBClient.java b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBClient.java
index 0d7f88644f..1c2a0496fb 100644
--- a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBClient.java
+++ b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBClient.java
@@ -21,8 +21,6 @@ package org.apache.skywalking.oap.server.storage.plugin.iotdb;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
import lombok.extern.slf4j.Slf4j;
import org.apache.iotdb.rpc.IoTDBConnectionException;
import org.apache.iotdb.rpc.StatementExecutionException;
@@ -32,18 +30,16 @@ import org.apache.iotdb.session.pool.SessionPool;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.read.common.Field;
import org.apache.iotdb.tsfile.read.common.RowRecord;
-import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
-import org.apache.skywalking.oap.server.core.analysis.manual.log.LogRecord;
-import org.apache.skywalking.oap.server.core.browser.manual.errorlog.BrowserErrorLogRecord;
-import org.apache.skywalking.oap.server.core.management.ui.template.UITemplate;
import org.apache.skywalking.oap.server.core.storage.StorageData;
-import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
+import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.client.healthcheck.DelegatedHealthChecker;
import org.apache.skywalking.oap.server.library.client.healthcheck.HealthCheckable;
import org.apache.skywalking.oap.server.library.util.HealthChecker;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.base.IoTDBInsertRequest;
+import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBDataConverter;
+import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
public class IoTDBClient implements Client, HealthCheckable {
@@ -69,11 +65,11 @@ public class IoTDBClient implements Client, HealthCheckable {
public void connect() throws IoTDBConnectionException, StatementExecutionException {
try {
final int sessionPoolSize = config.getSessionPoolSize() == 0 ?
- Runtime.getRuntime().availableProcessors() * 2 : config.getSessionPoolSize();
+ Runtime.getRuntime().availableProcessors() * 2 : config.getSessionPoolSize();
log.info("SessionPool Size: {}", sessionPoolSize);
- sessionPool = new SessionPool(config.getHost(), config.getRpcPort(), config.getUsername(),
- config.getPassword(), sessionPoolSize, false, false
- );
+ sessionPool = new SessionPool(config.getHost(), config.getRpcPort(),
+ config.getUsername(), config.getPassword(),
+ sessionPoolSize, false, false);
sessionPool.setStorageGroup(storageGroup);
healthChecker.health();
@@ -119,9 +115,11 @@ public class IoTDBClient implements Client, HealthCheckable {
devicePath.append(storageGroup).append(IoTDBClient.DOT).append(request.getModelName());
try {
// make an index value as a layer name of the storage path
+ // every storage path has a fix order which has been set in IoTDBTableMetaInfo
if (!request.getIndexes().isEmpty()) {
- request.getIndexValues().forEach(value -> devicePath.append(IoTDBClient.DOT)
- .append(indexValue2LayerName(value)));
+ request.getIndexValues().forEach(
+ value -> devicePath.append(IoTDBClient.DOT)
+ .append(IoTDBUtils.indexValue2LayerName(value)));
}
sessionPool.insertRecord(devicePath.toString(), request.getTime(),
request.getMeasurements(), request.getMeasurementTypes(),
@@ -158,8 +156,9 @@ public class IoTDBClient implements Client, HealthCheckable {
devicePath.append(storageGroup).append(IoTDBClient.DOT).append(request.getModelName());
// make an index value as a layer name of the storage path
if (!request.getIndexes().isEmpty()) {
- request.getIndexValues().forEach(value -> devicePath.append(IoTDBClient.DOT)
- .append(indexValue2LayerName(value)));
+ request.getIndexValues().forEach(
+ value -> devicePath.append(IoTDBClient.DOT)
+ .append(IoTDBUtils.indexValue2LayerName(value)));
}
devicePathList.add(devicePath.toString());
timeList.add(request.getTime());
@@ -188,7 +187,7 @@ public class IoTDBClient implements Client, HealthCheckable {
*/
public List super StorageData> filterQuery(String modelName, String querySQL,
StorageBuilder extends StorageData> storageBuilder)
- throws IOException {
+ throws IOException {
if (!querySQL.contains("align by device")) {
throw new IOException("querySQL must contain \"align by device\"");
}
@@ -200,53 +199,14 @@ public class IoTDBClient implements Client, HealthCheckable {
log.debug("SQL: {}, columnNames: {}", querySQL, wrapper.getColumnNames());
}
- List columnNames = wrapper.getColumnNames();
IoTDBTableMetaInfo tableMetaInfo = IoTDBTableMetaInfo.get(modelName);
List indexes = tableMetaInfo.getIndexes();
+ List columnNames = wrapper.getColumnNames();
while (wrapper.hasNext()) {
- Map map = new ConcurrentHashMap<>();
RowRecord rowRecord = wrapper.next();
- List fields = rowRecord.getFields();
- // transform timestamp to time_bucket
- if (!UITemplate.INDEX_NAME.equals(modelName)) {
- map.put(IoTDBClient.TIME_BUCKET, TimeBucket.getTimeBucket(
- rowRecord.getTimestamp(),
- tableMetaInfo.getModel().getDownsampling()
- ));
- }
- // field.get(0) -> Device, transform layerName to indexValue
- String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
- for (int i = 0; i < indexes.size(); i++) {
- map.put(indexes.get(i), layerName2IndexValue(layerNames[i + 1]));
- }
- for (int i = 0; i < columnNames.size() - 2; i++) {
- String columnName = columnNames.get(i + 2);
- Field field = fields.get(i + 1);
- if (field.getDataType() == null) {
- continue;
- }
- if (field.getDataType().equals(TSDataType.TEXT)) {
- map.put(columnName, field.getStringValue());
- } else {
- map.put(columnName, field.getObjectValue(field.getDataType()));
- }
- }
- if (map.containsKey(IoTDBIndexes.LAYER_IDX)) {
- String layer = (String) map.get(IoTDBIndexes.LAYER_IDX);
- map.put(IoTDBIndexes.LAYER_IDX, Integer.valueOf(layer));
- }
- if (modelName.equals(BrowserErrorLogRecord.INDEX_NAME) || modelName.equals(LogRecord.INDEX_NAME)) {
- map.put(IoTDBClient.TIMESTAMP, map.get("\"" + IoTDBClient.TIMESTAMP + "\""));
- }
- for (Map.Entry entry : map.entrySet()) {
- // remove double quotes
- String key = entry.getKey();
- if (key.contains(".")) {
- map.put(key.substring(1, key.length() - 1), entry.getValue());
- }
- }
-
- storageDataList.add(storageBuilder.storage2Entity(new HashMapConverter.ToEntity(map)));
+ Convert2Entity convert2Entity =
+ new IoTDBDataConverter.ToEntity(tableMetaInfo, indexes, columnNames, rowRecord);
+ storageDataList.add(storageBuilder.storage2Entity(convert2Entity));
}
healthChecker.health();
} catch (IoTDBConnectionException | StatementExecutionException e) {
@@ -316,35 +276,7 @@ public class IoTDBClient implements Client, HealthCheckable {
}
}
- public String indexValue2LayerName(String indexValue) {
- return "\"" + indexValue + "\"";
- }
-
- public String layerName2IndexValue(String layerName) {
- return layerName.substring(0, layerName.length() - 1);
- }
-
- public StringBuilder addQueryIndexValue(String modelName,
- StringBuilder query,
- Map indexAndValueMap) {
- List indexes = IoTDBTableMetaInfo.get(modelName).getIndexes();
- indexes.forEach(index -> {
- if (indexAndValueMap.containsKey(index)) {
- query.append(IoTDBClient.DOT).append(indexValue2LayerName(indexAndValueMap.get(index)));
- } else {
- query.append(IoTDBClient.DOT).append("*");
- }
- });
- return query;
- }
-
- public StringBuilder addQueryAsterisk(String modelName, StringBuilder query) {
- List indexes = IoTDBTableMetaInfo.get(modelName).getIndexes();
- indexes.forEach(index -> query.append(IoTDBClient.DOT).append("*"));
- return query;
- }
-
- public StringBuilder addModelPath(StringBuilder query, String modelName) {
- return query.append(storageGroup).append(IoTDBClient.DOT).append(modelName);
+ public String getStorageGroup() {
+ return storageGroup;
}
}
diff --git a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBStorageProvider.java b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBStorageProvider.java
index d9deec790b..fd0e5455ea 100644
--- a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBStorageProvider.java
+++ b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBStorageProvider.java
@@ -117,10 +117,10 @@ public class IoTDBStorageProvider extends ModuleProvider {
this.registerServiceImplementation(UITemplateManagementDAO.class, new IoTDBUITemplateManagementDAO(client));
this.registerServiceImplementation(IProfileTaskLogQueryDAO.class,
- new IoTDBProfileTaskLogQueryDAO(client, config.getFetchTaskLogMaxSize()));
+ new IoTDBProfileTaskLogQueryDAO(client, config.getFetchTaskLogMaxSize()));
this.registerServiceImplementation(IProfileTaskQueryDAO.class, new IoTDBProfileTaskQueryDAO(client));
this.registerServiceImplementation(IProfileThreadSnapshotQueryDAO.class,
- new IoTDBProfileThreadSnapshotQueryDAO(client));
+ new IoTDBProfileThreadSnapshotQueryDAO(client));
this.registerServiceImplementation(IAggregationQueryDAO.class, new IoTDBAggregationQueryDAO(client));
this.registerServiceImplementation(IAlarmQueryDAO.class, new IoTDBAlarmQueryDAO(client));
@@ -141,8 +141,8 @@ public class IoTDBStorageProvider extends ModuleProvider {
@Override
public void start() throws ServiceNotProvidedException, ModuleStartException {
MetricsCreator metricCreator = getManager().find(TelemetryModule.NAME)
- .provider()
- .getService(MetricsCreator.class);
+ .provider()
+ .getService(MetricsCreator.class);
HealthCheckMetrics healthChecker = metricCreator.createHealthCheckerGauge(
"storage_iotdb", MetricsTag.EMPTY_KEY, MetricsTag.EMPTY_VALUE);
client.registerChecker(healthChecker);
diff --git a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBTableMetaInfo.java b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBTableMetaInfo.java
index 4c17e4913c..2567cde7f3 100644
--- a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBTableMetaInfo.java
+++ b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/IoTDBTableMetaInfo.java
@@ -86,8 +86,11 @@ public class IoTDBTableMetaInfo {
indexes.add(IoTDBIndexes.AGENT_ID_INX);
}
- final IoTDBTableMetaInfo tableMetaInfo = IoTDBTableMetaInfo.builder().model(model)
- .columnAndTypeMap(columnAndTypeMap).indexes(indexes).build();
+ final IoTDBTableMetaInfo tableMetaInfo = IoTDBTableMetaInfo.builder()
+ .model(model)
+ .columnAndTypeMap(columnAndTypeMap)
+ .indexes(indexes)
+ .build();
TABLE_META_INFOS.put(model.getName(), tableMetaInfo);
}
diff --git a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/base/IoTDBInsertRequest.java b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/base/IoTDBInsertRequest.java
index 56de4f8f7a..7f8f7e1e98 100644
--- a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/base/IoTDBInsertRequest.java
+++ b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/base/IoTDBInsertRequest.java
@@ -19,23 +19,19 @@
package org.apache.skywalking.oap.server.storage.plugin.iotdb.base;
import java.util.ArrayList;
-import java.util.Iterator;
import java.util.List;
-import java.util.Map;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
import lombok.extern.slf4j.Slf4j;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.skywalking.oap.server.core.Const;
import org.apache.skywalking.oap.server.core.storage.StorageData;
-import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
-import org.apache.skywalking.oap.server.core.storage.type.StorageDataComplexObject;
import org.apache.skywalking.oap.server.library.client.request.InsertRequest;
import org.apache.skywalking.oap.server.library.client.request.UpdateRequest;
-import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
-import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBIndexes;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBTableMetaInfo;
+import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBDataConverter;
@Getter
@Setter
@@ -50,59 +46,24 @@ public class IoTDBInsertRequest implements InsertRequest, UpdateRequest {
private List measurementTypes;
private List