Refactor `StorageData#id` to the new StorageID object from a String type. (#10157)
* Refactor kernel id0()/id() methods to return StorageID object, rather than a built string.
This commit is contained in:
parent
960aa9ffe1
commit
9a7bdc8eae
|
|
@ -40,6 +40,7 @@
|
|||
* Add global-specific settings used to override global configurations (e.g `segmentIntervalDays`, `blockIntervalHours`) in BanyanDB.
|
||||
* Use TTL-driven interval settings for the `measure-default` group in BanyanDB.
|
||||
* Fix wrong group of non time-relative metadata in BanyanDB.
|
||||
* Refactor `StorageData#id` to the new StorageID object from a String type.
|
||||
|
||||
#### UI
|
||||
|
||||
|
|
|
|||
|
|
@ -79,6 +79,6 @@ public class ProcessRegistry {
|
|||
traffic.setTimeBucket(timeBucket);
|
||||
traffic.setLastPingTimestamp(timeBucket);
|
||||
MetricsStreamProcessor.getInstance().in(traffic);
|
||||
return traffic.id();
|
||||
return traffic.id().build();
|
||||
}
|
||||
}
|
||||
|
|
@ -101,7 +101,7 @@ public class KafkaLogExporter extends KafkaExportProducer implements LogExportSe
|
|||
LogData logData = transLogData(logRecord);
|
||||
ProducerRecord<String, Bytes> record = new ProducerRecord<>(
|
||||
setting.getKafkaTopicLog(),
|
||||
logRecord.id(),
|
||||
logRecord.id().build(),
|
||||
Bytes.wrap(logData.toByteArray())
|
||||
);
|
||||
super.getProducer().send(record, (metadata, ex) -> {
|
||||
|
|
|
|||
|
|
@ -20,12 +20,13 @@ package org.apache.skywalking.oap.server.exporter.provider.grpc;
|
|||
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
|
||||
public class MockMetrics extends Metrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return "mock-metrics";
|
||||
protected StorageID id0() {
|
||||
return new StorageID().append("", "mock-metrics");
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -1,10 +1,9 @@
|
|||
protected String id0() {
|
||||
StringBuilder splitJointId = new StringBuilder(String.valueOf(getTimeBucket()));
|
||||
protected org.apache.skywalking.oap.server.core.storage.StorageID id0() {
|
||||
org.apache.skywalking.oap.server.core.storage.StorageID id = new org.apache.skywalking.oap.server.core.storage.StorageID().append(TIME_BUCKET, getTimeBucket());
|
||||
<#list fieldsFromSource as sourceField>
|
||||
<#if sourceField.isID()>
|
||||
splitJointId.append(org.apache.skywalking.oap.server.core.Const.ID_CONNECTOR)
|
||||
.append(${sourceField.fieldName});
|
||||
id.append("${sourceField.columnName}", ${sourceField.fieldName});
|
||||
</#if>
|
||||
</#list>
|
||||
return splitJointId.toString();
|
||||
return id;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -37,6 +37,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.MultiIntValuesHolder;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.joda.time.LocalDateTime;
|
||||
import org.joda.time.format.DateTimeFormat;
|
||||
import org.joda.time.format.DateTimeFormatter;
|
||||
|
|
@ -429,7 +430,7 @@ public class RunningRuleTest {
|
|||
private int value;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
@ -486,7 +487,7 @@ public class RunningRuleTest {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
@ -538,7 +539,7 @@ public class RunningRuleTest {
|
|||
private DataTable value;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -21,13 +21,13 @@ package org.apache.skywalking.oap.server.core.alarm;
|
|||
import java.util.List;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
|
|
@ -35,6 +35,7 @@ import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
|||
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.StorageBuilder;
|
||||
|
||||
import static org.apache.skywalking.oap.server.core.analysis.record.Record.TIME_BUCKET;
|
||||
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.ALARM;
|
||||
|
||||
|
|
@ -59,8 +60,12 @@ public class AlarmRecord extends Record {
|
|||
public static final String TAGS_RAW_DATA = "tags_raw_data";
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + ruleName + Const.ID_CONNECTOR + id0 + Const.ID_CONNECTOR + id1;
|
||||
public StorageID id() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(RULE_NAME, ruleName)
|
||||
.append(ID0, id0)
|
||||
.append(ID1, id1);
|
||||
}
|
||||
|
||||
@Column(columnName = SCOPE)
|
||||
|
|
|
|||
|
|
@ -24,13 +24,14 @@ import java.util.LinkedList;
|
|||
import java.util.List;
|
||||
import org.apache.skywalking.oap.server.core.storage.ComparableStorageData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
|
||||
/**
|
||||
* LimitedSizeBufferedData is a thread no safe implementation of {@link BufferedData}. It collects limited records of
|
||||
* each {@link StorageData#id()}.
|
||||
*/
|
||||
public class LimitedSizeBufferedData<STORAGE_DATA extends ComparableStorageData & StorageData> implements BufferedData<STORAGE_DATA> {
|
||||
private final HashMap<String, LinkedList<STORAGE_DATA>> data;
|
||||
private final HashMap<StorageID, LinkedList<STORAGE_DATA>> data;
|
||||
private final int limitedSize;
|
||||
|
||||
public LimitedSizeBufferedData(int limitedSize) {
|
||||
|
|
@ -40,7 +41,7 @@ public class LimitedSizeBufferedData<STORAGE_DATA extends ComparableStorageData
|
|||
|
||||
@Override
|
||||
public void accept(final STORAGE_DATA data) {
|
||||
final String id = data.id();
|
||||
final StorageID id = data.id();
|
||||
LinkedList<STORAGE_DATA> storageDataList = this.data.get(id);
|
||||
if (storageDataList == null) {
|
||||
storageDataList = new LinkedList<>();
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@ import java.util.List;
|
|||
import java.util.Map;
|
||||
import java.util.stream.Collectors;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
|
||||
/**
|
||||
* MergableBufferedData is a thread no safe implementation of {@link BufferedData}. {@link Metrics} in this cache would
|
||||
|
|
@ -31,7 +32,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
* Concurrency {@link #accept(Metrics)}s and {@link #read()} while {@link #accept(Metrics)} are both not recommended.
|
||||
*/
|
||||
public class MergableBufferedData<METRICS extends Metrics> implements BufferedData<METRICS> {
|
||||
private Map<String, METRICS> buffer;
|
||||
private Map<StorageID, METRICS> buffer;
|
||||
|
||||
public MergableBufferedData() {
|
||||
buffer = new HashMap<>();
|
||||
|
|
@ -46,7 +47,7 @@ public class MergableBufferedData<METRICS extends Metrics> implements BufferedDa
|
|||
*/
|
||||
@Override
|
||||
public void accept(final METRICS data) {
|
||||
final String id = data.id();
|
||||
final StorageID id = data.id();
|
||||
final METRICS existed = buffer.get(id);
|
||||
if (existed == null) {
|
||||
buffer.put(id, data);
|
||||
|
|
|
|||
|
|
@ -25,6 +25,7 @@ import org.apache.skywalking.oap.server.core.analysis.Stream;
|
|||
import org.apache.skywalking.oap.server.core.analysis.topn.TopN;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.TopNStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -47,8 +48,8 @@ public class TopNCacheReadCommand extends TopN {
|
|||
private String command;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return id;
|
||||
public StorageID id() {
|
||||
return new StorageID().appendMutant(null, id);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ import org.apache.skywalking.oap.server.core.analysis.Stream;
|
|||
import org.apache.skywalking.oap.server.core.analysis.topn.TopN;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.TopNStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -48,8 +49,8 @@ public class TopNCacheWriteCommand extends TopN {
|
|||
private String command;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return id;
|
||||
public StorageID id() {
|
||||
return new StorageID().appendMutant(null, id);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -25,6 +25,7 @@ import org.apache.skywalking.oap.server.core.analysis.Stream;
|
|||
import org.apache.skywalking.oap.server.core.analysis.topn.TopN;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.TopNStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -47,8 +48,8 @@ public class TopNDatabaseStatement extends TopN {
|
|||
private String statement;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return id;
|
||||
public StorageID id() {
|
||||
return new StorageID().appendMutant(null, id);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -65,11 +66,18 @@ public class EndpointTraffic extends Metrics {
|
|||
private String name = Const.EMPTY_STRING;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
// Downgrade the time bucket to day level only.
|
||||
// supportDownSampling == false for this entity.
|
||||
return IDManager.EndpointID.buildId(
|
||||
this.getServiceId(), this.getName());
|
||||
return new StorageID()
|
||||
.appendMutant(
|
||||
new String[] {
|
||||
SERVICE_ID,
|
||||
NAME
|
||||
},
|
||||
IDManager.EndpointID.buildId(
|
||||
this.getServiceId(), this.getName())
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -131,8 +132,12 @@ public class InstanceTraffic extends Metrics {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return IDManager.ServiceInstanceID.buildId(serviceId, name);
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.appendMutant(new String[] {
|
||||
SERVICE_ID,
|
||||
NAME
|
||||
}, IDManager.ServiceInstanceID.buildId(serviceId, name));
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<InstanceTraffic> {
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag;
|
|||
import org.apache.skywalking.oap.server.core.analysis.record.LongText;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.query.type.ContentType;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
|
|
@ -106,7 +107,7 @@ public abstract class AbstractLogRecord extends Record {
|
|||
private List<String> tagsInString;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
public StorageID id() {
|
||||
throw new UnexpectedException("AbstractLogRecord doesn't provide id()");
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@ import org.apache.skywalking.oap.server.core.analysis.Stream;
|
|||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -50,8 +51,8 @@ public class LogRecord extends AbstractLogRecord {
|
|||
private String uniqueId;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return uniqueId;
|
||||
public StorageID id() {
|
||||
return new StorageID().append(UNIQUE_ID, uniqueId);
|
||||
}
|
||||
|
||||
public static class Builder extends AbstractLogRecord.Builder<LogRecord> {
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -87,8 +88,9 @@ public class NetworkAddressAlias extends Metrics {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return IDManager.NetworkAddressAliasDefine.buildId(address);
|
||||
protected StorageID id0() {
|
||||
return new StorageID().appendMutant(
|
||||
new String[] {ADDRESS}, IDManager.NetworkAddressAliasDefine.buildId(address));
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,6 +21,7 @@ package org.apache.skywalking.oap.server.core.analysis.manual.process;
|
|||
import com.google.gson.Gson;
|
||||
import com.google.gson.JsonElement;
|
||||
import com.google.gson.JsonObject;
|
||||
import java.util.Map;
|
||||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
|
|
@ -32,6 +33,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -39,8 +41,6 @@ import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
|
|||
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
|
||||
import org.apache.skywalking.oap.server.library.util.StringUtil;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.PROCESS;
|
||||
|
||||
@Stream(name = ProcessTraffic.INDEX_NAME, scopeId = PROCESS,
|
||||
|
|
@ -183,11 +183,14 @@ public class ProcessTraffic extends Metrics {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
if (processId != null) {
|
||||
return processId;
|
||||
protected StorageID id0() {
|
||||
if (processId == null) {
|
||||
processId = IDManager.ProcessID.buildId(instanceId, name);
|
||||
}
|
||||
return IDManager.ProcessID.buildId(instanceId, name);
|
||||
return new StorageID().appendMutant(new String[] {
|
||||
INSTANCE_ID,
|
||||
NAME
|
||||
}, processId);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<ProcessTraffic> {
|
||||
|
|
|
|||
|
|
@ -20,22 +20,20 @@ package org.apache.skywalking.oap.server.core.analysis.manual.process;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.MetricsExtension;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
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.StorageBuilder;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.Base64;
|
||||
|
||||
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.SERVICE_LABEL;
|
||||
|
||||
/**
|
||||
|
|
@ -47,7 +45,7 @@ import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.SE
|
|||
@Getter
|
||||
@Stream(name = ServiceLabelRecord.INDEX_NAME, scopeId = SERVICE_LABEL,
|
||||
builder = ServiceLabelRecord.Builder.class, processor = MetricsStreamProcessor.class)
|
||||
@MetricsExtension(supportDownSampling = false, supportUpdate = false)
|
||||
@MetricsExtension(supportDownSampling = false, supportUpdate = false, timeRelativeID = false)
|
||||
@EqualsAndHashCode(of = {
|
||||
"serviceId",
|
||||
"label"
|
||||
|
|
@ -59,8 +57,10 @@ public class ServiceLabelRecord extends Metrics {
|
|||
public static final String SERVICE_ID = "service_id";
|
||||
public static final String LABEL = "label";
|
||||
|
||||
@BanyanDB.SeriesID(index = 0)
|
||||
@Column(columnName = SERVICE_ID)
|
||||
private String serviceId;
|
||||
@BanyanDB.SeriesID(index = 1)
|
||||
@Column(columnName = LABEL, length = 50)
|
||||
private String label;
|
||||
|
||||
|
|
@ -84,9 +84,10 @@ public class ServiceLabelRecord extends Metrics {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return this.serviceId + Const.ID_CONNECTOR + new String(Base64.getEncoder()
|
||||
.encode(label.getBytes(StandardCharsets.UTF_8)), StandardCharsets.UTF_8);
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(SERVICE_ID, serviceId)
|
||||
.append(LABEL, label);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.core.analysis.manual.relation.endpoint;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.MetricsExtension;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
|
|
@ -29,6 +28,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -72,10 +72,10 @@ public class EndpointRelationServerSideMetrics extends Metrics {
|
|||
private String entityId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
String splitJointId = String.valueOf(getTimeBucket());
|
||||
splitJointId += Const.ID_CONNECTOR + entityId;
|
||||
return splitJointId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,13 +21,13 @@ package org.apache.skywalking.oap.server.core.analysis.manual.relation.instance;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -78,8 +78,9 @@ public class ServiceInstanceRelationClientSideMetrics extends Metrics {
|
|||
private String entityId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID().append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, entityId);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,13 +21,13 @@ package org.apache.skywalking.oap.server.core.analysis.manual.relation.instance;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -80,8 +80,10 @@ public class ServiceInstanceRelationServerSideMetrics extends Metrics {
|
|||
private String entityId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.core.analysis.manual.relation.process;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.MetricsExtension;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
|
|
@ -29,6 +28,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -73,8 +73,10 @@ public class ProcessRelationClientSideMetrics extends Metrics {
|
|||
private int componentId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.core.analysis.manual.relation.process;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.MetricsExtension;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
|
|
@ -29,6 +28,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -72,8 +72,9 @@ public class ProcessRelationServerSideMetrics extends Metrics {
|
|||
private int componentId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID().append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, entityId);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,13 +21,13 @@ package org.apache.skywalking.oap.server.core.analysis.manual.relation.service;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -68,8 +68,9 @@ public class ServiceRelationClientSideMetrics extends Metrics {
|
|||
private String entityId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID().append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, entityId);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.core.analysis.manual.relation.service;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.MetricsExtension;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
|
|
@ -29,6 +28,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -72,8 +72,10 @@ public class ServiceRelationServerSideMetrics extends Metrics {
|
|||
private String entityId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -28,6 +28,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -87,8 +88,12 @@ public class TagAutocompleteData extends Metrics {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return toTimeBucketInDay() + "-" + tagType + "-" + tagKey + "=" + tagValue;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.appendMutant(new String[] {TIME_BUCKET}, toTimeBucketInDay())
|
||||
.append(TAG_TYPE, tagType)
|
||||
.append(TAG_KEY, tagKey)
|
||||
.append(TAG_VALUE, tagValue);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
|||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -106,8 +107,8 @@ public class SegmentRecord extends Record {
|
|||
private List<String> tags;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return segmentId;
|
||||
public StorageID id() {
|
||||
return new StorageID().append(SEGMENT_ID, segmentId);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<SegmentRecord> {
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -98,12 +99,18 @@ public class ServiceTraffic extends Metrics {
|
|||
* @return Base64 encode(serviceName) + "." + layer.value
|
||||
*/
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
String id;
|
||||
if (layer != null) {
|
||||
return encode(name) + Const.POINT + layer.value();
|
||||
id = encode(name) + Const.POINT + layer.value();
|
||||
} else {
|
||||
return encode(name) + Const.POINT + Layer.UNDEFINED.value();
|
||||
id = encode(name) + Const.POINT + Layer.UNDEFINED.value();
|
||||
}
|
||||
return new StorageID().appendMutant(new String[] {
|
||||
NAME,
|
||||
LAYER
|
||||
}, id);
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -20,11 +20,11 @@ package org.apache.skywalking.oap.server.core.analysis.manual.spanattach;
|
|||
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -81,8 +81,12 @@ public class SpanAttachedEventRecord extends Record {
|
|||
private long timestamp;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return traceSegmentId + Const.ID_CONNECTOR + startTimeSecond + Const.ID_CONNECTOR + startTimeNanos + Const.ID_CONNECTOR + event;
|
||||
public StorageID id() {
|
||||
return new StorageID()
|
||||
.append(TRACE_SEGMENT_ID, traceSegmentId)
|
||||
.append(START_TIME_SECOND, startTimeSecond)
|
||||
.append(START_TIME_NANOS, startTimeNanos)
|
||||
.append(EVENT, event);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<SpanAttachedEventRecord> {
|
||||
|
|
|
|||
|
|
@ -20,12 +20,12 @@ package org.apache.skywalking.oap.server.core.analysis.manual.trace;
|
|||
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.topn.TopN;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -67,8 +67,11 @@ public class SampledSlowTraceRecord extends Record {
|
|||
private long timestamp;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId + Const.ID_CONNECTOR + traceId;
|
||||
public StorageID id() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, entityId)
|
||||
.append(TRACE_ID, traceId);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<SampledSlowTraceRecord> {
|
||||
|
|
|
|||
|
|
@ -20,12 +20,12 @@ package org.apache.skywalking.oap.server.core.analysis.manual.trace;
|
|||
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.topn.TopN;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -68,8 +68,11 @@ public class SampledStatus4xxTraceRecord extends Record {
|
|||
private long timestamp;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId + Const.ID_CONNECTOR + traceId;
|
||||
public StorageID id() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, entityId)
|
||||
.append(TRACE_ID, traceId);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<SampledStatus4xxTraceRecord> {
|
||||
|
|
|
|||
|
|
@ -20,12 +20,12 @@ package org.apache.skywalking.oap.server.core.analysis.manual.trace;
|
|||
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.topn.TopN;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -68,8 +68,11 @@ public class SampledStatus5xxTraceRecord extends Record {
|
|||
private long timestamp;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId + Const.ID_CONNECTOR + traceId;
|
||||
public StorageID id() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, entityId)
|
||||
.append(TRACE_ID, traceId);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<SampledStatus5xxTraceRecord> {
|
||||
|
|
|
|||
|
|
@ -23,7 +23,6 @@ import lombok.Getter;
|
|||
import lombok.Setter;
|
||||
import lombok.ToString;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.MeterEntity;
|
||||
|
|
@ -31,6 +30,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.DataTable;
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.query.type.Bucket;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -141,8 +141,10 @@ public abstract class HistogramFunction extends Meter implements AcceptableValue
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -26,7 +26,6 @@ import lombok.Getter;
|
|||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.Setter;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.MeterEntity;
|
||||
|
|
@ -37,6 +36,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.MultiIntValuesHold
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.PercentileMetrics;
|
||||
import org.apache.skywalking.oap.server.core.query.type.Bucket;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
|
|
@ -251,8 +251,10 @@ public abstract class PercentileFunction extends Meter implements AcceptableValu
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -22,7 +22,6 @@ import java.util.Objects;
|
|||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import lombok.ToString;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.instance.InstanceTraffic;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
|
|
@ -36,6 +35,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
|
||||
import org.apache.skywalking.oap.server.core.query.sql.Function;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -151,8 +151,10 @@ public abstract class AvgFunction extends Meter implements AcceptableValue<Long>
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -23,7 +23,6 @@ import lombok.Getter;
|
|||
import lombok.Setter;
|
||||
import lombok.ToString;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.MeterEntity;
|
||||
|
|
@ -34,6 +33,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.DataTable;
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.query.type.Bucket;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
|
|
@ -172,8 +172,10 @@ public abstract class AvgHistogramFunction extends Meter implements AcceptableVa
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -30,7 +30,6 @@ import java.util.stream.IntStream;
|
|||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.MeterEntity;
|
||||
|
|
@ -43,6 +42,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.MultiIntValuesHolder;
|
||||
import org.apache.skywalking.oap.server.core.query.type.Bucket;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
|
|
@ -325,8 +325,10 @@ public abstract class AvgHistogramPercentileFunction extends Meter implements Ac
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -23,7 +23,6 @@ import java.util.Set;
|
|||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import lombok.ToString;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.instance.InstanceTraffic;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
|
|
@ -34,6 +33,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.DataTable;
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.LabeledValueHolder;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
|
|
@ -157,8 +157,10 @@ public abstract class AvgLabeledFunction extends Meter implements AcceptableValu
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -22,7 +22,6 @@ import java.util.Objects;
|
|||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import lombok.ToString;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.instance.InstanceTraffic;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
|
|
@ -35,6 +34,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
|
||||
import org.apache.skywalking.oap.server.core.query.sql.Function;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -136,8 +136,10 @@ public abstract class LatestFunction extends Meter implements AcceptableValue<Lo
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -22,7 +22,6 @@ import java.util.Objects;
|
|||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import lombok.ToString;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.instance.InstanceTraffic;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
|
|
@ -35,6 +34,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
|
||||
import org.apache.skywalking.oap.server.core.query.sql.Function;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -126,8 +126,10 @@ public abstract class SumFunction extends Meter implements AcceptableValue<Long>
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + getEntityId();
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -21,10 +21,14 @@ package org.apache.skywalking.oap.server.core.analysis.meter.function.sum;
|
|||
import com.google.common.base.Strings;
|
||||
import io.vavr.Tuple;
|
||||
import io.vavr.Tuple2;
|
||||
import java.util.Comparator;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.stream.Collector;
|
||||
import java.util.stream.IntStream;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.MeterEntity;
|
||||
|
|
@ -37,6 +41,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.MultiIntValuesHolder;
|
||||
import org.apache.skywalking.oap.server.core.query.type.Bucket;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
|
||||
|
|
@ -44,12 +49,6 @@ 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.StorageBuilder;
|
||||
|
||||
import java.util.Comparator;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.stream.Collector;
|
||||
import java.util.stream.IntStream;
|
||||
|
||||
import static java.util.stream.Collectors.groupingBy;
|
||||
import static java.util.stream.Collectors.mapping;
|
||||
|
||||
|
|
@ -290,8 +289,10 @@ public abstract class SumHistogramPercentileFunction extends Meter implements Ac
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -18,10 +18,10 @@
|
|||
|
||||
package org.apache.skywalking.oap.server.core.analysis.meter.function.sumpermin;
|
||||
|
||||
import java.util.Objects;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import lombok.ToString;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.instance.InstanceTraffic;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
|
|
@ -34,14 +34,13 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
|
||||
import org.apache.skywalking.oap.server.core.query.sql.Function;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
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.StorageBuilder;
|
||||
|
||||
import java.util.Objects;
|
||||
|
||||
@ToString
|
||||
@MeterFunction(functionName = "sumPerMin")
|
||||
public abstract class SumPerMinFunction extends Meter implements AcceptableValue<Long>, LongValueHolder {
|
||||
|
|
@ -133,8 +132,10 @@ public abstract class SumPerMinFunction extends Meter implements AcceptableValue
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + getEntityId();
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -18,9 +18,9 @@
|
|||
|
||||
package org.apache.skywalking.oap.server.core.analysis.meter.function.sumpermin;
|
||||
|
||||
import java.util.Objects;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.UnexpectedException;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.instance.InstanceTraffic;
|
||||
import org.apache.skywalking.oap.server.core.analysis.meter.Meter;
|
||||
|
|
@ -33,14 +33,13 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
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.StorageBuilder;
|
||||
|
||||
import java.util.Objects;
|
||||
|
||||
@MeterFunction(functionName = "sumPerMinLabeled")
|
||||
public abstract class SumPerMinLabeledFunction extends Meter implements AcceptableValue<DataTable>, LabeledValueHolder {
|
||||
|
||||
|
|
@ -124,8 +123,10 @@ public abstract class SumPerMinLabeledFunction extends Meter implements Acceptab
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + getEntityId();
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, getEntityId());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -18,9 +18,10 @@
|
|||
|
||||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EVENT;
|
||||
import static org.apache.skywalking.oap.server.library.util.StringUtil.isNotBlank;
|
||||
import com.google.common.base.Strings;
|
||||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Layer;
|
||||
import org.apache.skywalking.oap.server.core.analysis.MetricsExtension;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
|
|
@ -29,14 +30,15 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
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.StorageBuilder;
|
||||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
|
||||
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EVENT;
|
||||
import static org.apache.skywalking.oap.server.library.util.StringUtil.isNotBlank;
|
||||
|
||||
@Getter
|
||||
@Setter
|
||||
|
|
@ -77,8 +79,8 @@ public class Event extends Metrics {
|
|||
private static final int PARAMETER_MAX_LENGTH = 2000;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getUuid();
|
||||
protected StorageID id0() {
|
||||
return new StorageID().append(UUID, getUuid());
|
||||
}
|
||||
|
||||
@Column(columnName = UUID)
|
||||
|
|
|
|||
|
|
@ -25,6 +25,7 @@ import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
|
|||
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
||||
|
|
@ -149,15 +150,15 @@ public abstract class Metrics extends StreamData implements StorageData {
|
|||
return TimeBucket.isDayBucket(timeBucket);
|
||||
}
|
||||
|
||||
private volatile String id;
|
||||
private volatile StorageID id;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
public StorageID id() {
|
||||
if (id == null) {
|
||||
id = id0();
|
||||
}
|
||||
return id;
|
||||
}
|
||||
|
||||
protected abstract String id0();
|
||||
protected abstract StorageID id0();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -118,9 +118,6 @@ public class MetricsAggregateWorker extends AbstractWorker<Metrics> {
|
|||
if (currentTime - lastSendTime > l1FlushPeriod) {
|
||||
mergeDataCache.read().forEach(
|
||||
data -> {
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug(data.toString());
|
||||
}
|
||||
nextWorker.in(data);
|
||||
}
|
||||
);
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
|||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -50,8 +51,8 @@ public class BrowserErrorLogRecord extends Record {
|
|||
public static final String DATA_BINARY = "data_binary";
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return uniqueId;
|
||||
public StorageID id() {
|
||||
return new StorageID().append(UNIQUE_ID, uniqueId);
|
||||
}
|
||||
|
||||
@Setter
|
||||
|
|
|
|||
|
|
@ -25,6 +25,7 @@ import org.apache.skywalking.oap.server.core.analysis.Stream;
|
|||
import org.apache.skywalking.oap.server.core.analysis.management.ManagementData;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.ManagementStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
|
||||
|
|
@ -59,8 +60,8 @@ public class UITemplate extends ManagementData {
|
|||
private int disabled;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return templateId;
|
||||
public StorageID id() {
|
||||
return new StorageID().append(TEMPLATE_ID, templateId);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<UITemplate> {
|
||||
|
|
|
|||
|
|
@ -19,6 +19,16 @@
|
|||
package org.apache.skywalking.oap.server.core.profiling.ebpf;
|
||||
|
||||
import com.google.gson.Gson;
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.stream.Collectors;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.skywalking.oap.server.core.CoreModuleConfig;
|
||||
|
|
@ -53,17 +63,6 @@ import org.apache.skywalking.oap.server.library.module.Service;
|
|||
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
|
||||
import org.apache.skywalking.oap.server.library.util.StringUtil;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@Slf4j
|
||||
@RequiredArgsConstructor
|
||||
public class EBPFProfilingQueryService implements Service {
|
||||
|
|
@ -204,8 +203,8 @@ public class EBPFProfilingQueryService implements Service {
|
|||
final List<Metrics> processes = getProcessMetricsDAO().multiGet(processModel, processMetrics);
|
||||
|
||||
final Map<String, Process> processMap = processes.stream()
|
||||
.map(t -> (ProcessTraffic) t)
|
||||
.collect(Collectors.toMap(Metrics::id, this::convertProcess));
|
||||
.map(t -> (ProcessTraffic) t)
|
||||
.collect(Collectors.toMap(m -> m.id().build(), this::convertProcess));
|
||||
schedules.forEach(p -> p.setProcess(processMap.get(p.getProcessId())));
|
||||
}
|
||||
return schedules;
|
||||
|
|
@ -219,7 +218,7 @@ public class EBPFProfilingQueryService implements Service {
|
|||
|
||||
private Process convertProcess(ProcessTraffic traffic) {
|
||||
final Process process = new Process();
|
||||
process.setId(traffic.id());
|
||||
process.setId(traffic.id().build());
|
||||
process.setName(traffic.getName());
|
||||
final String serviceId = traffic.getServiceId();
|
||||
process.setServiceId(serviceId);
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ import lombok.Data;
|
|||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -64,12 +65,19 @@ public class EBPFProfilingDataRecord extends Record {
|
|||
private long uploadTime;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return Hashing.sha256().newHasher()
|
||||
.putString(scheduleId, Charsets.UTF_8)
|
||||
.putString(stackIdList, Charsets.UTF_8)
|
||||
.putLong(uploadTime)
|
||||
.hash().toString();
|
||||
public StorageID id() {
|
||||
return new StorageID().appendMutant(
|
||||
new String[] {
|
||||
SCHEDULE_ID,
|
||||
STACK_ID_LIST,
|
||||
UPLOAD_TIME
|
||||
},
|
||||
Hashing.sha256().newHasher()
|
||||
.putString(scheduleId, Charsets.UTF_8)
|
||||
.putString(stackIdList, Charsets.UTF_8)
|
||||
.putLong(uploadTime)
|
||||
.hash().toString()
|
||||
);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<EBPFProfilingDataRecord> {
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
|||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -43,12 +44,12 @@ import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EB
|
|||
@Setter
|
||||
@Getter
|
||||
@Stream(name = EBPFProfilingScheduleRecord.INDEX_NAME, scopeId = EBPF_PROFILING_SCHEDULE,
|
||||
builder = EBPFProfilingScheduleRecord.Builder.class, processor = MetricsStreamProcessor.class)
|
||||
builder = EBPFProfilingScheduleRecord.Builder.class, processor = MetricsStreamProcessor.class)
|
||||
@MetricsExtension(supportDownSampling = false, supportUpdate = true)
|
||||
@EqualsAndHashCode(of = {
|
||||
"taskId",
|
||||
"processId",
|
||||
"startTime",
|
||||
"taskId",
|
||||
"processId",
|
||||
"startTime",
|
||||
})
|
||||
@SQLDatabase.Sharding(shardingAlgorithm = ShardingAlgorithm.NO_SHARDING)
|
||||
public class EBPFProfilingScheduleRecord extends Metrics {
|
||||
|
|
@ -95,8 +96,8 @@ public class EBPFProfilingScheduleRecord extends Metrics {
|
|||
}
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return scheduleId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID().append(EBPF_PROFILING_SCHEDULE_ID, scheduleId);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -25,6 +25,7 @@ import org.apache.skywalking.oap.server.core.analysis.Stream;
|
|||
import org.apache.skywalking.oap.server.core.analysis.config.NoneStream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.NoneStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -39,7 +40,7 @@ import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EB
|
|||
@Data
|
||||
@ScopeDeclaration(id = EBPF_PROFILING_TASK, name = "EBPFProfilingTask")
|
||||
@Stream(name = EBPFProfilingTaskRecord.INDEX_NAME, scopeId = EBPF_PROFILING_TASK,
|
||||
builder = EBPFProfilingTaskRecord.Builder.class, processor = NoneStreamProcessor.class)
|
||||
builder = EBPFProfilingTaskRecord.Builder.class, processor = NoneStreamProcessor.class)
|
||||
@BanyanDB.TimestampColumn(EBPFProfilingTaskRecord.CREATE_TIME)
|
||||
public class EBPFProfilingTaskRecord extends NoneStream {
|
||||
public static final String INDEX_NAME = "ebpf_profiling_task";
|
||||
|
|
@ -83,11 +84,17 @@ public class EBPFProfilingTaskRecord extends NoneStream {
|
|||
private String extensionConfigJson;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return Hashing.sha256().newHasher()
|
||||
.putString(logicalId, Charsets.UTF_8)
|
||||
.putLong(createTime)
|
||||
.hash().toString();
|
||||
public StorageID id() {
|
||||
return new StorageID().appendMutant(
|
||||
new String[] {
|
||||
LOGICAL_ID,
|
||||
CREATE_TIME
|
||||
},
|
||||
Hashing.sha256().newHasher()
|
||||
.putString(logicalId, Charsets.UTF_8)
|
||||
.putLong(createTime)
|
||||
.hash().toString()
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -95,10 +102,10 @@ public class EBPFProfilingTaskRecord extends NoneStream {
|
|||
*/
|
||||
public void generateLogicalId() {
|
||||
this.logicalId = Hashing.sha256().newHasher()
|
||||
.putString(serviceId, Charsets.UTF_8)
|
||||
.putString(processLabelsJson, Charsets.UTF_8)
|
||||
.putLong(startTime)
|
||||
.hash().toString();
|
||||
.putString(serviceId, Charsets.UTF_8)
|
||||
.putString(processLabelsJson, Charsets.UTF_8)
|
||||
.putLong(startTime)
|
||||
.hash().toString();
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<EBPFProfilingTaskRecord> {
|
||||
|
|
|
|||
|
|
@ -20,11 +20,11 @@ package org.apache.skywalking.oap.server.core.profiling.trace;
|
|||
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -65,8 +65,12 @@ public class ProfileTaskLogRecord extends Record {
|
|||
private long timestamp;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return getTaskId() + Const.ID_CONNECTOR + getInstanceId() + Const.ID_CONNECTOR + getOperationType() + Const.ID_CONNECTOR + getOperationTime();
|
||||
public StorageID id() {
|
||||
return new StorageID()
|
||||
.append(TASK_ID, getTaskId())
|
||||
.append(INSTANCE_ID, getInstanceId())
|
||||
.append(OPERATION_TYPE, getOperationType())
|
||||
.append(OPERATION_TIME, getOperationTime());
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<ProfileTaskLogRecord> {
|
||||
|
|
|
|||
|
|
@ -97,7 +97,7 @@ public class ProfileTaskMutationService implements Service {
|
|||
task.setTimeBucket(TimeBucket.getMinuteTimeBucket(taskStartTime));
|
||||
NoneStreamProcessor.getInstance().in(task);
|
||||
|
||||
return ProfileTaskCreationResult.builder().id(task.id()).build();
|
||||
return ProfileTaskCreationResult.builder().id(task.id().build()).build();
|
||||
}
|
||||
|
||||
private String checkDataSuccess(final String serviceId,
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ import org.apache.skywalking.oap.server.core.analysis.Stream;
|
|||
import org.apache.skywalking.oap.server.core.analysis.config.NoneStream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.NoneStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -54,8 +55,8 @@ public class ProfileTaskRecord extends NoneStream {
|
|||
public static final String MAX_SAMPLING_COUNT = "max_sampling_count";
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return taskId;
|
||||
public StorageID id() {
|
||||
return new StorageID().append(TASK_ID, taskId);
|
||||
}
|
||||
|
||||
@Column(columnName = SERVICE_ID)
|
||||
|
|
|
|||
|
|
@ -20,11 +20,11 @@ package org.apache.skywalking.oap.server.core.profiling.trace;
|
|||
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -67,8 +67,11 @@ public class ProfileThreadSnapshotRecord extends Record {
|
|||
private byte[] stackBinary;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return getTaskId() + Const.ID_CONNECTOR + getSegmentId() + Const.ID_CONNECTOR + getSequence() + Const.ID_CONNECTOR;
|
||||
public StorageID id() {
|
||||
return new StorageID()
|
||||
.append(TASK_ID, getTaskId())
|
||||
.append(SEGMENT_ID, getSegmentId())
|
||||
.append(SEQUENCE, getSequence());
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<ProfileThreadSnapshotRecord> {
|
||||
|
|
|
|||
|
|
@ -18,12 +18,20 @@
|
|||
|
||||
package org.apache.skywalking.oap.server.core.query;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.Stream;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.skywalking.oap.server.core.CoreModule;
|
||||
import org.apache.skywalking.oap.server.core.analysis.IDManager;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.process.ProcessDetectType;
|
||||
import org.apache.skywalking.oap.server.core.analysis.manual.process.ProcessTraffic;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.config.IComponentLibraryCatalogService;
|
||||
import org.apache.skywalking.oap.server.core.query.type.Call;
|
||||
import org.apache.skywalking.oap.server.core.query.type.ProcessNode;
|
||||
|
|
@ -36,16 +44,6 @@ import org.apache.skywalking.oap.server.core.storage.model.Model;
|
|||
import org.apache.skywalking.oap.server.core.storage.model.StorageModels;
|
||||
import org.apache.skywalking.oap.server.library.module.ModuleManager;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.stream.Collectors;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
@Slf4j
|
||||
public class ProcessTopologyBuilder {
|
||||
private final IComponentLibraryCatalogService componentLibraryCatalogService;
|
||||
|
|
@ -88,7 +86,7 @@ public class ProcessTopologyBuilder {
|
|||
return p;
|
||||
}).collect(Collectors.toList())).stream()
|
||||
.map(t -> (ProcessTraffic) t)
|
||||
.collect(Collectors.toMap(Metrics::id, this::buildNode));
|
||||
.collect(Collectors.toMap(m -> m.id().build(), this::buildNode));
|
||||
|
||||
for (Call.CallDetail clientCall : clientCalls) {
|
||||
if (!callMap.containsKey(clientCall.getId())) {
|
||||
|
|
@ -128,7 +126,7 @@ public class ProcessTopologyBuilder {
|
|||
|
||||
private ProcessNode buildNode(ProcessTraffic traffic) {
|
||||
ProcessNode processNode = new ProcessNode();
|
||||
processNode.setId(traffic.id());
|
||||
processNode.setId(traffic.id().build());
|
||||
processNode.setServiceId(traffic.getServiceId());
|
||||
processNode.setServiceName(IDManager.ServiceID.analysisId(traffic.getServiceId()).getName());
|
||||
processNode.setServiceInstanceId(traffic.getInstanceId());
|
||||
|
|
|
|||
|
|
@ -25,5 +25,5 @@ public interface StorageData {
|
|||
/**
|
||||
* @return the unique id used in any storage option.
|
||||
*/
|
||||
String id();
|
||||
StorageID id();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,169 @@
|
|||
/*
|
||||
* 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.core.storage;
|
||||
|
||||
import com.google.common.base.Joiner;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.library.util.StringUtil;
|
||||
|
||||
/**
|
||||
* StorageID represents an identification for the metric or the record.
|
||||
* Typically, an ID is composited by two parts
|
||||
* 1. Time bucket based on downsampling.
|
||||
* 2. The encoded entity ID, such as Service ID.
|
||||
*
|
||||
* In the SQL database and ElasticSearch, the string ID is preferred.
|
||||
* In the BanyanDB, time series and entity ID(series ID) would be treated separately.
|
||||
*
|
||||
* @since 9.4.0 StorageID replaced the `string id()` method in the StorageData. An object-oriented ID provides a more
|
||||
* friendly interface for various database implementation.
|
||||
*/
|
||||
@EqualsAndHashCode(of = {
|
||||
"fragments"
|
||||
})
|
||||
public class StorageID {
|
||||
private final List<Fragment> fragments;
|
||||
/**
|
||||
* Once the storage ID was {@link #build()} or {@link #read()},
|
||||
* this object would switch to the sealed status, no more append is allowed.
|
||||
*/
|
||||
private boolean sealed = false;
|
||||
/**
|
||||
* The string ID would only be built once.
|
||||
*/
|
||||
private String builtID;
|
||||
|
||||
public StorageID() {
|
||||
fragments = new ArrayList<>(2);
|
||||
}
|
||||
|
||||
public StorageID append(String name, String value) {
|
||||
if (StringUtil.isBlank(name)) {
|
||||
throw new IllegalArgumentException("The name of storage ID should not be null or empty.");
|
||||
}
|
||||
if (sealed) {
|
||||
throw new IllegalStateException("The storage ID is sealed. Can't append a new fragment, name=" + name);
|
||||
}
|
||||
fragments.add(new Fragment(new String[] {name}, String.class, false, value));
|
||||
return this;
|
||||
}
|
||||
|
||||
public StorageID append(String name, long value) {
|
||||
if (StringUtil.isBlank(name)) {
|
||||
throw new IllegalArgumentException("The name of storage ID should not be null or empty.");
|
||||
}
|
||||
if (sealed) {
|
||||
throw new IllegalStateException("The storage ID is sealed. Can't append a new fragment, name=" + name);
|
||||
}
|
||||
fragments.add(new Fragment(new String[] {name}, Long.class, false, value));
|
||||
return this;
|
||||
}
|
||||
|
||||
public StorageID append(String name, int value) {
|
||||
if (StringUtil.isBlank(name)) {
|
||||
throw new IllegalArgumentException("The name of storage ID should not be null or empty.");
|
||||
}
|
||||
if (sealed) {
|
||||
throw new IllegalStateException("The storage ID is sealed. Can't append a new fragment, name=" + name);
|
||||
}
|
||||
fragments.add(new Fragment(new String[] {name}, Integer.class, false, value));
|
||||
return this;
|
||||
}
|
||||
|
||||
public StorageID appendMutant(String[] source, long value) {
|
||||
if (sealed) {
|
||||
throw new IllegalStateException("The storage ID is sealed. Can't append a new fragment, source=" + Arrays.toString(source));
|
||||
}
|
||||
fragments.add(new Fragment(source, Long.class, true, value));
|
||||
return this;
|
||||
}
|
||||
|
||||
public StorageID appendMutant(final String[] source, final String value) {
|
||||
if (sealed) {
|
||||
throw new IllegalStateException("The storage ID is sealed. Can't append a new fragment, source=" + Arrays.toString(source));
|
||||
}
|
||||
fragments.add(new Fragment(source, String.class, true, value));
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return the string ID concatenating the values of {@link #fragments} by the underline(_).
|
||||
*/
|
||||
public String build() {
|
||||
sealed = true;
|
||||
if (builtID == null) {
|
||||
builtID = Joiner.on(Const.ID_CONNECTOR).join(fragments);
|
||||
}
|
||||
return builtID;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return a read-only list to avoid unexpected change for metric ID.
|
||||
*/
|
||||
public List<Fragment> read() {
|
||||
sealed = true;
|
||||
return Collections.unmodifiableList(fragments);
|
||||
}
|
||||
|
||||
@RequiredArgsConstructor
|
||||
@Getter
|
||||
@EqualsAndHashCode(of = {
|
||||
"name",
|
||||
"value"
|
||||
}, doNotUseGetters = true)
|
||||
public static class Fragment {
|
||||
/**
|
||||
* The column name of the value, or the original column names of the mutate value.
|
||||
*
|
||||
* The names could be
|
||||
* 1. Always one column if this is not {@link #mutate} and from a certain persistent column.
|
||||
* 2. Be null if {@link #mutate} is true and no relative column, such as the original value is not in
|
||||
* the persistence.
|
||||
* 3. One or multi-values if the value is built through a symmetrical or asymmetrical encoding algorithm.
|
||||
*/
|
||||
private final String[] name;
|
||||
/**
|
||||
* Represent the class type of the {@link #value}.
|
||||
*/
|
||||
private final Class<?> type;
|
||||
/**
|
||||
* If true, the field was from {@link #name}, and value is changed by internal rules.
|
||||
* Such as time bucket downsampling, use a day-level time-bucket to build the ID for a minute dimension metric.
|
||||
*/
|
||||
private final boolean mutate;
|
||||
private final Object value;
|
||||
|
||||
public Optional<String[]> getName() {
|
||||
return Optional.ofNullable(name);
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return value.toString();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.core.zipkin;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.MetricsExtension;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
|
|
@ -29,6 +28,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -59,8 +59,10 @@ public class ZipkinServiceRelationTraffic extends Metrics {
|
|||
private String remoteServiceName;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return serviceName + Const.ID_CONNECTOR + remoteServiceName;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(SERVICE_NAME, serviceName)
|
||||
.append(REMOTE_SERVICE_NAME, remoteServiceName);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -60,8 +61,10 @@ public class ZipkinServiceSpanTraffic extends Metrics {
|
|||
private String spanName = Const.EMPTY_STRING;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return serviceName + Const.ID_CONNECTOR + spanName;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(SERVICE_NAME, serviceName)
|
||||
.append(SPAN_NAME, spanName);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
|
|||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
|
||||
|
|
@ -53,8 +54,8 @@ public class ZipkinServiceTraffic extends Metrics {
|
|||
private String serviceName = Const.EMPTY_STRING;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return serviceName;
|
||||
protected StorageID id0() {
|
||||
return new StorageID().append(SERVICE_NAME, serviceName);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ import org.apache.skywalking.oap.server.core.analysis.record.Record;
|
|||
import org.apache.skywalking.oap.server.core.analysis.worker.RecordStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
|
||||
|
|
@ -167,8 +168,8 @@ public class ZipkinSpanRecord extends Record {
|
|||
private List<String> query;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return traceId + Const.LINE + spanId;
|
||||
public StorageID id() {
|
||||
return new StorageID().append(TRACE_ID, traceId).append(SPAN_ID, spanId);
|
||||
}
|
||||
|
||||
public static class Builder implements StorageBuilder<ZipkinSpanRecord> {
|
||||
|
|
|
|||
|
|
@ -20,6 +20,7 @@ package org.apache.skywalking.oap.server.core.analysis.data;
|
|||
|
||||
import java.util.Objects;
|
||||
import org.apache.skywalking.oap.server.core.storage.ComparableStorageData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -63,8 +64,8 @@ public class LimitedSizeBufferedDataTest {
|
|||
}
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return "id";
|
||||
public StorageID id() {
|
||||
return new StorageID().append("ID", "id");
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -104,7 +105,7 @@ public class ApdexMetricsTest {
|
|||
public class ApdexMetricsImpl extends ApdexMetrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -56,7 +57,7 @@ public class CountMetricsTest {
|
|||
|
||||
public class CountMetricsImpl extends CountMetrics {
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -88,7 +89,7 @@ public class HeatMapMetricsTest {
|
|||
public class HistogramMetricsMocker extends HistogramMetrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -52,7 +53,7 @@ public class LongAvgMetricsTest {
|
|||
public class LongAvgMetricsImpl extends LongAvgMetrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -54,7 +55,7 @@ public class MaxLongMetricsTest {
|
|||
public class MaxLongMetricsImpl extends MaxLongMetrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -73,7 +74,7 @@ public class MetricsTest {
|
|||
public class MetricsMocker extends Metrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -59,7 +60,7 @@ public class MinLongMetricsTest {
|
|||
public class MinLongMetricsImpl extends MinLongMetrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -20,6 +20,7 @@ package org.apache.skywalking.oap.server.core.analysis.metrics;
|
|||
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.expression.StringMatch;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -67,7 +68,7 @@ public class PercentMetricsTest {
|
|||
public class PercentMetricsImpl extends PercentMetrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@
|
|||
package org.apache.skywalking.oap.server.core.analysis.metrics;
|
||||
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
|
|
@ -118,7 +119,7 @@ public class PercentileMetricsTest {
|
|||
public class PercentileMetricsMocker extends PercentileMetrics {
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
protected StorageID id0() {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -19,8 +19,8 @@ package org.apache.skywalking.oap.server.core.analysis.metrics.expression;
|
|||
|
||||
import org.junit.Test;
|
||||
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.assertFalse;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
public class NumberMatchTest {
|
||||
|
||||
|
|
|
|||
|
|
@ -116,8 +116,8 @@ public class PersistenceTimerTest {
|
|||
private final String id;
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return id;
|
||||
public StorageID id() {
|
||||
return new StorageID().append("ID", id);
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,46 @@
|
|||
/*
|
||||
* 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.core.storage;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
public class StorageIDTest {
|
||||
@Test
|
||||
public void testRawBuild() {
|
||||
StorageID id = new StorageID();
|
||||
id.append("time_bucket", 202212141438L) //2022-12-14 14:38
|
||||
.append("entity_id", "encoded-service-name");
|
||||
|
||||
Assert.assertEquals("202212141438_encoded-service-name", id.build());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEqual() {
|
||||
StorageID id = new StorageID();
|
||||
id.append("time_bucket", 202212141438L) //2022-12-14 14:38
|
||||
.append("entity_id", "encoded-service-name");
|
||||
|
||||
StorageID id2 = new StorageID();
|
||||
id2.append("time_bucket", 202212141438L) //2022-12-14 14:38
|
||||
.append("entity_id", "encoded-service-name");
|
||||
|
||||
Assert.assertEquals(true, id.equals(id2));
|
||||
}
|
||||
}
|
||||
|
|
@ -36,7 +36,7 @@ public class BulkConsumePool implements ConsumerPool {
|
|||
|
||||
public BulkConsumePool(String name, int size, long consumeCycle) {
|
||||
size = EnvUtil.getInt(name + "_THREAD", size);
|
||||
allConsumers = new ArrayList<MultipleChannelsConsumer>(size);
|
||||
allConsumers = new ArrayList<>(size);
|
||||
for (int i = 0; i < size; i++) {
|
||||
MultipleChannelsConsumer multipleChannelsConsumer = new MultipleChannelsConsumer("DataCarrier." + name + ".BulkConsumePool." + i + ".Thread", consumeCycle);
|
||||
multipleChannelsConsumer.setDaemon(true);
|
||||
|
|
|
|||
|
|
@ -95,7 +95,7 @@ public class MultipleChannelsConsumer extends Thread {
|
|||
public void addNewTarget(Channels channels, IConsumer consumer) {
|
||||
Group group = new Group(channels, consumer);
|
||||
// Recreate the new list to avoid change list while the list is used in consuming.
|
||||
ArrayList<Group> newList = new ArrayList<Group>();
|
||||
ArrayList<Group> newList = new ArrayList<>();
|
||||
for (Group target : consumeTargets) {
|
||||
newList.add(target);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -47,7 +47,7 @@ public class BanyanDBNoneStreamDAO extends AbstractDAO<BanyanDBStorageClient> im
|
|||
StreamWrite streamWrite = new StreamWrite(
|
||||
schema.getMetadata().getGroup(), // group name
|
||||
schema.getMetadata().name(), // stream-name
|
||||
noneStream.id() // identity
|
||||
noneStream.id().build() // identity
|
||||
); // set timestamp inside `BanyanDBConverter.StreamToStorage`
|
||||
Convert2Storage<StreamWrite> convert2Storage = new BanyanDBConverter.StreamToStorage(schema, streamWrite);
|
||||
storageBuilder.entity2Storage(noneStream, convert2Storage);
|
||||
|
|
|
|||
|
|
@ -101,11 +101,11 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
|
|||
this.getClient().define(applyStatus(uiTemplate));
|
||||
return TemplateChangeStatus.builder()
|
||||
.status(true)
|
||||
.id(uiTemplate.id())
|
||||
.id(uiTemplate.id().build())
|
||||
.build();
|
||||
} catch (IOException ioEx) {
|
||||
log.error("fail to disable the template", ioEx);
|
||||
return TemplateChangeStatus.builder().status(false).id(uiTemplate.id()).message("Can't disable the template")
|
||||
return TemplateChangeStatus.builder().status(false).id(uiTemplate.id().build()).message("Can't disable the template")
|
||||
.build();
|
||||
}
|
||||
}
|
||||
|
|
@ -133,7 +133,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
|
|||
}
|
||||
|
||||
public Property applyAll(UITemplate uiTemplate) {
|
||||
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id())
|
||||
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id().build())
|
||||
.addTag(TagAndValue.newStringTag(UITemplate.CONFIGURATION, uiTemplate.getConfiguration()))
|
||||
.addTag(TagAndValue.newLongTag(UITemplate.DISABLED, uiTemplate.getDisabled()))
|
||||
.addTag(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()))
|
||||
|
|
@ -147,7 +147,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
|
|||
* @return new property (patch) to be applied
|
||||
*/
|
||||
public Property applyStatus(UITemplate uiTemplate) {
|
||||
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id())
|
||||
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id().build())
|
||||
.addTag(TagAndValue.newLongTag(UITemplate.DISABLED, uiTemplate.getDisabled()))
|
||||
.addTag(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()))
|
||||
.build();
|
||||
|
|
@ -160,7 +160,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
|
|||
* @return new property (patch) to be applied
|
||||
*/
|
||||
public Property applyConfiguration(UITemplate uiTemplate) {
|
||||
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id())
|
||||
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id().build())
|
||||
.addTag(TagAndValue.newStringTag(UITemplate.CONFIGURATION, uiTemplate.getConfiguration()))
|
||||
.addTag(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()))
|
||||
.build();
|
||||
|
|
|
|||
|
|
@ -62,7 +62,7 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements IMetricsD
|
|||
@Override
|
||||
protected void apply(MeasureQuery query) {
|
||||
for (final Metrics missCachedMetric : metrics) {
|
||||
query.or(id(missCachedMetric.id()));
|
||||
query.or(id(missCachedMetric.id().build()));
|
||||
}
|
||||
}
|
||||
});
|
||||
|
|
@ -89,7 +89,7 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements IMetricsD
|
|||
TimeBucket.getTimestamp(metrics.getTimeBucket(), model.getDownsampling())); // timestamp
|
||||
final BanyanDBConverter.MeasureToStorage toStorage = new BanyanDBConverter.MeasureToStorage(schema, measureWrite);
|
||||
storageBuilder.entity2Storage(metrics, toStorage);
|
||||
toStorage.acceptID(metrics.id());
|
||||
toStorage.acceptID(metrics.id().build());
|
||||
return new BanyanDBMeasureInsertRequest(toStorage.obtain(), callback);
|
||||
}
|
||||
|
||||
|
|
@ -105,7 +105,7 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements IMetricsD
|
|||
TimeBucket.getTimestamp(metrics.getTimeBucket(), model.getDownsampling())); // timestamp
|
||||
final BanyanDBConverter.MeasureToStorage toStorage = new BanyanDBConverter.MeasureToStorage(schema, measureWrite);
|
||||
storageBuilder.entity2Storage(metrics, toStorage);
|
||||
toStorage.acceptID(metrics.id());
|
||||
toStorage.acceptID(metrics.id().build());
|
||||
return new BanyanDBMeasureUpdateRequest(toStorage.obtain());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -44,7 +44,7 @@ public class BanyanDBRecordDAO implements IRecordDAO {
|
|||
StreamWrite streamWrite = new StreamWrite(
|
||||
schema.getMetadata().getGroup(), // group name
|
||||
model.getName(), // index-name
|
||||
record.id() // identity
|
||||
record.id().build() // identity
|
||||
); // set timestamp inside `BanyanDBConverter.StreamToStorage`
|
||||
Convert2Storage<StreamWrite> convert2Storage = new BanyanDBConverter.StreamToStorage(schema, streamWrite);
|
||||
storageBuilder.entity2Storage(record, convert2Storage);
|
||||
|
|
|
|||
|
|
@ -38,7 +38,7 @@ public class ManagementEsDAO extends EsDAO implements IManagementDAO {
|
|||
@Override
|
||||
public void insert(Model model, ManagementData managementData) throws IOException {
|
||||
String tableName = IndexController.INSTANCE.getTableName(model);
|
||||
String docId = IndexController.INSTANCE.generateDocId(model, managementData.id());
|
||||
String docId = IndexController.INSTANCE.generateDocId(model, managementData.id().build());
|
||||
final boolean exist = getClient().existDoc(tableName, docId);
|
||||
if (exist) {
|
||||
return;
|
||||
|
|
|
|||
|
|
@ -65,7 +65,7 @@ public class MetricsEsDAO extends EsDAO implements IMetricsDAO {
|
|||
Map<String, List<String>> indexIdsGroup = new HashMap<>();
|
||||
groupIndices.forEach((tableName, metricList) -> {
|
||||
List<String> ids = metricList.stream()
|
||||
.map(item -> IndexController.INSTANCE.generateDocId(model, item.id()))
|
||||
.map(item -> IndexController.INSTANCE.generateDocId(model, item.id().build()))
|
||||
.collect(Collectors.toList());
|
||||
indexIdsGroup.put(tableName, ids);
|
||||
});
|
||||
|
|
@ -85,7 +85,7 @@ public class MetricsEsDAO extends EsDAO implements IMetricsDAO {
|
|||
});
|
||||
groupIndices.forEach((tableName, metricList) -> {
|
||||
List<String> ids = metricList.stream()
|
||||
.map(item -> IndexController.INSTANCE.generateDocId(model, item.id()))
|
||||
.map(item -> IndexController.INSTANCE.generateDocId(model, item.id().build()))
|
||||
.collect(Collectors.toList());
|
||||
final SearchResponse response = getClient().searchIDs(tableName, ids);
|
||||
response.getHits().getHits().forEach(hit -> {
|
||||
|
|
@ -103,7 +103,7 @@ public class MetricsEsDAO extends EsDAO implements IMetricsDAO {
|
|||
storageBuilder.entity2Storage(metrics, toStorage);
|
||||
Map<String, Object> builder = IndexController.INSTANCE.appendTableColumn(model, toStorage.obtain());
|
||||
String modelName = TimeSeriesUtils.writeIndexName(model, metrics.getTimeBucket());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, metrics.id());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, metrics.id().build());
|
||||
return new MetricIndexRequestWrapper(getClient().prepareInsert(modelName, id, builder), callback);
|
||||
}
|
||||
|
||||
|
|
@ -114,7 +114,7 @@ public class MetricsEsDAO extends EsDAO implements IMetricsDAO {
|
|||
Map<String, Object> builder =
|
||||
IndexController.INSTANCE.appendTableColumn(model, toStorage.obtain());
|
||||
String modelName = TimeSeriesUtils.writeIndexName(model, metrics.getTimeBucket());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, metrics.id());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, metrics.id().build());
|
||||
return new MetricIndexUpdateWrapper(getClient().prepareUpdate(modelName, id, builder), callback);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -45,7 +45,7 @@ public class NoneStreamEsDAO extends EsDAO implements INoneStreamDAO {
|
|||
Map<String, Object> builder =
|
||||
IndexController.INSTANCE.appendTableColumn(model, toStorage.obtain());
|
||||
String modelName = TimeSeriesUtils.writeIndexName(model, noneStream.getTimeBucket());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, noneStream.id());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, noneStream.id().build());
|
||||
getClient().forceInsert(modelName, id, builder);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -42,7 +42,7 @@ public class RecordEsDAO extends EsDAO implements IRecordDAO {
|
|||
storageBuilder.entity2Storage(record, toStorage);
|
||||
Map<String, Object> builder = IndexController.INSTANCE.appendTableColumn(model, toStorage.obtain());
|
||||
String modelName = TimeSeriesUtils.writeIndexName(model, record.getTimeBucket());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, record.id());
|
||||
String id = IndexController.INSTANCE.generateDocId(model, record.id().build());
|
||||
return getClient().prepareInsert(modelName, id, builder);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -228,7 +228,7 @@ public class MetadataQueryEsDAO extends EsDAO implements IMetadataQueryDAO {
|
|||
new EndpointTraffic.Builder().storage2Entity(new ElasticSearchConverter.ToEntity(EndpointTraffic.INDEX_NAME, sourceAsMap));
|
||||
|
||||
Endpoint endpoint = new Endpoint();
|
||||
endpoint.setId(endpointTraffic.id());
|
||||
endpoint.setId(endpointTraffic.id().build());
|
||||
endpoint.setName(endpointTraffic.getName());
|
||||
endpoints.add(endpoint);
|
||||
}
|
||||
|
|
@ -388,7 +388,7 @@ public class MetadataQueryEsDAO extends EsDAO implements IMetadataQueryDAO {
|
|||
new InstanceTraffic.Builder().storage2Entity(new ElasticSearchConverter.ToEntity(InstanceTraffic.INDEX_NAME, sourceAsMap));
|
||||
|
||||
ServiceInstance serviceInstance = new ServiceInstance();
|
||||
serviceInstance.setId(instanceTraffic.id());
|
||||
serviceInstance.setId(instanceTraffic.id().build());
|
||||
serviceInstance.setName(instanceTraffic.getName());
|
||||
serviceInstance.setInstanceUUID(serviceInstance.getId());
|
||||
|
||||
|
|
@ -420,7 +420,7 @@ public class MetadataQueryEsDAO extends EsDAO implements IMetadataQueryDAO {
|
|||
new ProcessTraffic.Builder().storage2Entity(new ElasticSearchConverter.ToEntity(ProcessTraffic.INDEX_NAME, sourceAsMap));
|
||||
|
||||
Process process = new Process();
|
||||
process.setId(processTraffic.id());
|
||||
process.setId(processTraffic.id().build());
|
||||
process.setName(processTraffic.getName());
|
||||
final String serviceId = processTraffic.getServiceId();
|
||||
process.setServiceId(serviceId);
|
||||
|
|
|
|||
|
|
@ -105,7 +105,7 @@ public class UITemplateManagementEsDAO extends EsDAO implements UITemplateManage
|
|||
final UITemplate.Builder builder = new UITemplate.Builder();
|
||||
final UITemplate uiTemplate = setting.toEntity();
|
||||
|
||||
final boolean exist = getClient().existDoc(UITemplate.INDEX_NAME, uiTemplate.id());
|
||||
final boolean exist = getClient().existDoc(UITemplate.INDEX_NAME, uiTemplate.id().build());
|
||||
if (exist) {
|
||||
return TemplateChangeStatus.builder().status(false).id(setting.getId()).message("Template exists")
|
||||
.build();
|
||||
|
|
@ -113,7 +113,7 @@ public class UITemplateManagementEsDAO extends EsDAO implements UITemplateManage
|
|||
|
||||
final ElasticSearchConverter.ToStorage toStorage = new ElasticSearchConverter.ToStorage(UITemplate.INDEX_NAME);
|
||||
builder.entity2Storage(uiTemplate, toStorage);
|
||||
getClient().forceInsert(UITemplate.INDEX_NAME, uiTemplate.id(), toStorage.obtain());
|
||||
getClient().forceInsert(UITemplate.INDEX_NAME, uiTemplate.id().build(), toStorage.obtain());
|
||||
return TemplateChangeStatus.builder().status(true).id(uiTemplate.getTemplateId()).build();
|
||||
} catch (Exception e) {
|
||||
log.error(e.getMessage(), e);
|
||||
|
|
@ -128,7 +128,7 @@ public class UITemplateManagementEsDAO extends EsDAO implements UITemplateManage
|
|||
final UITemplate.Builder builder = new UITemplate.Builder();
|
||||
final UITemplate uiTemplate = setting.toEntity();
|
||||
|
||||
final boolean exist = getClient().existDoc(UITemplate.INDEX_NAME, uiTemplate.id());
|
||||
final boolean exist = getClient().existDoc(UITemplate.INDEX_NAME, uiTemplate.id().build());
|
||||
if (!exist) {
|
||||
return TemplateChangeStatus.builder().status(false).id(setting.getId())
|
||||
.message("Can't find the template").build();
|
||||
|
|
@ -136,7 +136,7 @@ public class UITemplateManagementEsDAO extends EsDAO implements UITemplateManage
|
|||
|
||||
final ElasticSearchConverter.ToStorage toStorage = new ElasticSearchConverter.ToStorage(UITemplate.INDEX_NAME);
|
||||
builder.entity2Storage(uiTemplate, toStorage);
|
||||
getClient().forceUpdate(UITemplate.INDEX_NAME, uiTemplate.id(), toStorage.obtain());
|
||||
getClient().forceUpdate(UITemplate.INDEX_NAME, uiTemplate.id().build(), toStorage.obtain());
|
||||
return TemplateChangeStatus.builder().status(true).id(setting.getId()).build();
|
||||
} catch (Exception e) {
|
||||
log.error(e.getMessage(), e);
|
||||
|
|
@ -156,7 +156,7 @@ public class UITemplateManagementEsDAO extends EsDAO implements UITemplateManage
|
|||
|
||||
final ElasticSearchConverter.ToStorage toStorage = new ElasticSearchConverter.ToStorage(UITemplate.INDEX_NAME);
|
||||
builder.entity2Storage(uiTemplate, toStorage);
|
||||
getClient().forceUpdate(UITemplate.INDEX_NAME, uiTemplate.id(), toStorage.obtain());
|
||||
getClient().forceUpdate(UITemplate.INDEX_NAME, uiTemplate.id().build(), toStorage.obtain());
|
||||
return TemplateChangeStatus.builder().status(true).id(id).build();
|
||||
} else {
|
||||
return TemplateChangeStatus.builder().status(false).id(id).message("Can't find the template")
|
||||
|
|
|
|||
|
|
@ -42,7 +42,7 @@ public class JDBCManagementDAO extends JDBCSQLExecutor implements IManagementDAO
|
|||
@Override
|
||||
public void insert(Model model, ManagementData storageData) throws IOException {
|
||||
try (Connection connection = jdbcClient.getConnection()) {
|
||||
final StorageData data = getByID(jdbcClient, model.getName(), storageData.id(), storageBuilder);
|
||||
final StorageData data = getByID(jdbcClient, model.getName(), storageData.id().build(), storageBuilder);
|
||||
if (data != null) {
|
||||
return;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -40,7 +40,7 @@ public class JDBCMetricsDAO extends JDBCSQLExecutor implements IMetricsDAO {
|
|||
|
||||
@Override
|
||||
public List<Metrics> multiGet(Model model, List<Metrics> metrics) throws IOException {
|
||||
String[] ids = metrics.stream().map(Metrics::id).collect(Collectors.toList()).toArray(new String[] {});
|
||||
String[] ids = metrics.stream().map(m -> m.id().build()).collect(Collectors.toList()).toArray(new String[] {});
|
||||
List<StorageData> storageDataList = getByIDs(jdbcClient, model.getName(), ids, storageBuilder);
|
||||
List<Metrics> result = new ArrayList<>(storageDataList.size());
|
||||
for (StorageData storageData : storageDataList) {
|
||||
|
|
|
|||
|
|
@ -154,7 +154,7 @@ public class JDBCSQLExecutor {
|
|||
SQLBuilder sqlBuilder = new SQLBuilder("INSERT INTO " + tableName + " VALUES");
|
||||
List<Object> param = new ArrayList<>();
|
||||
sqlBuilder.append("(?,");
|
||||
param.add(metrics.id());
|
||||
param.add(metrics.id().build());
|
||||
for (int i = 0; i < columns.size(); i++) {
|
||||
ModelColumn column = columns.get(i);
|
||||
sqlBuilder.append("?");
|
||||
|
|
@ -185,7 +185,7 @@ public class JDBCSQLExecutor {
|
|||
SQLBuilder sqlBuilder = new SQLBuilder("INSERT INTO " + tableName + " VALUES");
|
||||
List<Object> param = new ArrayList<>();
|
||||
sqlBuilder.append("(?,");
|
||||
param.add(metrics.id());
|
||||
param.add(metrics.id().build());
|
||||
int position = 0;
|
||||
List valueList = new ArrayList();
|
||||
for (int i = 0; i < columns.size(); i++) {
|
||||
|
|
@ -257,7 +257,7 @@ public class JDBCSQLExecutor {
|
|||
}
|
||||
sqlBuilder.replace(sqlBuilder.length() - 1, sqlBuilder.length(), "");
|
||||
sqlBuilder.append(" WHERE id = ?");
|
||||
param.add(metrics.id());
|
||||
param.add(metrics.id().build());
|
||||
|
||||
return new SQLExecutor(sqlBuilder.toString(), param, callback);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -21,13 +21,13 @@ package org.apache.skywalking.oap.server.storage.plugin.jdbc.shardingsphere;
|
|||
import lombok.EqualsAndHashCode;
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.analysis.Stream;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.CPMMetrics;
|
||||
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
|
||||
import org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProcessor;
|
||||
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
|
||||
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
|
||||
import org.apache.skywalking.oap.server.core.storage.StorageID;
|
||||
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
|
||||
|
||||
@Stream(
|
||||
|
|
@ -51,8 +51,10 @@ public class ServiceCpmMetrics extends CPMMetrics {
|
|||
private String entityId;
|
||||
|
||||
@Override
|
||||
protected String id0() {
|
||||
return getTimeBucket() + Const.ID_CONNECTOR + entityId;
|
||||
protected StorageID id0() {
|
||||
return new StorageID()
|
||||
.append(TIME_BUCKET, getTimeBucket())
|
||||
.append(ENTITY_ID, entityId);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
Loading…
Reference in New Issue