Refactor IoTDB storage plugin and bump up iotdb-session to 0.12.5 (#8755)

This commit is contained in:
刘威 2022-03-30 19:19:49 +08:00 committed by GitHub
parent 2a8bf4deff
commit 1554c087f4
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
39 changed files with 755 additions and 434 deletions

View File

@ -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

View File

@ -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

View File

@ -77,7 +77,7 @@
<awaitility.version>3.0.0</awaitility.version>
<httpcore.version>4.4.13</httpcore.version>
<commons-compress.version>1.21</commons-compress.version>
<iotdb-session.version>0.12.4</iotdb-session.version>
<iotdb-session.version>0.12.5</iotdb-session.version>
<lz4-java.version>1.6.0</lz4-java.version>
</properties>

View File

@ -17,7 +17,8 @@
~
-->
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>server-storage-plugin</artifactId>
<groupId>org.apache.skywalking</groupId>
@ -62,8 +63,8 @@
</exclusions>
</dependency>
<dependency>
<groupId>org.lz4</groupId>
<artifactId>lz4-java</artifactId>
<groupId>org.lz4</groupId>
<artifactId>lz4-java</artifactId>
</dependency>
</dependencies>
</project>

View File

@ -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<String> columnNames = wrapper.getColumnNames();
IoTDBTableMetaInfo tableMetaInfo = IoTDBTableMetaInfo.get(modelName);
List<String> indexes = tableMetaInfo.getIndexes();
List<String> columnNames = wrapper.getColumnNames();
while (wrapper.hasNext()) {
Map<String, Object> map = new ConcurrentHashMap<>();
RowRecord rowRecord = wrapper.next();
List<Field> 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<String, Object> 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<String, String> indexAndValueMap) {
List<String> 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<String> 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;
}
}

View File

@ -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);

View File

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

View File

@ -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<TSDataType> measurementTypes;
private List<Object> measurementValues;
public <T extends StorageData> IoTDBInsertRequest(String modelName, long time, T storageData,
StorageBuilder<T> storageBuilder) {
public IoTDBInsertRequest(String modelName, long time) {
this.modelName = modelName;
this.time = time;
indexes = IoTDBTableMetaInfo.get(modelName).getIndexes();
indexValues = new ArrayList<>(indexes.size());
final HashMapConverter.ToStorage toStorage = new HashMapConverter.ToStorage();
this.indexes = IoTDBTableMetaInfo.get(modelName).getIndexes();
this.indexValues = new ArrayList<>(indexes.size());
this.indexes.forEach(index -> this.indexValues.add(Const.EMPTY_STRING));
int measurementsSize = IoTDBTableMetaInfo.get(modelName).getColumnAndTypeMap().size();
this.measurements = new ArrayList<>(measurementsSize);
this.measurementTypes = new ArrayList<>(measurementsSize);
this.measurementValues = new ArrayList<>(measurementsSize);
}
public static <T extends StorageData> IoTDBInsertRequest buildRequest(String modelName, long time, T storageData,
StorageBuilder<T> storageBuilder) {
String id = storageData.id();
IoTDBDataConverter.ToStorage toStorage = new IoTDBDataConverter.ToStorage(modelName, time, id);
storageBuilder.entity2Storage(storageData, toStorage);
Map<String, Object> storageMap = toStorage.obtain();
indexes.forEach(index -> {
if (index.equals(IoTDBIndexes.ID_IDX)) {
indexValues.add(storageData.id());
} else if (storageMap.containsKey(index)) {
// avoid indexValue be "null" when inserting
if (storageMap.get(index) == null) {
indexValues.add("");
} else {
indexValues.add(String.valueOf(storageMap.get(index)));
}
storageMap.remove(index);
}
});
// time_bucket has changed to time before calling this method, so remove it from measurements
storageMap.remove(IoTDBClient.TIME_BUCKET);
// processing value to make it suitable for storage
Iterator<Map.Entry<String, Object>> entryIterator = storageMap.entrySet().iterator();
while (entryIterator.hasNext()) {
Map.Entry<String, Object> entry = entryIterator.next();
// IoTDB doesn't allow insert null value.
if (entry.getValue() == null) {
entryIterator.remove();
}
if (entry.getValue() instanceof StorageDataComplexObject) {
storageMap.put(entry.getKey(), ((StorageDataComplexObject) entry.getValue()).toStorageData());
}
}
measurements = new ArrayList<>(storageMap.keySet());
Map<String, TSDataType> columnAndTypeMap = IoTDBTableMetaInfo.get(modelName).getColumnAndTypeMap();
measurementTypes = new ArrayList<>(measurements.size());
for (String measurement : measurements) {
measurementTypes.add(columnAndTypeMap.get(measurement));
}
measurementValues = new ArrayList<>(storageMap.values());
// IoTDB doesn't allow a measurement named `timestamp` or contains `.`
for (String key : storageMap.keySet()) {
if (key.equals(IoTDBClient.TIMESTAMP) || key.contains(".")) {
int idx = measurements.indexOf(key);
measurements.set(idx, "\"" + key + "\"");
}
}
return toStorage.obtain();
}
}

View File

@ -36,7 +36,8 @@ public class IoTDBManagementDAO implements IManagementDAO {
@Override
public void insert(Model model, ManagementData storageData) throws IOException {
IoTDBInsertRequest request = new IoTDBInsertRequest(model.getName(), 1L, storageData, storageBuilder);
IoTDBInsertRequest request =
IoTDBInsertRequest.buildRequest(model.getName(), 1L, storageData, storageBuilder);
client.write(request);
}
}

View File

@ -32,6 +32,7 @@ import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
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.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -45,8 +46,8 @@ public class IoTDBMetricsDAO implements IMetricsDAO {
query.append("select * from ");
for (Metrics metric : metrics) {
query.append(", ");
query = client.addModelPath(query, model.getName());
query.append(IoTDBClient.DOT).append(client.indexValue2LayerName(metric.id()));
IoTDBUtils.addModelPath(client.getStorageGroup(), query, model.getName());
query.append(IoTDBClient.DOT).append(IoTDBUtils.indexValue2LayerName(metric.id()));
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
String queryString = query.toString().replaceFirst(", ", "");
@ -59,7 +60,7 @@ public class IoTDBMetricsDAO implements IMetricsDAO {
@Override
public InsertRequest prepareBatchInsert(Model model, Metrics metrics) {
final long timestamp = TimeBucket.getTimestamp(metrics.getTimeBucket(), model.getDownsampling());
return new IoTDBInsertRequest(model.getName(), timestamp, metrics, storageBuilder);
return IoTDBInsertRequest.buildRequest(model.getName(), timestamp, metrics, storageBuilder);
}
@Override

View File

@ -35,7 +35,8 @@ public class IoTDBNoneStreamDAO implements INoneStreamDAO {
@Override
public void insert(Model model, NoneStream noneStream) throws IOException {
final long timestamp = TimeBucket.getTimestamp(noneStream.getTimeBucket(), model.getDownsampling());
final IoTDBInsertRequest request = new IoTDBInsertRequest(model.getName(), timestamp, noneStream, storageBuilder);
final IoTDBInsertRequest request =
IoTDBInsertRequest.buildRequest(model.getName(), timestamp, noneStream, storageBuilder);
client.write(request);
}
}

View File

@ -40,7 +40,8 @@ public class IoTDBRecordDAO implements IRecordDAO {
@Override
public InsertRequest prepareBatchInsert(Model model, Record record) {
final long timestamp = TimeBucket.getTimestamp(record.getTimeBucket(), model.getDownsampling());
IoTDBInsertRequest request = new IoTDBInsertRequest(model.getName(), timestamp, record, storageBuilder);
IoTDBInsertRequest request =
IoTDBInsertRequest.buildRequest(model.getName(), timestamp, record, storageBuilder);
// transform tags of SegmentRecord, LogRecord, AlarmRecord to tag1, tag2, ...
List<String> measurements = request.getMeasurements();

View File

@ -27,6 +27,7 @@ import org.apache.skywalking.oap.server.core.analysis.manual.networkalias.Networ
import org.apache.skywalking.oap.server.core.storage.StorageData;
import org.apache.skywalking.oap.server.core.storage.cache.INetworkAddressAliasDAO;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -38,14 +39,14 @@ public class IoTDBNetworkAddressAliasDAO implements INetworkAddressAliasDAO {
public List<NetworkAddressAlias> loadLastUpdate(long timeBucket) {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, NetworkAddressAlias.INDEX_NAME);
query = client.addQueryAsterisk(NetworkAddressAlias.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, NetworkAddressAlias.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(NetworkAddressAlias.INDEX_NAME, query);
query.append(" where ").append(NetworkAddressAlias.LAST_UPDATE_TIME_BUCKET).append(" >= ").append(timeBucket)
.append(IoTDBClient.ALIGN_BY_DEVICE);
.append(IoTDBClient.ALIGN_BY_DEVICE);
try {
List<? super StorageData> storageDataList = client.filterQuery(NetworkAddressAlias.INDEX_NAME,
query.toString(), storageBuilder);
query.toString(), storageBuilder);
List<NetworkAddressAlias> networkAddressAliases = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> networkAddressAliases.add((NetworkAddressAlias) storageData));
return networkAddressAliases;

View File

@ -37,6 +37,7 @@ import org.apache.skywalking.oap.server.library.util.StringUtil;
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.base.IoTDBInsertRequest;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -52,14 +53,14 @@ public class IoTDBUITemplateManagementDAO implements UITemplateManagementDAO {
}
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, UITemplate.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, UITemplate.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ID_IDX, id);
query = client.addQueryIndexValue(UITemplate.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(UITemplate.INDEX_NAME, query, indexAndValueMap);
query.append(" limit 1").append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(
UITemplate.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList =
client.filterQuery(UITemplate.INDEX_NAME, query.toString(), storageBuilder);
if (storageDataList.size() > 0) {
return new DashboardConfiguration().fromEntity((UITemplate) storageDataList.get(0));
}
@ -70,20 +71,19 @@ public class IoTDBUITemplateManagementDAO implements UITemplateManagementDAO {
public List<DashboardConfiguration> getAllTemplates(Boolean includingDisabled) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, UITemplate.INDEX_NAME);
query = client.addQueryAsterisk(UITemplate.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, UITemplate.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(UITemplate.INDEX_NAME, query);
if (!includingDisabled) {
query.append(" where ").append(UITemplate.DISABLED).append(" = ").append(BooleanUtils.FALSE);
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(UITemplate.INDEX_NAME, query.toString(),
storageBuilder
);
storageBuilder);
List<DashboardConfiguration> dashboardConfigurationList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData ->
dashboardConfigurationList.add(
new DashboardConfiguration().fromEntity((UITemplate) storageData)));
storageDataList.forEach(
storageData -> dashboardConfigurationList.add(
new DashboardConfiguration().fromEntity((UITemplate) storageData)));
return dashboardConfigurationList;
}
@ -91,9 +91,9 @@ public class IoTDBUITemplateManagementDAO implements UITemplateManagementDAO {
public TemplateChangeStatus addTemplate(DashboardSetting setting) throws IOException {
final UITemplate uiTemplate = setting.toEntity();
IoTDBInsertRequest request = new IoTDBInsertRequest(UITemplate.INDEX_NAME, UI_TEMPLATE_TIMESTAMP,
uiTemplate, storageBuilder
);
IoTDBInsertRequest request =
IoTDBInsertRequest.buildRequest(UITemplate.INDEX_NAME, UI_TEMPLATE_TIMESTAMP,
uiTemplate, storageBuilder);
client.write(request);
return TemplateChangeStatus.builder().status(true).id(setting.getId()).build();
}
@ -104,11 +104,11 @@ public class IoTDBUITemplateManagementDAO implements UITemplateManagementDAO {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, UITemplate.INDEX_NAME);
query.append(IoTDBClient.DOT).append(client.indexValue2LayerName(uiTemplate.id()))
IoTDBUtils.addModelPath(client.getStorageGroup(), query, UITemplate.INDEX_NAME);
query.append(IoTDBClient.DOT).append(IoTDBUtils.indexValue2LayerName(uiTemplate.id()))
.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> queryResult = client.filterQuery(
UITemplate.INDEX_NAME, query.toString(), storageBuilder);
UITemplate.INDEX_NAME, query.toString(), storageBuilder);
if (queryResult.size() == 0) {
return TemplateChangeStatus.builder()
.status(false)
@ -116,9 +116,9 @@ public class IoTDBUITemplateManagementDAO implements UITemplateManagementDAO {
.message("Can't find the template")
.build();
} else {
IoTDBInsertRequest request = new IoTDBInsertRequest(UITemplate.INDEX_NAME, UI_TEMPLATE_TIMESTAMP,
uiTemplate, storageBuilder
);
IoTDBInsertRequest request =
IoTDBInsertRequest.buildRequest(UITemplate.INDEX_NAME, UI_TEMPLATE_TIMESTAMP,
uiTemplate, storageBuilder);
client.write(request);
return TemplateChangeStatus.builder().status(true).id(setting.getId()).build();
}
@ -128,20 +128,20 @@ public class IoTDBUITemplateManagementDAO implements UITemplateManagementDAO {
public TemplateChangeStatus disableTemplate(String id) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, UITemplate.INDEX_NAME);
query.append(IoTDBClient.DOT).append(client.indexValue2LayerName(id))
IoTDBUtils.addModelPath(client.getStorageGroup(), query, UITemplate.INDEX_NAME);
query.append(IoTDBClient.DOT).append(IoTDBUtils.indexValue2LayerName(id))
.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> queryResult = client.filterQuery(
UITemplate.INDEX_NAME, query.toString(), storageBuilder);
UITemplate.INDEX_NAME, query.toString(), storageBuilder);
if (queryResult.size() == 0) {
return TemplateChangeStatus.builder().status(false).id(id).message("Can't find the template").build();
} else {
final UITemplate uiTemplate = (UITemplate) queryResult.get(0);
uiTemplate.setDisabled(BooleanUtils.TRUE);
IoTDBInsertRequest request = new IoTDBInsertRequest(UITemplate.INDEX_NAME, UI_TEMPLATE_TIMESTAMP,
uiTemplate, storageBuilder
);
IoTDBInsertRequest request =
IoTDBInsertRequest.buildRequest(UITemplate.INDEX_NAME, UI_TEMPLATE_TIMESTAMP,
uiTemplate, storageBuilder);
client.write(request);
return TemplateChangeStatus.builder().status(true).id(id).build();
}

View File

@ -30,6 +30,7 @@ import org.apache.skywalking.oap.server.core.storage.StorageData;
import org.apache.skywalking.oap.server.core.storage.profiling.trace.IProfileTaskLogQueryDAO;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -42,28 +43,31 @@ public class IoTDBProfileTaskLogQueryDAO implements IProfileTaskLogQueryDAO {
public List<ProfileTaskLog> getTaskLogList() throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProfileTaskLogRecord.INDEX_NAME);
query = client.addQueryAsterisk(ProfileTaskLogRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProfileTaskLogRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(ProfileTaskLogRecord.INDEX_NAME, query);
query.append(" limit ").append(fetchTaskLogMaxSize).append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProfileTaskLogRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ProfileTaskLogRecord.INDEX_NAME,
query.toString(), storageBuilder);
List<ProfileTaskLogRecord> profileTaskLogRecordList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> profileTaskLogRecordList.add((ProfileTaskLogRecord) storageData));
storageDataList.forEach(storageData ->
profileTaskLogRecordList.add((ProfileTaskLogRecord) storageData));
// resort by self, because of the query result order by time.
profileTaskLogRecordList.sort((ProfileTaskLogRecord r1, ProfileTaskLogRecord r2) ->
Long.compare(r2.getOperationTime(), r1.getOperationTime()));
Long.compare(r2.getOperationTime(), r1.getOperationTime()));
List<ProfileTaskLog> profileTaskLogList = new ArrayList<>(profileTaskLogRecordList.size());
profileTaskLogRecordList.forEach(profileTaskLogRecord -> profileTaskLogList.add(parseLog(profileTaskLogRecord)));
profileTaskLogRecordList.forEach(profileTaskLogRecord ->
profileTaskLogList.add(parseLog(profileTaskLogRecord)));
return profileTaskLogList;
}
private ProfileTaskLog parseLog(ProfileTaskLogRecord record) {
return ProfileTaskLog.builder()
.id(record.id())
.taskId(record.getTaskId())
.instanceId(record.getInstanceId())
.operationType(ProfileTaskLogOperationType.parse(record.getOperationType()))
.operationTime(record.getOperationTime())
.build();
.id(record.id())
.taskId(record.getTaskId())
.instanceId(record.getInstanceId())
.operationType(ProfileTaskLogOperationType.parse(record.getOperationType()))
.operationTime(record.getOperationTime())
.build();
}
}

View File

@ -34,6 +34,7 @@ import org.apache.skywalking.oap.server.core.storage.profiling.trace.IProfileTas
import org.apache.skywalking.oap.server.library.util.StringUtil;
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.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -46,22 +47,28 @@ public class IoTDBProfileTaskQueryDAO implements IProfileTaskQueryDAO {
Long endTimeBucket, Integer limit) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProfileTaskRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProfileTaskRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(serviceId)) {
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
}
query = client.addQueryIndexValue(ProfileTaskRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(ProfileTaskRecord.INDEX_NAME, query, indexAndValueMap);
StringBuilder where = new StringBuilder(" where ");
if (StringUtil.isNotEmpty(endpointName)) {
where.append(ProfileTaskRecord.ENDPOINT_NAME).append(" = \"").append(endpointName).append("\"").append(" and ");
where.append(ProfileTaskRecord.ENDPOINT_NAME).append(" = \"")
.append(endpointName).append("\"")
.append(" and ");
}
if (Objects.nonNull(startTimeBucket)) {
where.append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(startTimeBucket)).append(" and ");
where.append(IoTDBClient.TIME).append(" >= ")
.append(TimeBucket.getTimestamp(startTimeBucket))
.append(" and ");
}
if (Objects.nonNull(endTimeBucket)) {
where.append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTimeBucket)).append(" and ");
where.append(IoTDBClient.TIME).append(" <= ")
.append(TimeBucket.getTimestamp(endTimeBucket))
.append(" and ");
}
if (where.length() > 7) {
int length = where.length();
@ -73,9 +80,11 @@ public class IoTDBProfileTaskQueryDAO implements IProfileTaskQueryDAO {
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProfileTaskRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ProfileTaskRecord.INDEX_NAME,
query.toString(), storageBuilder);
List<ProfileTask> profileTaskList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> profileTaskList.add(record2ProfileTask((ProfileTaskRecord) storageData)));
storageDataList.forEach(storageData ->
profileTaskList.add(record2ProfileTask((ProfileTaskRecord) storageData)));
return profileTaskList;
}
@ -86,27 +95,28 @@ public class IoTDBProfileTaskQueryDAO implements IProfileTaskQueryDAO {
}
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProfileTaskRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProfileTaskRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ID_IDX, id);
query = client.addQueryIndexValue(ProfileTaskRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(ProfileTaskRecord.INDEX_NAME, query, indexAndValueMap);
query.append(" limit 1").append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProfileTaskRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ProfileTaskRecord.INDEX_NAME,
query.toString(), storageBuilder);
return record2ProfileTask((ProfileTaskRecord) storageDataList.get(0));
}
private static ProfileTask record2ProfileTask(ProfileTaskRecord record) {
return ProfileTask.builder()
.id(record.id())
.serviceId(record.getServiceId())
.endpointName(record.getEndpointName())
.startTime(record.getStartTime())
.createTime(record.getCreateTime())
.duration(record.getDuration())
.minDurationThreshold(record.getMinDurationThreshold())
.dumpPeriod(record.getDumpPeriod())
.maxSamplingCount(record.getMaxSamplingCount())
.build();
.id(record.id())
.serviceId(record.getServiceId())
.endpointName(record.getEndpointName())
.startTime(record.getStartTime())
.createTime(record.getCreateTime())
.duration(record.getDuration())
.minDurationThreshold(record.getMinDurationThreshold())
.dumpPeriod(record.getDumpPeriod())
.maxSamplingCount(record.getMaxSamplingCount())
.build();
}
}

View File

@ -32,28 +32,32 @@ import org.apache.skywalking.oap.server.core.storage.profiling.trace.IProfileThr
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.library.util.BooleanUtils;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@RequiredArgsConstructor
public class IoTDBProfileThreadSnapshotQueryDAO implements IProfileThreadSnapshotQueryDAO {
private final IoTDBClient client;
private final StorageBuilder<ProfileThreadSnapshotRecord> profileThreadSnapshotRecordBuilder = new ProfileThreadSnapshotRecord.Builder();
private final StorageBuilder<ProfileThreadSnapshotRecord> profileThreadSnapshotRecordBuilder =
new ProfileThreadSnapshotRecord.Builder();
private final StorageBuilder<SegmentRecord> segmentRecordBuilder = new SegmentRecord.Builder();
@Override
public List<BasicTrace> queryProfiledSegments(String taskId) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProfileThreadSnapshotRecord.INDEX_NAME);
query = client.addQueryAsterisk(ProfileThreadSnapshotRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProfileThreadSnapshotRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(ProfileThreadSnapshotRecord.INDEX_NAME, query);
query.append(" where ").append(ProfileThreadSnapshotRecord.TASK_ID).append(" = \"").append(taskId).append("\"")
.append(" and ").append(ProfileThreadSnapshotRecord.SEQUENCE).append(" = 0")
.append(IoTDBClient.ALIGN_BY_DEVICE);
.append(" and ").append(ProfileThreadSnapshotRecord.SEQUENCE).append(" = 0")
.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProfileThreadSnapshotRecord.INDEX_NAME,
query.toString(), profileThreadSnapshotRecordBuilder);
query.toString(),
profileThreadSnapshotRecordBuilder);
// We can insure the size of List, so use ArrayList to improve visit speed. (Other storage plugin use LinkedList)
final List<String> segmentIds = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> segmentIds.add(((ProfileThreadSnapshotRecord) storageData).getSegmentId()));
storageDataList.forEach(
storageData -> segmentIds.add(((ProfileThreadSnapshotRecord) storageData).getSegmentId()));
if (segmentIds.isEmpty()) {
return Collections.emptyList();
}
@ -62,8 +66,8 @@ public class IoTDBProfileThreadSnapshotQueryDAO implements IProfileThreadSnapsho
// https://github.com/apache/iotdb/discussions/3888
query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, SegmentRecord.INDEX_NAME);
query = client.addQueryAsterisk(SegmentRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, SegmentRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(SegmentRecord.INDEX_NAME, query);
query.append(" where ").append(SegmentRecord.SEGMENT_ID).append(" in (");
for (String segmentId : segmentIds) {
query.append("\"").append(segmentId).append("\"").append(", ");
@ -74,14 +78,16 @@ public class IoTDBProfileThreadSnapshotQueryDAO implements IProfileThreadSnapsho
List<SegmentRecord> segmentRecordList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> segmentRecordList.add((SegmentRecord) storageData));
// resort by self, because of the select query result order by time.
segmentRecordList.sort((SegmentRecord r1, SegmentRecord r2) -> Long.compare(r2.getStartTime(), r1.getStartTime()));
segmentRecordList.sort((SegmentRecord r1, SegmentRecord r2) ->
Long.compare(r2.getStartTime(), r1.getStartTime()));
List<BasicTrace> result = new ArrayList<>(segmentRecordList.size());
segmentRecordList.forEach(segmentRecord -> {
BasicTrace basicTrace = new BasicTrace();
basicTrace.setSegmentId(segmentRecord.getSegmentId());
basicTrace.setStart(String.valueOf(segmentRecord.getStartTime()));
basicTrace.getEndpointNames().add(IDManager.EndpointID.analysisId(segmentRecord.getEndpointId()).getEndpointName());
basicTrace.getEndpointNames().add(IDManager.EndpointID.analysisId(segmentRecord.getEndpointId())
.getEndpointName());
basicTrace.setDuration(segmentRecord.getLatency());
basicTrace.setError(BooleanUtils.valueToBoolean(segmentRecord.getIsError()));
basicTrace.getTraceIds().add(segmentRecord.getTraceId());
@ -101,20 +107,23 @@ public class IoTDBProfileThreadSnapshotQueryDAO implements IProfileThreadSnapsho
}
@Override
public List<ProfileThreadSnapshotRecord> queryRecords(String segmentId, int minSequence, int maxSequence) throws IOException {
public List<ProfileThreadSnapshotRecord> queryRecords(String segmentId, int minSequence, int maxSequence)
throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProfileThreadSnapshotRecord.INDEX_NAME);
query = client.addQueryAsterisk(ProfileThreadSnapshotRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProfileThreadSnapshotRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(ProfileThreadSnapshotRecord.INDEX_NAME, query);
query.append(" where ").append(ProfileThreadSnapshotRecord.SEGMENT_ID).append(" = \"").append(segmentId).append("\"")
.append(" and ").append(ProfileThreadSnapshotRecord.SEQUENCE).append(" >= ").append(minSequence)
.append(" and ").append(ProfileThreadSnapshotRecord.SEQUENCE).append(" <= ").append(maxSequence)
.append(IoTDBClient.ALIGN_BY_DEVICE);
.append(" and ").append(ProfileThreadSnapshotRecord.SEQUENCE).append(" >= ").append(minSequence)
.append(" and ").append(ProfileThreadSnapshotRecord.SEQUENCE).append(" <= ").append(maxSequence)
.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProfileThreadSnapshotRecord.INDEX_NAME,
query.toString(), profileThreadSnapshotRecordBuilder);
query.toString(),
profileThreadSnapshotRecordBuilder);
List<ProfileThreadSnapshotRecord> profileThreadSnapshotRecordList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> profileThreadSnapshotRecordList.add((ProfileThreadSnapshotRecord) storageData));
storageDataList.forEach(
storageData -> profileThreadSnapshotRecordList.add((ProfileThreadSnapshotRecord) storageData));
return profileThreadSnapshotRecordList;
}
@ -122,13 +131,13 @@ public class IoTDBProfileThreadSnapshotQueryDAO implements IProfileThreadSnapsho
public SegmentRecord getProfiledSegment(String segmentId) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, SegmentRecord.INDEX_NAME);
query = client.addQueryAsterisk(SegmentRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, SegmentRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(SegmentRecord.INDEX_NAME, query);
query.append(" where ").append(SegmentRecord.SEGMENT_ID).append(" = \"").append(segmentId).append("\"")
.append(IoTDBClient.ALIGN_BY_DEVICE);
.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(SegmentRecord.INDEX_NAME,
query.toString(), segmentRecordBuilder);
query.toString(), segmentRecordBuilder);
if (storageDataList.isEmpty()) {
return null;
}
@ -140,15 +149,15 @@ public class IoTDBProfileThreadSnapshotQueryDAO implements IProfileThreadSnapsho
// See https://github.com/apache/iotdb/discussions/3907
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProfileThreadSnapshotRecord.INDEX_NAME);
query = client.addQueryAsterisk(ProfileThreadSnapshotRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProfileThreadSnapshotRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(ProfileThreadSnapshotRecord.INDEX_NAME, query);
query.append(" where ").append(ProfileThreadSnapshotRecord.SEGMENT_ID).append(" = \"").append(segmentId).append("\"")
.append(" and ").append(ProfileThreadSnapshotRecord.DUMP_TIME).append(" >= ").append(start)
.append(" and ").append(ProfileThreadSnapshotRecord.DUMP_TIME).append(" <= ").append(end)
.append(IoTDBClient.ALIGN_BY_DEVICE);
.append(" and ").append(ProfileThreadSnapshotRecord.DUMP_TIME).append(" >= ").append(start)
.append(" and ").append(ProfileThreadSnapshotRecord.DUMP_TIME).append(" <= ").append(end)
.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProfileThreadSnapshotRecord.INDEX_NAME,
query.toString(), profileThreadSnapshotRecordBuilder);
query.toString(),
profileThreadSnapshotRecordBuilder);
if (aggType.equals("min_value")) {
int minValue = Integer.MAX_VALUE;
for (Object storageData : storageDataList) {

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.oap.server.storage.plugin.iotdb.query;
import com.google.common.base.Splitter;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Comparator;
@ -40,6 +41,7 @@ import org.apache.skywalking.oap.server.core.query.type.SelectedRecord;
import org.apache.skywalking.oap.server.core.storage.query.IAggregationQueryDAO;
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.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -49,11 +51,12 @@ public class IoTDBAggregationQueryDAO implements IAggregationQueryDAO {
@Override
public List<SelectedRecord> sortMetrics(TopNCondition condition, String valueColumnName, Duration duration,
List<KeyValue> additionalConditions) throws IOException {
// This method maybe have poor efficiency. It queries all data which meets a condition without aggregation function.
// This method maybe have poor efficiency.
// It queries all data which meets a condition without aggregation function.
// https://github.com/apache/iotdb/issues/4006
StringBuilder query = new StringBuilder();
query.append(String.format("select %s from ", valueColumnName));
query = client.addModelPath(query, condition.getName());
IoTDBUtils.addModelPath(client.getStorageGroup(), query, condition.getName());
Map<String, String> indexAndValueMap = new HashMap<>();
List<KeyValue> measurementConditions = new ArrayList<>();
@ -68,17 +71,17 @@ public class IoTDBAggregationQueryDAO implements IAggregationQueryDAO {
}
}
if (!indexAndValueMap.isEmpty()) {
query = client.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
} else {
query = client.addQueryAsterisk(condition.getName(), query);
IoTDBUtils.addQueryAsterisk(condition.getName(), query);
}
query.append(" where ").append(IoTDBClient.TIME).append(" >= ").append(duration.getStartTimestamp())
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(duration.getEndTimestamp());
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(duration.getEndTimestamp());
if (!measurementConditions.isEmpty()) {
for (KeyValue measurementCondition : measurementConditions) {
query.append(" and ").append(measurementCondition.getKey()).append(" = \"")
.append(measurementCondition.getValue()).append("\"");
.append(measurementCondition.getValue()).append("\"");
}
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
@ -97,8 +100,8 @@ public class IoTDBAggregationQueryDAO implements IAggregationQueryDAO {
while (wrapper.hasNext()) {
RowRecord rowRecord = wrapper.next();
List<Field> fields = rowRecord.getFields();
String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
String entityId = client.layerName2IndexValue(layerNames[2]);
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"").splitToList(fields.get(0).getStringValue());
String entityId = IoTDBUtils.layerName2IndexValue(layerNames.get(2));
double value = Double.parseDouble(fields.get(1).getStringValue());
entityIdAndSumMap.merge(entityId, value, Double::sum);
entityIdAndCountMap.merge(entityId, 1, Integer::sum);
@ -122,7 +125,8 @@ public class IoTDBAggregationQueryDAO implements IAggregationQueryDAO {
if (condition.getOrder().equals(Order.DES)) {
topEntities.sort((SelectedRecord t1, SelectedRecord t2) ->
Double.compare(Double.parseDouble(t2.getValue()), Double.parseDouble(t1.getValue())));
Double.compare(Double.parseDouble(t2.getValue()),
Double.parseDouble(t1.getValue())));
} else {
topEntities.sort(Comparator.comparingDouble((SelectedRecord t) -> Double.parseDouble(t.getValue())));
}

View File

@ -34,6 +34,7 @@ import org.apache.skywalking.oap.server.core.storage.query.IAlarmQueryDAO;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@RequiredArgsConstructor
public class IoTDBAlarmQueryDAO implements IAlarmQueryDAO {
@ -41,13 +42,15 @@ public class IoTDBAlarmQueryDAO implements IAlarmQueryDAO {
private final StorageBuilder<AlarmRecord> storageBuilder = new AlarmRecord.Builder();
@Override
public Alarms getAlarm(Integer scopeId, String keyword, int limit, int from, long startTB, long endTB, List<Tag> tags) throws IOException {
public Alarms getAlarm(Integer scopeId, String keyword, int limit, int from,
long startTB, long endTB, List<Tag> tags)
throws IOException {
StringBuilder query = new StringBuilder();
// This method maybe have poor efficiency. It queries all data which meets a condition without select function.
// https://github.com/apache/iotdb/discussions/3888
query.append("select * from ");
query = client.addModelPath(query, AlarmRecord.INDEX_NAME);
query = client.addQueryAsterisk(AlarmRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, AlarmRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(AlarmRecord.INDEX_NAME, query);
StringBuilder where = new StringBuilder(" where ");
if (Objects.nonNull(scopeId)) {
@ -73,7 +76,8 @@ public class IoTDBAlarmQueryDAO implements IAlarmQueryDAO {
query.append(IoTDBClient.ALIGN_BY_DEVICE);
Alarms alarms = new Alarms();
List<? super StorageData> storageDataList = client.filterQuery(AlarmRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(AlarmRecord.INDEX_NAME,
query.toString(), storageBuilder);
int limitCount = 0;
for (int i = from; i < storageDataList.size(); i++) {
if (limitCount < limit) {
@ -84,7 +88,8 @@ public class IoTDBAlarmQueryDAO implements IAlarmQueryDAO {
}
alarms.setTotal(storageDataList.size());
// resort by self, because of the select query result order by time.
alarms.getMsgs().sort((AlarmMessage m1, AlarmMessage m2) -> Long.compare(m2.getStartTime(), m1.getStartTime()));
alarms.getMsgs().sort((AlarmMessage m1, AlarmMessage m2) ->
Long.compare(m2.getStartTime(), m1.getStartTime()));
return alarms;
}

View File

@ -40,6 +40,7 @@ import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import org.apache.skywalking.oap.server.library.util.StringUtil;
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.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -50,31 +51,42 @@ public class IoTDBBrowserLogQueryDAO implements IBrowserLogQueryDAO {
@Override
public BrowserErrorLogs queryBrowserErrorLogs(String serviceId, String serviceVersionId, String pagePathId,
BrowserErrorCategory category, long startSecondTB,
long endSecondTB, int limit, int from) throws IOException {
long endSecondTB, int limit, int from)
throws IOException {
StringBuilder query = new StringBuilder();
// This method maybe have poor efficiency. It queries all data which meets a condition without select function.
// https://github.com/apache/iotdb/discussions/3888
query.append("select * from ");
query = client.addModelPath(query, BrowserErrorLogRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, BrowserErrorLogRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(serviceId)) {
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
}
query = client.addQueryIndexValue(BrowserErrorLogRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(BrowserErrorLogRecord.INDEX_NAME, query, indexAndValueMap);
StringBuilder where = new StringBuilder(" where ");
if (startSecondTB != 0 && endSecondTB != 0) {
where.append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(startSecondTB)).append(" and ");
where.append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endSecondTB)).append(" and ");
where.append(IoTDBClient.TIME).append(" >= ")
.append(TimeBucket.getTimestamp(startSecondTB))
.append(" and ");
where.append(IoTDBClient.TIME).append(" <= ")
.append(TimeBucket.getTimestamp(endSecondTB))
.append(" and ");
}
if (StringUtil.isNotEmpty(serviceVersionId)) {
where.append(BrowserErrorLogRecord.SERVICE_VERSION_ID).append(" = \"").append(serviceVersionId).append("\"").append(" and ");
where.append(BrowserErrorLogRecord.SERVICE_VERSION_ID)
.append(" = \"").append(serviceVersionId).append("\"")
.append(" and ");
}
if (StringUtil.isNotEmpty(pagePathId)) {
where.append(BrowserErrorLogRecord.PAGE_PATH_ID).append(" = \"").append(pagePathId).append("\"").append(" and ");
where.append(BrowserErrorLogRecord.PAGE_PATH_ID)
.append(" = \"").append(pagePathId).append("\"")
.append(" and ");
}
if (Objects.nonNull(category)) {
where.append(BrowserErrorLogRecord.ERROR_CATEGORY).append(" = ").append(category.getValue()).append(" and ");
where.append(BrowserErrorLogRecord.ERROR_CATEGORY)
.append(" = ").append(category.getValue())
.append(" and ");
}
if (where.length() > 7) {
int length = where.length();
@ -83,11 +95,13 @@ public class IoTDBBrowserLogQueryDAO implements IBrowserLogQueryDAO {
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(BrowserErrorLogRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(BrowserErrorLogRecord.INDEX_NAME,
query.toString(), storageBuilder);
List<BrowserErrorLogRecord> browserErrorLogRecordList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> browserErrorLogRecordList.add((BrowserErrorLogRecord) storageData));
// resort by self, because of the select query result order by time.
browserErrorLogRecordList.sort((BrowserErrorLogRecord b1, BrowserErrorLogRecord b2) -> Long.compare(b2.getTimestamp(), b1.getTimestamp()));
browserErrorLogRecordList.sort((BrowserErrorLogRecord b1, BrowserErrorLogRecord b2) ->
Long.compare(b2.getTimestamp(), b1.getTimestamp()));
BrowserErrorLogs logs = new BrowserErrorLogs();
int limitCount = 0;
for (int i = from; i < browserErrorLogRecordList.size(); i++) {

View File

@ -31,6 +31,7 @@ import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -39,12 +40,13 @@ public class IoTDBEBPFProfilingDataDAO implements IEBPFProfilingDataDAO {
private final StorageBuilder<EBPFProfilingDataRecord> storageBuilder = new EBPFProfilingDataRecord.Builder();
@Override
public List<EBPFProfilingDataRecord> queryData(String taskId, long beginTime, long endTime) throws IOException {
public List<EBPFProfilingDataRecord> queryData(String taskId, long beginTime, long endTime)
throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, EBPFProfilingDataRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, EBPFProfilingDataRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
query = client.addQueryIndexValue(EBPFProfilingDataRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(EBPFProfilingDataRecord.INDEX_NAME, query, indexAndValueMap);
StringBuilder where = new StringBuilder(" where ");
where.append(EBPFProfilingDataRecord.TASK_ID).append(" = \"").append(taskId).append("\" and ");
@ -57,7 +59,8 @@ public class IoTDBEBPFProfilingDataDAO implements IEBPFProfilingDataDAO {
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(EBPFProfilingDataRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(EBPFProfilingDataRecord.INDEX_NAME,
query.toString(), storageBuilder);
List<EBPFProfilingDataRecord> dataList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> dataList.add((EBPFProfilingDataRecord) storageData));
return dataList;

View File

@ -18,6 +18,9 @@
package org.apache.skywalking.oap.server.storage.plugin.iotdb.query;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
@ -27,10 +30,7 @@ import org.apache.skywalking.oap.server.core.storage.StorageData;
import org.apache.skywalking.oap.server.core.storage.profiling.ebpf.IEBPFProfilingScheduleDAO;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -39,11 +39,12 @@ public class IoTDBEBPFProfilingScheduleDAO implements IEBPFProfilingScheduleDAO
private final StorageBuilder<EBPFProfilingScheduleRecord> storageBuilder = new EBPFProfilingScheduleRecord.Builder();
@Override
public List<EBPFProfilingSchedule> querySchedules(String taskId, long startTimeBucket, long endTimeBucket) throws IOException {
public List<EBPFProfilingSchedule> querySchedules(String taskId, long startTimeBucket, long endTimeBucket)
throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, EBPFProfilingScheduleRecord.INDEX_NAME);
query = client.addQueryAsterisk(EBPFProfilingScheduleRecord.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, EBPFProfilingScheduleRecord.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(EBPFProfilingScheduleRecord.INDEX_NAME, query);
StringBuilder where = new StringBuilder(" where ");
where.append(EBPFProfilingScheduleRecord.TASK_ID).append(" = \"").append(taskId).append("\" and ");
@ -56,7 +57,8 @@ public class IoTDBEBPFProfilingScheduleDAO implements IEBPFProfilingScheduleDAO
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(EBPFProfilingScheduleRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(EBPFProfilingScheduleRecord.INDEX_NAME,
query.toString(), storageBuilder);
List<EBPFProfilingSchedule> scheduleList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> scheduleList.add(parseSchedule((EBPFProfilingScheduleRecord) storageData)));
return scheduleList;

View File

@ -40,6 +40,7 @@ import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -50,10 +51,11 @@ public class IoTDBEBPFProfilingTaskDAO implements IEBPFProfilingTaskDAO {
@Override
public List<EBPFProfilingTask> queryTasks(EBPFProfilingProcessFinder finder,
EBPFProfilingTargetType targetType,
long taskStartTime, long latestUpdateTime) throws IOException {
long taskStartTime, long latestUpdateTime)
throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, EBPFProfilingTaskRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, EBPFProfilingTaskRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(finder.getServiceId())) {
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, finder.getServiceId());
@ -69,18 +71,22 @@ public class IoTDBEBPFProfilingTaskDAO implements IEBPFProfilingTaskDAO {
}
final HashMap<String, String> indexWithProcessId = new HashMap<>(indexAndValueMap);
indexWithProcessId.put(IoTDBIndexes.PROCESS_ID_INX, processIdList.get(i));
query = client.addQueryIndexValue(EBPFProfilingTaskRecord.INDEX_NAME, query, indexWithProcessId);
IoTDBUtils.addQueryIndexValue(EBPFProfilingTaskRecord.INDEX_NAME, query, indexWithProcessId);
}
} else {
query = client.addQueryIndexValue(EBPFProfilingTaskRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(EBPFProfilingTaskRecord.INDEX_NAME, query, indexAndValueMap);
}
StringBuilder where = new StringBuilder(" where ");
if (taskStartTime > 0) {
where.append(EBPFProfilingTaskRecord.START_TIME).append(" >= ").append(taskStartTime).append(" and ");
where.append(EBPFProfilingTaskRecord.START_TIME)
.append(" >= ").append(taskStartTime)
.append(" and ");
}
if (latestUpdateTime > 0) {
where.append(EBPFProfilingTaskRecord.LAST_UPDATE_TIME).append(" > ").append(latestUpdateTime).append(" and ");
where.append(EBPFProfilingTaskRecord.LAST_UPDATE_TIME)
.append(" > ").append(latestUpdateTime)
.append(" and ");
}
if (where.length() > 7) {
int length = where.length();
@ -89,7 +95,8 @@ public class IoTDBEBPFProfilingTaskDAO implements IEBPFProfilingTaskDAO {
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(EBPFProfilingTaskRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(EBPFProfilingTaskRecord.INDEX_NAME,
query.toString(), storageBuilder);
List<EBPFProfilingTask> taskList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> taskList.add(parseTask((EBPFProfilingTaskRecord) storageData)));
return taskList;

View File

@ -36,6 +36,7 @@ import org.apache.skywalking.oap.server.core.storage.StorageData;
import org.apache.skywalking.oap.server.core.storage.query.IEventQueryDAO;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -49,15 +50,16 @@ public class IoTDBEventQueryDAO implements IEventQueryDAO {
// https://github.com/apache/iotdb/discussions/3888
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, Event.INDEX_NAME);
query = client.addQueryAsterisk(Event.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, Event.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(Event.INDEX_NAME, query);
StringBuilder where = whereSQL(condition);
if (where.length() > 0) {
query.append(" where ").append(where);
}
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(Event.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(Event.INDEX_NAME,
query.toString(), storageBuilder);
final Events events = new Events();
int limitCount = 0;
PaginationUtils.Page page = PaginationUtils.INSTANCE.exchange(condition.getPaging());
@ -89,8 +91,8 @@ public class IoTDBEventQueryDAO implements IEventQueryDAO {
// https://github.com/apache/iotdb/discussions/3888
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, Event.INDEX_NAME);
query = client.addQueryAsterisk(Event.INDEX_NAME, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, Event.INDEX_NAME);
IoTDBUtils.addQueryAsterisk(Event.INDEX_NAME, query);
StringBuilder where = whereSQL(conditionList);
if (where.length() > 0) {
query.append(" where ").append(where);
@ -132,28 +134,42 @@ public class IoTDBEventQueryDAO implements IEventQueryDAO {
final Source source = condition.getSource();
if (source != null) {
if (!Strings.isNullOrEmpty(source.getService())) {
where.append(Event.SERVICE).append(" = \"").append(source.getService()).append("\"").append(" and ");
where.append(Event.SERVICE).append(" = \"")
.append(source.getService()).append("\"")
.append(" and ");
}
if (!Strings.isNullOrEmpty(source.getServiceInstance())) {
where.append(Event.SERVICE_INSTANCE).append(" = \"").append(source.getServiceInstance()).append("\"").append(" and ");
where.append(Event.SERVICE_INSTANCE).append(" = \"")
.append(source.getServiceInstance()).append("\"")
.append(" and ");
}
if (!Strings.isNullOrEmpty(source.getEndpoint())) {
where.append(Event.ENDPOINT).append(" = \"").append(source.getEndpoint()).append("\"").append(" and ");
where.append(Event.ENDPOINT).append(" = \"")
.append(source.getEndpoint()).append("\"")
.append(" and ");
}
}
if (!Strings.isNullOrEmpty(condition.getName())) {
where.append(Event.NAME).append(" = \"").append(condition.getName()).append("\"").append(" and ");
where.append(Event.NAME).append(" = \"")
.append(condition.getName()).append("\"")
.append(" and ");
}
if (condition.getType() != null) {
where.append(Event.TYPE).append(" = \"").append(condition.getType().name()).append("\"").append(" and ");
where.append(Event.TYPE).append(" = \"")
.append(condition.getType().name()).append("\"")
.append(" and ");
}
final Duration time = condition.getTime();
if (time != null) {
if (time.getStartTimestamp() > 0) {
where.append(Event.START_TIME).append(" > ").append(time.getStartTimestamp()).append(" and ");
where.append(Event.START_TIME).append(" > ")
.append(time.getStartTimestamp())
.append(" and ");
}
if (time.getEndTimestamp() > 0) {
where.append(Event.END_TIME).append(" < ").append(time.getEndTimestamp()).append(" and ");
where.append(Event.END_TIME).append(" < ")
.append(time.getEndTimestamp())
.append(" and ");
}
}
if (where.length() > 0) {
@ -183,7 +199,8 @@ public class IoTDBEventQueryDAO implements IEventQueryDAO {
}
private org.apache.skywalking.oap.server.core.query.type.event.Event parseEvent(final Event event) {
final org.apache.skywalking.oap.server.core.query.type.event.Event resultEvent = new org.apache.skywalking.oap.server.core.query.type.event.Event();
final org.apache.skywalking.oap.server.core.query.type.event.Event resultEvent =
new org.apache.skywalking.oap.server.core.query.type.event.Event();
resultEvent.setUuid(event.getUuid());
resultEvent.setSource(new Source(event.getService(), event.getServiceInstance(), event.getEndpoint()));
resultEvent.setName(event.getName());

View File

@ -45,6 +45,7 @@ import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import org.apache.skywalking.oap.server.library.util.StringUtil;
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.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -61,7 +62,7 @@ public class IoTDBLogQueryDAO implements ILogQueryDAO {
// This method maybe have poor efficiency. It queries all data which meets a condition without select function.
// https://github.com/apache/iotdb/discussions/3888
query.append("select * from ");
query = client.addModelPath(query, LogRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, LogRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(serviceId)) {
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
@ -69,30 +70,44 @@ public class IoTDBLogQueryDAO implements ILogQueryDAO {
if (Objects.nonNull(relatedTrace) && StringUtil.isNotEmpty(relatedTrace.getTraceId())) {
indexAndValueMap.put(IoTDBIndexes.TRACE_ID_IDX, relatedTrace.getTraceId());
}
query = client.addQueryIndexValue(LogRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(LogRecord.INDEX_NAME, query, indexAndValueMap);
StringBuilder where = new StringBuilder(" where ");
if (startTB != 0 && endTB != 0) {
where.append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(startTB)).append(" and ");
where.append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTB)).append(" and ");
where.append(IoTDBClient.TIME).append(" >= ")
.append(TimeBucket.getTimestamp(startTB))
.append(" and ");
where.append(IoTDBClient.TIME).append(" <= ")
.append(TimeBucket.getTimestamp(endTB))
.append(" and ");
}
if (StringUtil.isNotEmpty(serviceInstanceId)) {
where.append(AbstractLogRecord.SERVICE_INSTANCE_ID).append(" = \"").append(serviceInstanceId).append("\"").append(" and ");
where.append(AbstractLogRecord.SERVICE_INSTANCE_ID)
.append(" = \"").append(serviceInstanceId).append("\"")
.append(" and ");
}
if (StringUtil.isNotEmpty(endpointId)) {
where.append(AbstractLogRecord.ENDPOINT_ID).append(" = \"").append(endpointId).append("\"").append(" and ");
where.append(AbstractLogRecord.ENDPOINT_ID)
.append(" = \"").append(endpointId).append("\"")
.append(" and ");
}
if (Objects.nonNull(relatedTrace)) {
if (StringUtil.isNotEmpty(relatedTrace.getSegmentId())) {
where.append(AbstractLogRecord.TRACE_SEGMENT_ID).append(" = \"").append(relatedTrace.getSegmentId()).append("\"").append(" and ");
where.append(AbstractLogRecord.TRACE_SEGMENT_ID)
.append(" = \"").append(relatedTrace.getSegmentId()).append("\"")
.append(" and ");
}
if (Objects.nonNull(relatedTrace.getSpanId())) {
where.append(AbstractLogRecord.SPAN_ID).append(" = ").append(relatedTrace.getSpanId()).append(" and ");
where.append(AbstractLogRecord.SPAN_ID)
.append(" = ").append(relatedTrace.getSpanId())
.append(" and ");
}
}
if (CollectionUtils.isNotEmpty(tags)) {
for (final Tag tag : tags) {
where.append(tag.getKey()).append(" = \"").append(tag.getValue()).append("\"").append(" and ");
where.append(tag.getKey()).append(" = \"")
.append(tag.getValue()).append("\"")
.append(" and ");
}
}
if (where.length() > 7) {
@ -103,7 +118,8 @@ public class IoTDBLogQueryDAO implements ILogQueryDAO {
query.append(IoTDBClient.ALIGN_BY_DEVICE);
Logs logs = new Logs();
List<? super StorageData> storageDataList = client.filterQuery(LogRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(LogRecord.INDEX_NAME,
query.toString(), storageBuilder);
int limitCount = 0;
for (int i = from; i < storageDataList.size(); i++) {
if (limitCount < limit) {

View File

@ -28,6 +28,7 @@ import java.util.List;
import java.util.Map;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.oap.server.core.Const;
import org.apache.skywalking.oap.server.core.analysis.IDManager;
import org.apache.skywalking.oap.server.core.analysis.Layer;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
@ -48,6 +49,7 @@ import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.library.util.StringUtil;
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.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -62,7 +64,7 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
public List<Service> listServices(final String layer, final String group) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ServiceTraffic.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ServiceTraffic.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(layer)) {
indexAndValueMap.put(IoTDBIndexes.LAYER_IDX, String.valueOf(Layer.valueOf(layer).value()));
@ -70,10 +72,11 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
if (StringUtil.isNotEmpty(group)) {
indexAndValueMap.put(IoTDBIndexes.GROUP_IDX, group);
}
query = client.addQueryIndexValue(ServiceTraffic.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(ServiceTraffic.INDEX_NAME, query, indexAndValueMap);
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ServiceTraffic.INDEX_NAME, query.toString(), serviceBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ServiceTraffic.INDEX_NAME,
query.toString(), serviceBuilder);
return buildServices(storageDataList);
}
@ -81,29 +84,32 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
public List<Service> getServices(final String serviceId) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ServiceTraffic.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ServiceTraffic.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
query = client.addQueryIndexValue(ServiceTraffic.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(ServiceTraffic.INDEX_NAME, query, indexAndValueMap);
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ServiceTraffic.INDEX_NAME, query.toString(), serviceBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ServiceTraffic.INDEX_NAME,
query.toString(), serviceBuilder);
return buildServices(storageDataList);
}
@Override
public List<ServiceInstance> listInstances(long startTimestamp, long endTimestamp, String serviceId) throws IOException {
public List<ServiceInstance> listInstances(long startTimestamp, long endTimestamp, String serviceId)
throws IOException {
final long minuteTimeBucket = TimeBucket.getMinuteTimeBucket(startTimestamp);
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, InstanceTraffic.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, InstanceTraffic.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
query = client.addQueryIndexValue(InstanceTraffic.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(InstanceTraffic.INDEX_NAME, query, indexAndValueMap);
query.append(" where ").append(InstanceTraffic.LAST_PING_TIME_BUCKET).append(" >= ").append(minuteTimeBucket)
.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(InstanceTraffic.INDEX_NAME, query.toString(), instanceBuilder);
List<? super StorageData> storageDataList = client.filterQuery(InstanceTraffic.INDEX_NAME,
query.toString(), instanceBuilder);
return buildInstances(storageDataList);
}
@ -111,13 +117,14 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
public ServiceInstance getInstance(final String instanceId) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ServiceTraffic.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ServiceTraffic.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ID_IDX, instanceId);
query = client.addQueryIndexValue(ServiceTraffic.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(ServiceTraffic.INDEX_NAME, query, indexAndValueMap);
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ServiceTraffic.INDEX_NAME, query.toString(), serviceBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ServiceTraffic.INDEX_NAME,
query.toString(), serviceBuilder);
final List<ServiceInstance> instances = buildInstances(storageDataList);
return instances.size() > 0 ? instances.get(0) : null;
}
@ -126,16 +133,17 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
public List<Endpoint> findEndpoint(String keyword, String serviceId, int limit) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, EndpointTraffic.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, EndpointTraffic.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
query = client.addQueryIndexValue(EndpointTraffic.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(EndpointTraffic.INDEX_NAME, query, indexAndValueMap);
if (!Strings.isNullOrEmpty(keyword)) {
query.append(" where ").append(EndpointTraffic.NAME).append(" like '%").append(keyword).append("%'");
}
query.append(" limit ").append(limit).append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(EndpointTraffic.INDEX_NAME, query.toString(), endpointBuilder);
List<? super StorageData> storageDataList = client.filterQuery(EndpointTraffic.INDEX_NAME,
query.toString(), endpointBuilder);
List<Endpoint> endpointList = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> {
EndpointTraffic endpointTraffic = (EndpointTraffic) storageData;
@ -148,10 +156,11 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
}
@Override
public List<Process> listProcesses(String serviceId, String instanceId, String agentId) throws IOException {
public List<Process> listProcesses(String serviceId, String instanceId, String agentId)
throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProcessTraffic.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProcessTraffic.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(serviceId)) {
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
@ -162,10 +171,11 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
if (StringUtil.isNotEmpty(agentId)) {
indexAndValueMap.put(IoTDBIndexes.AGENT_ID_INX, agentId);
}
query = client.addQueryIndexValue(ProcessTraffic.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(ProcessTraffic.INDEX_NAME, query, indexAndValueMap);
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProcessTraffic.INDEX_NAME, query.toString(), processBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ProcessTraffic.INDEX_NAME,
query.toString(), processBuilder);
return buildProcesses(storageDataList);
}
@ -173,13 +183,14 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
public Process getProcess(String processId) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, ProcessTraffic.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, ProcessTraffic.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ID_IDX, processId);
query = client.addQueryIndexValue(ProcessTraffic.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(ProcessTraffic.INDEX_NAME, query, indexAndValueMap);
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(ProcessTraffic.INDEX_NAME, query.toString(), processBuilder);
List<? super StorageData> storageDataList = client.filterQuery(ProcessTraffic.INDEX_NAME,
query.toString(), processBuilder);
final List<Process> processes = buildProcesses(storageDataList);
return processes.size() > 0 ? processes.get(0) : null;
}
@ -205,7 +216,7 @@ public class IoTDBMetadataQueryDAO implements IMetadataQueryDAO {
storageDataList.forEach(storageData -> {
InstanceTraffic instanceTraffic = (InstanceTraffic) storageData;
if (instanceTraffic.getName() == null) {
instanceTraffic.setName("");
instanceTraffic.setName(Const.EMPTY_STRING);
}
ServiceInstance serviceInstance = new ServiceInstance();
serviceInstance.setId(instanceTraffic.id());

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.oap.server.storage.plugin.iotdb.query;
import com.google.common.base.Splitter;
import java.io.IOException;
import java.util.ArrayList;
import java.util.HashMap;
@ -45,6 +46,7 @@ import org.apache.skywalking.oap.server.core.storage.annotation.ValueColumnMetad
import org.apache.skywalking.oap.server.core.storage.query.IMetricsQueryDAO;
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.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -52,7 +54,10 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
private final IoTDBClient client;
@Override
public long readMetricsValue(MetricsCondition condition, String valueColumnName, Duration duration) throws IOException {
public long readMetricsValue(MetricsCondition condition,
String valueColumnName,
Duration duration)
throws IOException {
final int defaultValue = ValueColumnMetadata.INSTANCE.getDefaultValue(condition.getName());
final Function function = ValueColumnMetadata.INSTANCE.getValueFunction(condition.getName());
if (function == Function.Latest) {
@ -67,18 +72,20 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
op = "sum";
}
query.append(String.format("select %s(%s) from ", op, valueColumnName));
query = client.addModelPath(query, condition.getName());
IoTDBUtils.addModelPath(client.getStorageGroup(), query, condition.getName());
final String entityId = condition.getEntity().buildId();
if (entityId != null) {
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ENTITY_ID_IDX, entityId);
query = client.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
} else {
query = client.addQueryAsterisk(condition.getName(), query);
IoTDBUtils.addQueryAsterisk(condition.getName(), query);
}
query.append(" where ").append(String.format("%s >= %s and %s <= %s",
IoTDBClient.TIME, duration.getStartTimestamp(), IoTDBClient.TIME, duration.getEndTimestamp()))
.append(" group by level = 3");
query.append(" where ")
.append(String.format("%s >= %s and %s <= %s",
IoTDBClient.TIME, duration.getStartTimestamp(),
IoTDBClient.TIME, duration.getEndTimestamp()))
.append(" group by level = 3");
List<Double> results = client.queryWithAgg(query.toString());
if (results.size() > 0) {
@ -100,10 +107,10 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
StringBuilder query = new StringBuilder();
query.append("select ").append(valueColumnName).append(" from ");
for (String id : ids) {
query = client.addModelPath(query, condition.getName());
IoTDBUtils.addModelPath(client.getStorageGroup(), query, condition.getName());
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ID_IDX, id);
query = client.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
query.append(", ");
}
String queryString = query.toString();
@ -126,8 +133,9 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
while (wrapper.hasNext()) {
RowRecord rowRecord = wrapper.next();
List<Field> fields = rowRecord.getFields();
String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
String id = client.layerName2IndexValue(layerNames[1]);
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue());
String id = IoTDBUtils.layerName2IndexValue(layerNames.get(1));
Field valueField = fields.get(1);
TSDataType valueType = valueField.getDataType();
@ -150,12 +158,18 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
sessionPool.closeResultSet(wrapper);
}
}
metricsValues.setValues(Util.sortValues(intValues, ids, ValueColumnMetadata.INSTANCE.getDefaultValue(condition.getName())));
metricsValues.setValues(
Util.sortValues(intValues, ids,
ValueColumnMetadata.INSTANCE.getDefaultValue(condition.getName())));
return metricsValues;
}
@Override
public List<MetricsValues> readLabeledMetricsValues(MetricsCondition condition, String valueColumnName, List<String> labels, Duration duration) throws IOException {
public List<MetricsValues> readLabeledMetricsValues(MetricsCondition condition,
String valueColumnName,
List<String> labels,
Duration duration)
throws IOException {
final List<PointOfTime> pointOfTimes = duration.assembleDurationPoints();
List<String> ids = new ArrayList<>(pointOfTimes.size());
pointOfTimes.forEach(pointOfTime -> ids.add(pointOfTime.id(condition.getEntity().buildId())));
@ -163,10 +177,10 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
StringBuilder query = new StringBuilder();
query.append("select ").append(valueColumnName).append(" from ");
for (String id : ids) {
query = client.addModelPath(query, condition.getName());
IoTDBUtils.addModelPath(client.getStorageGroup(), query, condition.getName());
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ID_IDX, id);
query = client.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
query.append(", ");
}
String queryString = query.toString();
@ -187,8 +201,9 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
while (wrapper.hasNext()) {
RowRecord rowRecord = wrapper.next();
List<Field> fields = rowRecord.getFields();
String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
String id = client.layerName2IndexValue(layerNames[1]);
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue());
String id = IoTDBUtils.layerName2IndexValue(layerNames.get(1));
DataTable multipleValues = new DataTable(5);
multipleValues.toObject(fields.get(1).getStringValue());
@ -205,7 +220,10 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
}
@Override
public HeatMap readHeatMap(MetricsCondition condition, String valueColumnName, Duration duration) throws IOException {
public HeatMap readHeatMap(MetricsCondition condition,
String valueColumnName,
Duration duration)
throws IOException {
final List<PointOfTime> pointOfTimes = duration.assembleDurationPoints();
List<String> ids = new ArrayList<>(pointOfTimes.size());
pointOfTimes.forEach(pointOfTime -> ids.add(pointOfTime.id(condition.getEntity().buildId())));
@ -213,10 +231,10 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
StringBuilder query = new StringBuilder();
query.append("select ").append(valueColumnName).append(" from ");
for (String id : ids) {
query = client.addModelPath(query, condition.getName());
IoTDBUtils.addModelPath(client.getStorageGroup(), query, condition.getName());
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.ID_IDX, id);
query = client.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
query.append(", ");
}
String queryString = query.toString();
@ -238,8 +256,9 @@ public class IoTDBMetricsQueryDAO implements IMetricsQueryDAO {
while (wrapper.hasNext()) {
RowRecord rowRecord = wrapper.next();
List<Field> fields = rowRecord.getFields();
String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
String id = client.layerName2IndexValue(layerNames[1]);
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue());
String id = IoTDBUtils.layerName2IndexValue(layerNames.get(1));
heatMap.buildColumn(id, fields.get(1).getStringValue(), defaultValue);
}

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.oap.server.storage.plugin.iotdb.query;
import com.google.common.base.Splitter;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Comparator;
@ -45,6 +46,7 @@ import org.apache.skywalking.oap.server.library.util.StringUtil;
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.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -52,22 +54,27 @@ public class IoTDBTopNRecordsQueryDAO implements ITopNRecordsQueryDAO {
private final IoTDBClient client;
@Override
public List<SelectedRecord> readSampledRecords(TopNCondition condition, String valueColumnName, Duration duration) throws IOException {
public List<SelectedRecord> readSampledRecords(TopNCondition condition,
String valueColumnName,
Duration duration)
throws IOException {
StringBuilder query = new StringBuilder();
query.append("select ").append(TopN.STATEMENT).append(", ").append(valueColumnName)
.append(" from ");
query = client.addModelPath(query, condition.getName());
.append(" from ");
IoTDBUtils.addModelPath(client.getStorageGroup(), query, condition.getName());
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(condition.getParentService())) {
final String serviceId = IDManager.ServiceID.buildId(condition.getParentService(), condition.isNormal());
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
}
query = client.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(condition.getName(), query, indexAndValueMap);
StringBuilder where = new StringBuilder(" where ");
if (Objects.nonNull(duration)) {
where.append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(duration.getStartTimeBucketInSec())).append(" and ");
where.append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(duration.getEndTimeBucketInSec()));
where.append(IoTDBClient.TIME).append(" >= ")
.append(TimeBucket.getTimestamp(duration.getStartTimeBucketInSec())).append(" and ");
where.append(IoTDBClient.TIME).append(" <= ")
.append(TimeBucket.getTimestamp(duration.getEndTimeBucketInSec()));
}
if (where.length() > 7) {
query.append(where);
@ -92,8 +99,10 @@ public class IoTDBTopNRecordsQueryDAO implements ITopNRecordsQueryDAO {
List<Field> fields = rowRecord.getFields();
record.setName(fields.get(1).getStringValue());
String traceId = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"")[traceIdIdx + 1];
traceId = client.layerName2IndexValue(traceId);
String traceId = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue())
.get(traceIdIdx + 1);
traceId = IoTDBUtils.layerName2IndexValue(traceId);
record.setRefId(traceId);
record.setId(record.getRefId());
@ -111,7 +120,8 @@ public class IoTDBTopNRecordsQueryDAO implements ITopNRecordsQueryDAO {
// resort by self, because of the select query result order by time.
if (Order.DES.equals(condition.getOrder())) {
records.sort((SelectedRecord s1, SelectedRecord s2) ->
Long.compare(Long.parseLong(s2.getValue()), Long.parseLong(s1.getValue())));
Long.compare(Long.parseLong(s2.getValue()),
Long.parseLong(s1.getValue())));
} else {
records.sort(Comparator.comparingLong((SelectedRecord s) -> Long.parseLong(s.getValue())));
}

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.oap.server.storage.plugin.iotdb.query;
import com.google.common.base.Splitter;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
@ -39,6 +40,7 @@ import org.apache.skywalking.oap.server.core.query.type.Call;
import org.apache.skywalking.oap.server.core.source.DetectPoint;
import org.apache.skywalking.oap.server.core.storage.query.ITopologyQueryDAO;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.utils.IoTDBUtils;
@Slf4j
@RequiredArgsConstructor
@ -47,7 +49,8 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
@Override
public List<Call.CallDetail> loadServiceRelationsDetectedAtServerSide(long startTB, long endTB,
List<String> serviceIds) throws IOException {
List<String> serviceIds)
throws IOException {
return loadServiceCalls(
ServiceRelationServerSideMetrics.INDEX_NAME, startTB, endTB,
ServiceRelationServerSideMetrics.SOURCE_SERVICE_ID,
@ -57,7 +60,8 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
@Override
public List<Call.CallDetail> loadServiceRelationDetectedAtClientSide(long startTB, long endTB,
List<String> serviceIds) throws IOException {
List<String> serviceIds)
throws IOException {
return loadServiceCalls(
ServiceRelationClientSideMetrics.INDEX_NAME, startTB, endTB,
ServiceRelationClientSideMetrics.SOURCE_SERVICE_ID,
@ -66,7 +70,8 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
}
@Override
public List<Call.CallDetail> loadServiceRelationsDetectedAtServerSide(long startTB, long endTB) throws IOException {
public List<Call.CallDetail> loadServiceRelationsDetectedAtServerSide(long startTB, long endTB)
throws IOException {
return loadServiceCalls(
ServiceRelationServerSideMetrics.INDEX_NAME, startTB, endTB,
ServiceRelationServerSideMetrics.SOURCE_SERVICE_ID,
@ -75,7 +80,8 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
}
@Override
public List<Call.CallDetail> loadServiceRelationDetectedAtClientSide(long startTB, long endTB) throws IOException {
public List<Call.CallDetail> loadServiceRelationDetectedAtClientSide(long startTB, long endTB)
throws IOException {
return loadServiceCalls(
ServiceRelationClientSideMetrics.INDEX_NAME, startTB, endTB,
ServiceRelationClientSideMetrics.SOURCE_SERVICE_ID,
@ -86,7 +92,8 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
@Override
public List<Call.CallDetail> loadInstanceRelationDetectedAtServerSide(String clientServiceId,
String serverServiceId,
long startTB, long endTB) throws IOException {
long startTB, long endTB)
throws IOException {
return loadServiceInstanceCalls(
ServiceInstanceRelationServerSideMetrics.INDEX_NAME, startTB, endTB,
ServiceInstanceRelationServerSideMetrics.SOURCE_SERVICE_ID,
@ -97,7 +104,8 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
@Override
public List<Call.CallDetail> loadInstanceRelationDetectedAtClientSide(String clientServiceId,
String serverServiceId,
long startTB, long endTB) throws IOException {
long startTB, long endTB)
throws IOException {
return loadServiceInstanceCalls(
ServiceInstanceRelationClientSideMetrics.INDEX_NAME, startTB, endTB,
ServiceInstanceRelationClientSideMetrics.SOURCE_SERVICE_ID,
@ -106,7 +114,8 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
}
@Override
public List<Call.CallDetail> loadEndpointRelation(long startTB, long endTB, String destEndpointId) throws IOException {
public List<Call.CallDetail> loadEndpointRelation(long startTB, long endTB, String destEndpointId)
throws IOException {
List<Call.CallDetail> calls = loadEndpointFromSide(
EndpointRelationServerSideMetrics.INDEX_NAME, startTB, endTB,
EndpointRelationServerSideMetrics.SOURCE_ENDPOINT,
@ -122,20 +131,20 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
private List<Call.CallDetail> loadServiceCalls(String tableName, long startTB, long endTB,
String sourceCName, String destCName,
List<String> serviceIds, DetectPoint detectPoint) throws IOException {
List<String> serviceIds, DetectPoint detectPoint)
throws IOException {
// This method don't use "group by" like other storage plugin.
StringBuilder query = new StringBuilder();
query.append("select ").append(ServiceRelationServerSideMetrics.COMPONENT_ID).append(" from ");
query = client.addModelPath(query, tableName);
query = client.addQueryAsterisk(tableName, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, tableName);
IoTDBUtils.addQueryAsterisk(tableName, query);
query.append(" where ").append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(startTB))
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTB));
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTB));
if (serviceIds.size() > 0) {
query.append(" and (");
for (int i = 0; i < serviceIds.size(); i++) {
query.append(sourceCName).append(" = \"").append(serviceIds.get(i))
.append("\" or ")
.append(destCName).append(" = \"").append(serviceIds.get(i)).append("\"");
query.append(sourceCName).append(" = \"").append(serviceIds.get(i)).append("\" or ")
.append(destCName).append(" = \"").append(serviceIds.get(i)).append("\"");
if (i != serviceIds.size() - 1) {
query.append(" or ");
}
@ -157,8 +166,9 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
RowRecord rowRecord = wrapper.next();
List<Field> fields = rowRecord.getFields();
Call.CallDetail call = new Call.CallDetail();
String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
String entityId = client.layerName2IndexValue(layerNames[2]);
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue());
String entityId = IoTDBUtils.layerName2IndexValue(layerNames.get(2));
final int componentId = fields.get(1).getIntV();
call.buildFromServiceRelation(entityId, componentId, detectPoint);
calls.add(call);
@ -180,15 +190,15 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
// This method don't use "group by" like other storage plugin.
StringBuilder query = new StringBuilder();
query.append("select ").append(ServiceInstanceRelationServerSideMetrics.COMPONENT_ID).append(" from ");
query = client.addModelPath(query, tableName);
query = client.addQueryAsterisk(tableName, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, tableName);
IoTDBUtils.addQueryAsterisk(tableName, query);
query.append(" where ").append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(startTB))
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTB));
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTB));
query.append(" and ((").append(sourceCName).append(" = \"").append(sourceServiceId).append("\"")
.append(" and ").append(descCName).append(" = \"").append(destServiceId).append("\"")
.append(") or (").append(sourceCName).append(" = \"").append(destServiceId).append("\"")
.append(" and ").append(descCName).append(" = \"").append(sourceServiceId).append(" \"))")
.append(IoTDBClient.ALIGN_BY_DEVICE);
.append(" and ").append(descCName).append(" = \"").append(destServiceId).append("\"")
.append(") or (").append(sourceCName).append(" = \"").append(destServiceId).append("\"")
.append(" and ").append(descCName).append(" = \"").append(sourceServiceId).append(" \"))")
.append(IoTDBClient.ALIGN_BY_DEVICE);
SessionPool sessionPool = client.getSessionPool();
SessionDataSetWrapper wrapper = null;
@ -202,8 +212,9 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
RowRecord rowRecord = wrapper.next();
List<Field> fields = rowRecord.getFields();
Call.CallDetail call = new Call.CallDetail();
String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
String entityId = client.layerName2IndexValue(layerNames[2]);
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue());
String entityId = IoTDBUtils.layerName2IndexValue(layerNames.get(2));
final int componentId = fields.get(1).getIntV();
call.buildFromInstanceRelation(entityId, componentId, detectPoint);
calls.add(call);
@ -224,12 +235,12 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
// This method don't use "group by" like other storage plugin.
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, tableName);
query = client.addQueryAsterisk(tableName, query);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, tableName);
IoTDBUtils.addQueryAsterisk(tableName, query);
query.append(" where ").append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(startTB))
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTB));
.append(" and ").append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endTB));
query.append(" and ").append(isSourceId ? sourceCName : destCName).append(" = \"").append(id).append("\"")
.append(IoTDBClient.ALIGN_BY_DEVICE);
.append(IoTDBClient.ALIGN_BY_DEVICE);
SessionPool sessionPool = client.getSessionPool();
SessionDataSetWrapper wrapper = null;
@ -243,8 +254,9 @@ public class IoTDBTopologyQueryDAO implements ITopologyQueryDAO {
RowRecord rowRecord = wrapper.next();
List<Field> fields = rowRecord.getFields();
Call.CallDetail call = new Call.CallDetail();
String[] layerNames = fields.get(0).getStringValue().split("\\" + IoTDBClient.DOT + "\"");
String entityId = client.layerName2IndexValue(layerNames[2]);
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue());
String entityId = IoTDBUtils.layerName2IndexValue(layerNames.get(2));
call.buildFromEndpointRelation(entityId, DetectPoint.SERVER);
calls.add(call);
}

View File

@ -43,6 +43,7 @@ import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import org.apache.skywalking.oap.server.library.util.StringUtil;
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.utils.IoTDBUtils;
@RequiredArgsConstructor
public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
@ -59,7 +60,7 @@ public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
// This method maybe have poor efficiency. It queries all data which meets a condition without select function.
// https://github.com/apache/iotdb/discussions/3888
query.append("select * from ");
query = client.addModelPath(query, SegmentRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, SegmentRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
if (StringUtil.isNotEmpty(serviceId)) {
indexAndValueMap.put(IoTDBIndexes.SERVICE_ID_IDX, serviceId);
@ -67,12 +68,16 @@ public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
if (!Strings.isNullOrEmpty(traceId)) {
indexAndValueMap.put(IoTDBIndexes.TRACE_ID_IDX, traceId);
}
query = client.addQueryIndexValue(SegmentRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(SegmentRecord.INDEX_NAME, query, indexAndValueMap);
StringBuilder where = new StringBuilder(" where ");
if (startSecondTB != 0 && endSecondTB != 0) {
where.append(IoTDBClient.TIME).append(" >= ").append(TimeBucket.getTimestamp(startSecondTB)).append(" and ");
where.append(IoTDBClient.TIME).append(" <= ").append(TimeBucket.getTimestamp(endSecondTB)).append(" and ");
where.append(IoTDBClient.TIME).append(" >= ")
.append(TimeBucket.getTimestamp(startSecondTB))
.append(" and ");
where.append(IoTDBClient.TIME).append(" <= ")
.append(TimeBucket.getTimestamp(endSecondTB))
.append(" and ");
}
if (minDuration != 0) {
where.append(SegmentRecord.LATENCY).append(" >= ").append(minDuration).append(" and ");
@ -81,14 +86,20 @@ public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
where.append(SegmentRecord.LATENCY).append(" <= ").append(maxDuration).append(" and ");
}
if (StringUtil.isNotEmpty(serviceInstanceId)) {
where.append(SegmentRecord.SERVICE_INSTANCE_ID).append(" = \"").append(serviceInstanceId).append("\"").append(" and ");
where.append(SegmentRecord.SERVICE_INSTANCE_ID
).append(" = \"").append(serviceInstanceId).append("\"")
.append(" and ");
}
if (!Strings.isNullOrEmpty(endpointId)) {
where.append(SegmentRecord.ENDPOINT_ID).append(" = \"").append(endpointId).append("\"").append(" and ");
where.append(SegmentRecord.ENDPOINT_ID)
.append(" = \"").append(endpointId).append("\"")
.append(" and ");
}
if (CollectionUtils.isNotEmpty(tags)) {
for (final Tag tag : tags) {
where.append(tag.getKey()).append(" = \"").append(tag.getValue()).append("\"").append(" and ");
where.append(tag.getKey()).append(" = \"")
.append(tag.getValue()).append("\"")
.append(" and ");
}
}
switch (traceState) {
@ -107,7 +118,8 @@ public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
query.append(IoTDBClient.ALIGN_BY_DEVICE);
TraceBrief traceBrief = new TraceBrief();
List<? super StorageData> storageDataList = client.filterQuery(SegmentRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(SegmentRecord.INDEX_NAME,
query.toString(), storageBuilder);
int limitCount = 0;
for (int i = from; i < storageDataList.size(); i++) {
if (limitCount < limit) {
@ -116,7 +128,8 @@ public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
BasicTrace basicTrace = new BasicTrace();
basicTrace.setSegmentId(segmentRecord.getSegmentId());
basicTrace.setStart(String.valueOf(segmentRecord.getStartTime()));
basicTrace.getEndpointNames().add(IDManager.EndpointID.analysisId(segmentRecord.getEndpointId()).getEndpointName());
basicTrace.getEndpointNames().add(
IDManager.EndpointID.analysisId(segmentRecord.getEndpointId()).getEndpointName());
basicTrace.setDuration(segmentRecord.getLatency());
basicTrace.setError(BooleanUtils.valueToBoolean(segmentRecord.getIsError()));
basicTrace.getTraceIds().add(segmentRecord.getTraceId());
@ -128,11 +141,12 @@ public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
switch (queryOrder) {
case BY_START_TIME:
traceBrief.getTraces().sort((BasicTrace b1, BasicTrace b2) ->
Long.compare(Long.parseLong(b2.getStart()), Long.parseLong(b1.getStart())));
Long.compare(Long.parseLong(b2.getStart()),
Long.parseLong(b1.getStart())));
break;
case BY_DURATION:
traceBrief.getTraces().sort((BasicTrace b1, BasicTrace b2) ->
Integer.compare(b2.getDuration(), b1.getDuration()));
Integer.compare(b2.getDuration(), b1.getDuration()));
break;
}
return traceBrief;
@ -142,13 +156,14 @@ public class IoTDBTraceQueryDAO implements ITraceQueryDAO {
public List<SegmentRecord> queryByTraceId(String traceId) throws IOException {
StringBuilder query = new StringBuilder();
query.append("select * from ");
query = client.addModelPath(query, SegmentRecord.INDEX_NAME);
IoTDBUtils.addModelPath(client.getStorageGroup(), query, SegmentRecord.INDEX_NAME);
Map<String, String> indexAndValueMap = new HashMap<>();
indexAndValueMap.put(IoTDBIndexes.TRACE_ID_IDX, traceId);
query = client.addQueryIndexValue(SegmentRecord.INDEX_NAME, query, indexAndValueMap);
IoTDBUtils.addQueryIndexValue(SegmentRecord.INDEX_NAME, query, indexAndValueMap);
query.append(IoTDBClient.ALIGN_BY_DEVICE);
List<? super StorageData> storageDataList = client.filterQuery(SegmentRecord.INDEX_NAME, query.toString(), storageBuilder);
List<? super StorageData> storageDataList = client.filterQuery(SegmentRecord.INDEX_NAME,
query.toString(), storageBuilder);
List<SegmentRecord> segmentRecords = new ArrayList<>(storageDataList.size());
storageDataList.forEach(storageData -> segmentRecords.add((SegmentRecord) storageData));
return segmentRecords;

View File

@ -0,0 +1,189 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.storage.plugin.iotdb.utils;
import com.google.common.base.Splitter;
import java.util.ArrayList;
import java.util.Base64;
import java.util.List;
import java.util.function.Function;
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.Const;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.StorageDataComplexObject;
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
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.base.IoTDBInsertRequest;
public class IoTDBDataConverter {
public static class ToEntity implements Convert2Entity {
private final IoTDBTableMetaInfo tableMetaInfo;
private final List<String> indexValues;
private final List<String> columnNames;
private final long timestamp;
private final List<Field> fields;
public ToEntity(IoTDBTableMetaInfo tableMetaInfo, List<String> indexes,
List<String> columnNames, RowRecord rowRecord) {
this.tableMetaInfo = tableMetaInfo;
this.indexValues = new ArrayList<>(tableMetaInfo.getIndexes().size());
this.columnNames = columnNames;
this.timestamp = rowRecord.getTimestamp();
this.fields = rowRecord.getFields();
// field.get(0) -> Device, transform every layerName to indexValue
List<String> layerNames = Splitter.on(IoTDBClient.DOT + "\"")
.splitToList(fields.get(0).getStringValue());
for (int i = 0; i < indexes.size(); i++) {
indexValues.add(IoTDBUtils.layerName2IndexValue(layerNames.get(i + 1)));
}
}
@Override
public Object get(String fieldName) {
if (IoTDBIndexes.isIndex(fieldName)) {
if (tableMetaInfo.getIndexes().contains(fieldName)) {
String indexValue = indexValues.get(tableMetaInfo.getIndexes().indexOf(fieldName));
// convert String to Integer for LAYER_IDX value
return IoTDBIndexes.LAYER_IDX.equals(fieldName) ? Integer.valueOf(indexValue) : indexValue;
} else {
return null;
}
} else {
// time_bucket has changed to timestamp when writing to IoTDB.
if (IoTDBClient.TIME_BUCKET.equals(fieldName)) {
return TimeBucket.getTimeBucket(timestamp, tableMetaInfo.getModel().getDownsampling());
}
// IoTDB doesn't allow a measurement named `timestamp` or contains `.`,
// so we add double quotation mark to them.
// Also see IoTDBDataConverter.ToStorage.accept(String, Object)
if (IoTDBClient.TIMESTAMP.equals(fieldName) || fieldName.contains(".")) {
String columnName = IoTDBUtils.addQuotationMark(fieldName);
return columnNames.contains(columnName)
? IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(columnName) - 1))
: null;
} else {
return columnNames.contains(fieldName)
? IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(fieldName) - 1))
: null;
}
}
}
@Override
public <T, R> R getWith(final String fieldName, final Function<T, R> typeDecoder) {
if (columnNames.contains(fieldName)) {
final T value = (T) IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(fieldName) - 1));
return typeDecoder.apply(value);
} else {
return null;
}
}
}
public static class ToStorage implements Convert2Storage<IoTDBInsertRequest> {
private final IoTDBInsertRequest request;
private final IoTDBTableMetaInfo tableMetaInfo;
public ToStorage(String modelName, long time, String id) {
this.request = new IoTDBInsertRequest(modelName, time);
this.tableMetaInfo = IoTDBTableMetaInfo.get(modelName);
// set id as an index
this.request.getIndexValues().set(this.request.getIndexes().indexOf(IoTDBIndexes.ID_IDX), id);
}
@Override
public void accept(final String fieldName, final Object fieldValue) {
if (IoTDBIndexes.isIndex(fieldName)) {
List<String> indexes = request.getIndexes();
List<String> indexValues = request.getIndexValues();
// To avoid indexValue be "null" when inserting, replace null to empty string
if (indexes.contains(fieldName)) {
indexValues.set(indexes.indexOf(fieldName),
fieldValue == null ? Const.EMPTY_STRING : fieldValue.toString());
}
} else {
// time_bucket has changed to timestamp before calling this method,
// and IoTDB v0.12 doesn't allow insert null value,
// so they don't need to be stored.
if (fieldName.equals(IoTDBClient.TIME_BUCKET) || fieldValue == null) {
return;
}
List<String> measurements = request.getMeasurements();
List<TSDataType> measurementTypes = request.getMeasurementTypes();
List<Object> measurementValues = request.getMeasurementValues();
// IoTDB doesn't allow a measurement named `timestamp` or contains `.`
if (fieldName.equals(IoTDBClient.TIMESTAMP) || fieldName.contains(".")) {
measurements.add(IoTDBUtils.addQuotationMark(fieldName));
} else {
measurements.add(fieldName);
}
measurementTypes.add(tableMetaInfo.getColumnAndTypeMap().get(fieldName));
if (fieldValue instanceof StorageDataComplexObject) {
measurementValues.add(((StorageDataComplexObject) fieldValue).toStorageData());
} else {
measurementValues.add(fieldValue);
}
}
}
@Override
public void accept(final String fieldName, final byte[] fieldValue) {
request.getMeasurements().add(fieldName);
request.getMeasurementTypes().add(tableMetaInfo.getColumnAndTypeMap().get(fieldName));
if (CollectionUtils.isEmpty(fieldValue)) {
request.getMeasurementValues().add(Const.EMPTY_STRING);
} else {
request.getMeasurementValues().add(new String(Base64.getEncoder().encode(fieldValue)));
}
}
@Override
public void accept(final String fieldName, final List<String> fieldValue) {
request.getMeasurements().add(fieldName);
request.getMeasurementTypes().add(tableMetaInfo.getColumnAndTypeMap().get(fieldName));
request.getMeasurementValues().add(fieldValue);
}
@Override
public Object get(String fieldName) {
if (IoTDBIndexes.isIndex(fieldName)) {
return request.getIndexes().contains(fieldName)
? request.getIndexValues().get(request.getIndexes().indexOf(fieldName))
: null;
} else {
return request.getMeasurements().contains(fieldName)
? request.getMeasurementValues().get(request.getMeasurements().indexOf(fieldName))
: null;
}
}
@Override
public IoTDBInsertRequest obtain() {
return request;
}
}
}

View File

@ -0,0 +1,70 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.storage.plugin.iotdb.utils;
import java.util.List;
import java.util.Map;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.read.common.Field;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBTableMetaInfo;
public class IoTDBUtils {
public static String addQuotationMark(String string) {
return "\"" + string + "\"";
}
public static String indexValue2LayerName(String indexValue) {
return addQuotationMark(indexValue);
}
public static String layerName2IndexValue(String layerName) {
return layerName.substring(0, layerName.length() - 1);
}
public static void addQueryIndexValue(String modelName,
StringBuilder query,
Map<String, String> indexAndValueMap) {
List<String> 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("*");
}
});
}
public static void addQueryAsterisk(String modelName, StringBuilder query) {
List<String> indexes = IoTDBTableMetaInfo.get(modelName).getIndexes();
indexes.forEach(index -> query.append(IoTDBClient.DOT).append("*"));
}
public static void addModelPath(String storageGroup, StringBuilder query, String modelName) {
query.append(storageGroup).append(IoTDBClient.DOT).append(modelName);
}
public static Object getFieldValue(Field field) {
if (field.getDataType().equals(TSDataType.TEXT)) {
return field.getStringValue();
} else {
return field.getObjectValue(field.getDataType());
}
}
}

View File

@ -17,7 +17,7 @@ version: '2.1'
services:
iotdb:
image: apache/iotdb:0.12.4-node
image: apache/iotdb:0.12.5-node
expose:
- 6667
networks:

View File

@ -17,7 +17,7 @@ version: '3.8'
services:
iotdb:
image: apache/iotdb:0.12.4-node
image: apache/iotdb:0.12.5-node
expose:
- 6667
networks:

View File

@ -17,7 +17,7 @@ version: '2.1'
services:
iotdb:
image: apache/iotdb:0.12.4-node
image: apache/iotdb:0.12.5-node
expose:
- 6667
networks:

View File

@ -17,7 +17,7 @@ version: '2.1'
services:
iotdb:
image: apache/iotdb:0.12.4-node
image: apache/iotdb:0.12.5-node
expose:
- 6667
networks:

View File

@ -17,7 +17,7 @@ version: '2.1'
services:
iotdb:
image: apache/iotdb:0.12.4-node
image: apache/iotdb:0.12.5-node
expose:
- 6667
networks:

View File

@ -17,7 +17,7 @@ version: '2.1'
services:
iotdb:
image: apache/iotdb:0.12.4-node
image: apache/iotdb:0.12.5-node
expose:
- 6667
networks:

View File

@ -66,8 +66,8 @@ httpclient-4.5.13.jar
httpcore-4.4.13.jar
httpcore-nio-4.4.13.jar
influxdb-java-2.15.jar
iotdb-session-0.12.4.jar
iotdb-thrift-0.12.4.jar
iotdb-session-0.12.5.jar
iotdb-thrift-0.12.5.jar
j2objc-annotations-1.3.jar
jackson-annotations-2.12.2.jar
jackson-core-2.12.2.jar
@ -140,7 +140,7 @@ protobuf-java-3.19.2.jar
protobuf-java-util-3.19.2.jar
reactive-streams-1.0.2.jar
retrofit-2.5.0.jar
service-rpc-0.12.4.jar
service-rpc-0.12.5.jar
simpleclient-0.6.0.jar
simpleclient_common-0.6.0.jar
simpleclient_hotspot-0.6.0.jar
@ -149,7 +149,7 @@ slf4j-api-1.7.30.jar
snakeyaml-1.28.jar
snappy-java-1.1.7.3.jar
swagger-annotations-1.6.3.jar
tsfile-0.12.4.jar
tsfile-0.12.5.jar
vavr-0.10.3.jar
vavr-match-0.10.3.jar
zipkin-2.9.1.jar