diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java index d778e921e..17ab6e578 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/handler/TraceSegmentServiceHandler.java @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.agentstream.grpc.handler; import com.google.protobuf.InvalidProtocolBufferException; import io.grpc.stub.StreamObserver; import java.util.List; +import org.skywalking.apm.collector.agentstream.worker.segment.SegmentParse; import org.skywalking.apm.collector.server.grpc.GRPCHandler; import org.skywalking.apm.network.proto.Downstream; import org.skywalking.apm.network.proto.TraceSegmentObject; @@ -19,12 +20,15 @@ public class TraceSegmentServiceHandler extends TraceSegmentServiceGrpc.TraceSeg private final Logger logger = LoggerFactory.getLogger(TraceSegmentServiceHandler.class); + private SegmentParse segmentParse = new SegmentParse(); + @Override public StreamObserver collect(StreamObserver responseObserver) { return new StreamObserver() { @Override public void onNext(UpstreamSegment segment) { try { List traceIds = segment.getGlobalTraceIdsList(); TraceSegmentObject segmentObject = TraceSegmentObject.parseFrom(segment.getSegment()); + segmentParse.parse(traceIds, segmentObject); } catch (InvalidProtocolBufferException e) { logger.error(e.getMessage(), e); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java index 4f8b27ac4..4d5669c16 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/CommonTable.java @@ -4,6 +4,7 @@ package org.skywalking.apm.collector.agentstream.worker; * @author pengys5 */ public class CommonTable { + public static final String COLUMN_ID = "id"; public static final String COLUMN_AGG = "agg"; public static final String COLUMN_TIME_BUCKET = "time_bucket"; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggWorker.java index 005198b68..a3dec152e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentAggWorker.java @@ -1,5 +1,6 @@ package org.skywalking.apm.collector.agentstream.worker.node.component; +import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentDataDefine.java deleted file mode 100644 index 1d09a862e..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentDataDefine.java +++ /dev/null @@ -1,37 +0,0 @@ -package org.skywalking.apm.collector.agentstream.worker.node.component; - -import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; -import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; -import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; -import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; -import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; -import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; - -/** - * @author pengys5 - */ -public class NodeComponentDataDefine extends DataDefine { - - @Override public int defineId() { - return 0; - } - - @Override protected int initialCapacity() { - return 4; - } - - @Override protected void attributeDefine() { - addAttribute(0, new Attribute("id", AttributeType.STRING, new NonOperation())); - addAttribute(1, new Attribute("name", AttributeType.STRING, new CoverOperation())); - addAttribute(2, new Attribute("peers", AttributeType.STRING, new CoverOperation())); - addAttribute(3, new Attribute("aggregation", AttributeType.STRING, new CoverOperation())); - } - - @Override public Object deserialize(RemoteData remoteData) { - return null; - } - - @Override public RemoteData serialize(Object object) { - return null; - } -} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java new file mode 100644 index 000000000..ee9fca213 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java @@ -0,0 +1,52 @@ +package org.skywalking.apm.collector.agentstream.worker.node.component; + +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine; +import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; +import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; +import org.skywalking.apm.network.proto.SpanObject; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanListener, FirstSpanListener { + + private final Logger logger = LoggerFactory.getLogger(NodeComponentSpanListener.class); + + private List nodeComponents = new ArrayList<>(); + private long timeBucket; + + @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId) { + String peers = spanObject.getPeer(); + if (spanObject.getPeerId() == 0) { + peers = String.valueOf(spanObject.getPeerId()); + } + String agg = spanObject.getComponent() + Const.ID_SPLIT + peers; + nodeComponents.add(agg); + } + + @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId) { + String peers = String.valueOf(applicationId); + String agg = spanObject.getComponent() + Const.ID_SPLIT + peers; + nodeComponents.add(agg); + } + + @Override public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId) { + timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); + } + + @Override public void build() { + for (String agg : nodeComponents) { + NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); + nodeComponent.setId(timeBucket + Const.ID_SPLIT + agg); + nodeComponent.setAgg(agg); + nodeComponent.setTimeBucket(timeBucket); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java new file mode 100644 index 000000000..61941347b --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java @@ -0,0 +1,75 @@ +package org.skywalking.apm.collector.agentstream.worker.node.component.define; + +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; +import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; + +/** + * @author pengys5 + */ +public class NodeComponentDataDefine extends DataDefine { + + @Override public int defineId() { + return 101; + } + + @Override protected int initialCapacity() { + return 3; + } + + @Override protected void attributeDefine() { + addAttribute(0, new Attribute(NodeComponentTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); + addAttribute(1, new Attribute(NodeComponentTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); + addAttribute(2, new Attribute(NodeComponentTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + } + + @Override public Object deserialize(RemoteData remoteData) { + return null; + } + + @Override public RemoteData serialize(Object object) { + return null; + } + + public static class NodeComponent { + private String id; + private String agg; + private long timeBucket; + + public NodeComponent(String id, String agg, long timeBucket) { + this.id = id; + this.agg = agg; + this.timeBucket = timeBucket; + } + + public NodeComponent() { + } + + public String getId() { + return id; + } + + public String getAgg() { + return agg; + } + + public long getTimeBucket() { + return timeBucket; + } + + public void setId(String id) { + this.id = id; + } + + public void setAgg(String agg) { + this.agg = agg; + } + + public void setTimeBucket(long timeBucket) { + this.timeBucket = timeBucket; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java similarity index 70% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentEsTableDefine.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java index 01f76b785..700393edd 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java @@ -1,6 +1,5 @@ -package org.skywalking.apm.collector.agentstream.worker.node.define; +package org.skywalking.apm.collector.agentstream.worker.node.component.define; -import org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; @@ -26,8 +25,7 @@ public class NodeComponentEsTableDefine extends ElasticSearchTableDefine { } @Override public void initialize() { - addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); - addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_PEERS, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java similarity index 56% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentH2TableDefine.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java index 550890ba2..62b77b04e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeComponentH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java @@ -1,6 +1,5 @@ -package org.skywalking.apm.collector.agentstream.worker.node.define; +package org.skywalking.apm.collector.agentstream.worker.node.component.define; -import org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentTable; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; @@ -14,8 +13,8 @@ public class NodeComponentH2TableDefine extends H2TableDefine { } @Override public void initialize() { - addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_NAME, H2ColumnDefine.Type.Varchar.name())); - addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_PEERS, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_ID, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentTable.java similarity index 70% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentTable.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentTable.java index 211c8b5e6..1bce8df83 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentTable.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.node.component; +package org.skywalking.apm.collector.agentstream.worker.node.component.define; import org.skywalking.apm.collector.agentstream.worker.CommonTable; @@ -7,6 +7,4 @@ import org.skywalking.apm.collector.agentstream.worker.CommonTable; */ public class NodeComponentTable extends CommonTable { public static final String TABLE = "node_component"; - public static final String COLUMN_NAME = "name"; - public static final String COLUMN_PEERS = "peers"; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java new file mode 100644 index 000000000..6240656fb --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java @@ -0,0 +1,47 @@ +package org.skywalking.apm.collector.agentstream.worker.node.mapping; + +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine; +import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; +import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class NodeMappingSpanListener implements RefsListener, FirstSpanListener { + + private final Logger logger = LoggerFactory.getLogger(NodeMappingSpanListener.class); + + private List nodeMappings = new ArrayList<>(); + private long timeBucket; + + @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId) { + String peers = Const.PEERS_FRONT_SPLIT + reference.getNetworkAddressId() + Const.PEERS_BEHIND_SPLIT; + if (reference.getNetworkAddressId() == 0) { + peers = Const.PEERS_FRONT_SPLIT + reference.getNetworkAddress() + Const.PEERS_BEHIND_SPLIT; + } + + String agg = applicationId + Const.ID_SPLIT + peers; + nodeMappings.add(agg); + } + + @Override public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId) { + timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); + } + + @Override public void build() { + for (String agg : nodeMappings) { + NodeMappingDataDefine.NodeMapping nodeMapping = new NodeMappingDataDefine.NodeMapping(); + nodeMapping.setId(timeBucket + Const.ID_SPLIT + agg); + nodeMapping.setAgg(agg); + nodeMapping.setTimeBucket(timeBucket); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java new file mode 100644 index 000000000..f89b5d85e --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java @@ -0,0 +1,75 @@ +package org.skywalking.apm.collector.agentstream.worker.node.mapping.define; + +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; +import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; + +/** + * @author pengys5 + */ +public class NodeMappingDataDefine extends DataDefine { + + @Override public int defineId() { + return 102; + } + + @Override protected int initialCapacity() { + return 3; + } + + @Override protected void attributeDefine() { + addAttribute(0, new Attribute(NodeMappingTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); + addAttribute(1, new Attribute(NodeMappingTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); + addAttribute(2, new Attribute(NodeMappingTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + } + + @Override public Object deserialize(RemoteData remoteData) { + return null; + } + + @Override public RemoteData serialize(Object object) { + return null; + } + + public static class NodeMapping { + private String id; + private String agg; + private long timeBucket; + + public NodeMapping(String id, String agg, long timeBucket) { + this.id = id; + this.agg = agg; + this.timeBucket = timeBucket; + } + + public NodeMapping() { + } + + public String getId() { + return id; + } + + public String getAgg() { + return agg; + } + + public long getTimeBucket() { + return timeBucket; + } + + public void setId(String id) { + this.id = id; + } + + public void setAgg(String agg) { + this.agg = agg; + } + + public void setTimeBucket(long timeBucket) { + this.timeBucket = timeBucket; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java similarity index 73% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingEsTableDefine.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java index ef8a1a2ef..29223471a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.node.define; +package org.skywalking.apm.collector.agentstream.worker.node.mapping.define; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; @@ -25,8 +25,6 @@ public class NodeMappingEsTableDefine extends ElasticSearchTableDefine { } @Override public void initialize() { - addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); - addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_PEERS, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingH2TableDefine.java similarity index 67% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingH2TableDefine.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingH2TableDefine.java index f11217ccc..3805df48c 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingH2TableDefine.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.node.define; +package org.skywalking.apm.collector.agentstream.worker.node.mapping.define; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; @@ -13,8 +13,7 @@ public class NodeMappingH2TableDefine extends H2TableDefine { } @Override public void initialize() { - addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_NAME, H2ColumnDefine.Type.Varchar.name())); - addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_PEERS, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_ID, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingTable.java similarity index 53% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingTable.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingTable.java index 26f728fa9..811b297a2 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/define/NodeMappingTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingTable.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.node.define; +package org.skywalking.apm.collector.agentstream.worker.node.mapping.define; import org.skywalking.apm.collector.agentstream.worker.CommonTable; @@ -7,6 +7,4 @@ import org.skywalking.apm.collector.agentstream.worker.CommonTable; */ public class NodeMappingTable extends CommonTable { public static final String TABLE = "node_mapping"; - public static final String COLUMN_NAME = "name"; - public static final String COLUMN_PEERS = "peers"; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java new file mode 100644 index 000000000..19c76d606 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java @@ -0,0 +1,68 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.reference; + +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefDataDefine; +import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; +import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class NodeRefSpanListener implements EntrySpanListener, ExitSpanListener, FirstSpanListener, RefsListener { + + private final Logger logger = LoggerFactory.getLogger(NodeRefSpanListener.class); + + private List nodeExitReferences = new ArrayList<>(); + private List nodeEntryReferences = new ArrayList<>(); + private long timeBucket; + private boolean hasReference = false; + + @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId) { + String front = String.valueOf(applicationId); + String behind = String.valueOf(spanObject.getPeerId()); + if (spanObject.getPeerId() == 0) { + behind = spanObject.getPeer(); + } + + String agg = front + Const.ID_SPLIT + behind; + nodeExitReferences.add(agg); + } + + @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId) { + String behind = String.valueOf(applicationId); + String front = Const.USER_CODE; + + String agg = front + Const.ID_SPLIT + behind; + nodeEntryReferences.add(agg); + } + + @Override public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId) { + timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); + } + + @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId) { + hasReference = true; + } + + @Override public void build() { + if (!hasReference) { + nodeExitReferences.addAll(nodeEntryReferences); + } + + for (String agg : nodeExitReferences) { + NodeRefDataDefine.NodeReference nodeReference = new NodeRefDataDefine.NodeReference(); + nodeReference.setId(timeBucket + Const.ID_SPLIT + agg); + nodeReference.setAgg(agg); + nodeReference.setTimeBucket(timeBucket); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java new file mode 100644 index 000000000..09031c8a9 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java @@ -0,0 +1,85 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define; + +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; +import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; + +/** + * @author pengys5 + */ +public class NodeRefDataDefine extends DataDefine { + + public static final int DEFINE_ID = 201; + + @Override public int defineId() { + return DEFINE_ID; + } + + @Override protected int initialCapacity() { + return 3; + } + + @Override protected void attributeDefine() { + addAttribute(0, new Attribute(NodeRefTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); + addAttribute(1, new Attribute(NodeRefTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); + addAttribute(2, new Attribute(NodeRefTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + } + + @Override public Object deserialize(RemoteData remoteData) { + String id = remoteData.getDataStrings(0); + String agg = remoteData.getDataStrings(1); + long timeBucket = remoteData.getDataLongs(0); + return new NodeReference(id, agg, timeBucket); + } + + @Override public RemoteData serialize(Object object) { + NodeReference nodeReference = (NodeReference)object; + RemoteData.Builder builder = RemoteData.newBuilder(); + builder.addDataStrings(nodeReference.getId()); + builder.addDataStrings(nodeReference.getAgg()); + builder.addDataLongs(nodeReference.getTimeBucket()); + return builder.build(); + } + + public static class NodeReference { + private String id; + private String agg; + private long timeBucket; + + public NodeReference(String id, String agg, long timeBucket) { + this.id = id; + this.agg = agg; + this.timeBucket = timeBucket; + } + + public NodeReference() { + } + + public String getId() { + return id; + } + + public String getAgg() { + return agg; + } + + public long getTimeBucket() { + return timeBucket; + } + + public void setId(String id) { + this.id = id; + } + + public void setAgg(String agg) { + this.agg = agg; + } + + public void setTimeBucket(long timeBucket) { + this.timeBucket = timeBucket; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefEsTableDefine.java new file mode 100644 index 000000000..7ba2e614f --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefEsTableDefine.java @@ -0,0 +1,31 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define; + +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; + +/** + * @author pengys5 + */ +public class NodeRefEsTableDefine extends ElasticSearchTableDefine { + + public NodeRefEsTableDefine() { + super(NodeRefTable.TABLE); + } + + @Override public int refreshInterval() { + return 0; + } + + @Override public int numberOfShards() { + return 2; + } + + @Override public int numberOfReplicas() { + return 0; + } + + @Override public void initialize() { + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefH2TableDefine.java new file mode 100644 index 000000000..c2815b1de --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefH2TableDefine.java @@ -0,0 +1,19 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define; + +import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; +import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; + +/** + * @author pengys5 + */ +public class NodeRefH2TableDefine extends H2TableDefine { + + public NodeRefH2TableDefine() { + super(NodeRefTable.TABLE); + } + + @Override public void initialize() { + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefTable.java new file mode 100644 index 000000000..17c058b09 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefTable.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define; + +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + +/** + * @author pengys5 + */ +public class NodeRefTable extends CommonTable { + public static final String TABLE = "node_reference"; +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumDataDefine.java new file mode 100644 index 000000000..758ecc16a --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumDataDefine.java @@ -0,0 +1,164 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define; + +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; +import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; + +/** + * @author pengys5 + */ +public class NodeRefSumDataDefine extends DataDefine { + + public static final int DEFINE_ID = 202; + + @Override public int defineId() { + return DEFINE_ID; + } + + @Override protected int initialCapacity() { + return 9; + } + + @Override protected void attributeDefine() { + addAttribute(0, new Attribute(NodeRefSumTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); + addAttribute(1, new Attribute(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, AttributeType.LONG, new NonOperation())); + addAttribute(2, new Attribute(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, AttributeType.LONG, new NonOperation())); + addAttribute(3, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, AttributeType.LONG, new NonOperation())); + addAttribute(4, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, AttributeType.LONG, new NonOperation())); + addAttribute(5, new Attribute(NodeRefSumTable.COLUMN_ERROR, AttributeType.LONG, new NonOperation())); + addAttribute(6, new Attribute(NodeRefSumTable.COLUMN_SUMMARY, AttributeType.LONG, new NonOperation())); + addAttribute(7, new Attribute(NodeRefSumTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); + addAttribute(8, new Attribute(NodeRefSumTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + } + + @Override public Object deserialize(RemoteData remoteData) { + String id = remoteData.getDataStrings(0); + String agg = remoteData.getDataStrings(1); + Long oneSecondLess = remoteData.getDataLongs(0); + Long threeSecondLess = remoteData.getDataLongs(1); + Long fiveSecondLess = remoteData.getDataLongs(2); + Long fiveSecondGreater = remoteData.getDataLongs(3); + Long error = remoteData.getDataLongs(4); + Long summary = remoteData.getDataLongs(5); + long timeBucket = remoteData.getDataLongs(6); + return new NodeReferenceSum(id, oneSecondLess, threeSecondLess, fiveSecondLess, fiveSecondGreater, error, summary, agg, timeBucket); + } + + @Override public RemoteData serialize(Object object) { + NodeReferenceSum nodeReferenceSum = (NodeReferenceSum)object; + RemoteData.Builder builder = RemoteData.newBuilder(); + builder.addDataStrings(nodeReferenceSum.getId()); + builder.addDataStrings(nodeReferenceSum.getAgg()); + builder.addDataLongs(nodeReferenceSum.getOneSecondLess()); + builder.addDataLongs(nodeReferenceSum.getThreeSecondLess()); + builder.addDataLongs(nodeReferenceSum.getFiveSecondLess()); + builder.addDataLongs(nodeReferenceSum.getFiveSecondGreater()); + builder.addDataLongs(nodeReferenceSum.getError()); + builder.addDataLongs(nodeReferenceSum.getSummary()); + builder.addDataLongs(nodeReferenceSum.getTimeBucket()); + return builder.build(); + } + + public static class NodeReferenceSum { + private String id; + private Long oneSecondLess; + private Long threeSecondLess; + private Long fiveSecondLess; + private Long fiveSecondGreater; + private Long error; + private Long summary; + private String agg; + private long timeBucket; + + public NodeReferenceSum(String id, Long oneSecondLess, Long threeSecondLess, Long fiveSecondLess, + Long fiveSecondGreater, Long error, Long summary, String agg, long timeBucket) { + this.id = id; + this.oneSecondLess = oneSecondLess; + this.threeSecondLess = threeSecondLess; + this.fiveSecondLess = fiveSecondLess; + this.fiveSecondGreater = fiveSecondGreater; + this.error = error; + this.summary = summary; + this.agg = agg; + this.timeBucket = timeBucket; + } + + public NodeReferenceSum() { + } + + public String getId() { + return id; + } + + public void setId(String id) { + this.id = id; + } + + public Long getOneSecondLess() { + return oneSecondLess; + } + + public void setOneSecondLess(Long oneSecondLess) { + this.oneSecondLess = oneSecondLess; + } + + public Long getThreeSecondLess() { + return threeSecondLess; + } + + public void setThreeSecondLess(Long threeSecondLess) { + this.threeSecondLess = threeSecondLess; + } + + public Long getFiveSecondLess() { + return fiveSecondLess; + } + + public void setFiveSecondLess(Long fiveSecondLess) { + this.fiveSecondLess = fiveSecondLess; + } + + public Long getFiveSecondGreater() { + return fiveSecondGreater; + } + + public void setFiveSecondGreater(Long fiveSecondGreater) { + this.fiveSecondGreater = fiveSecondGreater; + } + + public Long getError() { + return error; + } + + public void setError(Long error) { + this.error = error; + } + + public Long getSummary() { + return summary; + } + + public void setSummary(Long summary) { + this.summary = summary; + } + + public String getAgg() { + return agg; + } + + public void setAgg(String agg) { + this.agg = agg; + } + + public long getTimeBucket() { + return timeBucket; + } + + public void setTimeBucket(long timeBucket) { + this.timeBucket = timeBucket; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumEsTableDefine.java new file mode 100644 index 000000000..c10b0d4a6 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumEsTableDefine.java @@ -0,0 +1,37 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define; + +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; + +/** + * @author pengys5 + */ +public class NodeRefSumEsTableDefine extends ElasticSearchTableDefine { + + public NodeRefSumEsTableDefine() { + super(NodeRefSumTable.TABLE); + } + + @Override public int refreshInterval() { + return 0; + } + + @Override public int numberOfShards() { + return 2; + } + + @Override public int numberOfReplicas() { + return 0; + } + + @Override public void initialize() { + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_ERROR, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_SUMMARY, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumH2TableDefine.java new file mode 100644 index 000000000..a31257102 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumH2TableDefine.java @@ -0,0 +1,25 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define; + +import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; +import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; + +/** + * @author pengys5 + */ +public class NodeRefSumH2TableDefine extends H2TableDefine { + + public NodeRefSumH2TableDefine() { + super(NodeRefSumTable.TABLE); + } + + @Override public void initialize() { + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_ERROR, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_SUMMARY, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumTable.java new file mode 100644 index 000000000..0c07ae9a8 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumTable.java @@ -0,0 +1,16 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define; + +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + +/** + * @author pengys5 + */ +public class NodeRefSumTable extends CommonTable { + public static final String TABLE = "node_reference_sum"; + public static final String COLUMN_ONE_SECOND_LESS = "one_second_less"; + public static final String COLUMN_THREE_SECOND_LESS = "three_second_less"; + public static final String COLUMN_FIVE_SECOND_LESS = "five_second_less"; + public static final String COLUMN_FIVE_SECOND_GREATER = "five_second_greater"; + public static final String COLUMN_ERROR = "error"; + public static final String COLUMN_SUMMARY = "summary"; +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/EntrySpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/EntrySpanListener.java new file mode 100644 index 000000000..1879681b6 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/EntrySpanListener.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.agentstream.worker.segment; + +import org.skywalking.apm.network.proto.SpanObject; + +/** + * @author pengys5 + */ +public interface EntrySpanListener extends SpanListener { + void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId); +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/ExitSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/ExitSpanListener.java new file mode 100644 index 000000000..baace12bd --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/ExitSpanListener.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.agentstream.worker.segment; + +import org.skywalking.apm.network.proto.SpanObject; + +/** + * @author pengys5 + */ +public interface ExitSpanListener extends SpanListener { + void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId); +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/FirstSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/FirstSpanListener.java new file mode 100644 index 000000000..b1617cfea --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/FirstSpanListener.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.agentstream.worker.segment; + +import org.skywalking.apm.network.proto.SpanObject; + +/** + * @author pengys5 + */ +public interface FirstSpanListener extends SpanListener { + void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId); +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/LocalSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/LocalSpanListener.java new file mode 100644 index 000000000..0486b3ffe --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/LocalSpanListener.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.agentstream.worker.segment; + +import org.skywalking.apm.network.proto.SpanObject; + +/** + * @author pengys5 + */ +public interface LocalSpanListener extends SpanListener { + void parseLocal(SpanObject spanObject, int applicationId, int applicationInstanceId); +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/RefsListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/RefsListener.java new file mode 100644 index 000000000..a931cef5a --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/RefsListener.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.agentstream.worker.segment; + +import org.skywalking.apm.network.proto.TraceSegmentReference; + +/** + * @author pengys5 + */ +public interface RefsListener extends SpanListener { + void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId); +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java index 09a90bfd9..2e6168b8e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java @@ -1,8 +1,107 @@ package org.skywalking.apm.collector.agentstream.worker.segment; +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentSpanListener; +import org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingSpanListener; +import org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefSpanListener; +import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.SpanType; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.skywalking.apm.network.proto.UniqueId; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + /** * @author pengys5 */ public class SegmentParse { + private final Logger logger = LoggerFactory.getLogger(SegmentParse.class); + + private List spanListeners; + private List refsListeners; + + public SegmentParse() { + spanListeners = new ArrayList<>(); + spanListeners.add(new NodeRefSpanListener()); + spanListeners.add(new NodeComponentSpanListener()); + spanListeners.add(new NodeMappingSpanListener()); + + refsListeners = new ArrayList<>(); + } + + public void parse(List traceIds, TraceSegmentObject segmentObject) { + for (UniqueId uniqueId : traceIds) { + uniqueId.getIdPartsList(); + } + + int applicationId = segmentObject.getApplicationId(); + int applicationInstanceId = segmentObject.getApplicationInstanceId(); + + for (TraceSegmentReference reference : segmentObject.getRefsList()) { + notifyRefsListener(reference, applicationId, applicationInstanceId); + } + + List spans = segmentObject.getSpansList(); + if (CollectionUtils.isNotEmpty(spans)) { + for (SpanObject spanObject : spans) { + if (spanObject.getSpanId() == 0) { + notifyFirstListener(spanObject, applicationId, applicationInstanceId); + } + + if (SpanType.Exit.equals(spanObject.getSpanType())) { + notifyExitListener(spanObject, applicationId, applicationInstanceId); + } else if (SpanType.Entry.equals(spanObject.getSpanType())) { + notifyEntryListener(spanObject, applicationId, applicationInstanceId); + } else if (SpanType.Local.equals(spanObject.getSpanType())) { + notifyLocalListener(spanObject, applicationId, applicationInstanceId); + } else { + logger.error("span type error, span type: {}", spanObject.getSpanType().name()); + } + } + } + } + + private void notifyExitListener(SpanObject spanObject, int applicationId, int applicationInstanceId) { + for (SpanListener listener : spanListeners) { + if (listener instanceof ExitSpanListener) { + ((ExitSpanListener)listener).parseExit(spanObject, applicationId, applicationInstanceId); + } + } + } + + private void notifyEntryListener(SpanObject spanObject, int applicationId, int applicationInstanceId) { + for (SpanListener listener : spanListeners) { + if (listener instanceof EntrySpanListener) { + ((EntrySpanListener)listener).parseEntry(spanObject, applicationId, applicationInstanceId); + } + } + } + + private void notifyLocalListener(SpanObject spanObject, int applicationId, int applicationInstanceId) { + for (SpanListener listener : spanListeners) { + if (listener instanceof LocalSpanListener) { + ((LocalSpanListener)listener).parseLocal(spanObject, applicationId, applicationInstanceId); + } + } + } + + private void notifyFirstListener(SpanObject spanObject, int applicationId, int applicationInstanceId) { + for (SpanListener listener : spanListeners) { + if (listener instanceof FirstSpanListener) { + ((FirstSpanListener)listener).parseFirst(spanObject, applicationId, applicationInstanceId); + } + } + } + + private void notifyRefsListener(TraceSegmentReference reference, int applicationId, int applicationInstanceId) { + for (SpanListener listener : refsListeners) { + if (listener instanceof RefsListener) { + ((RefsListener)listener).parseRef(reference, applicationId, applicationInstanceId); + } + } + } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SpanListener.java new file mode 100644 index 000000000..3dde04717 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SpanListener.java @@ -0,0 +1,8 @@ +package org.skywalking.apm.collector.agentstream.worker.segment; + +/** + * @author pengys5 + */ +public interface SpanListener { + void build(); +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/ServiceRefSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/ServiceRefSpanListener.java new file mode 100644 index 000000000..9c5cc6f87 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/ServiceRefSpanListener.java @@ -0,0 +1,79 @@ +package org.skywalking.apm.collector.agentstream.worker.serviceref.reference; + +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; +import org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefDataDefine; +import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class ServiceRefSpanListener implements FirstSpanListener, EntrySpanListener, ExitSpanListener, RefsListener { + + private final Logger logger = LoggerFactory.getLogger(ServiceRefSpanListener.class); + + private String front; + private List behinds = new ArrayList<>(); + private List fronts = new ArrayList<>(); + private long timeBucket; + + @Override public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId) { + timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); + } + + @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId) { + String entryService = String.valueOf(reference.getEntryServiceId()); + if (reference.getEntryServiceId() == 0) { + entryService = reference.getEntryServiceName(); + } + String parentService = String.valueOf(reference.getParentServiceId()); + if (reference.getParentServiceId() == 0) { + parentService = reference.getParentServiceName(); + } + fronts.add(new ServiceTemp(entryService, parentService)); + } + + @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId) { + front = String.valueOf(spanObject.getOperationNameId()); + if (spanObject.getOperationNameId() == 0) { + front = spanObject.getOperationName(); + } + } + + @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId) { + String behind = String.valueOf(spanObject.getOperationNameId()); + if (spanObject.getOperationNameId() == 0) { + behind = spanObject.getOperationName(); + } + behinds.add(behind); + } + + @Override public void build() { + for (String behind : behinds) { + String agg = front + Const.ID_SPLIT + behind; + ServiceRefDataDefine.ServiceReference serviceReference = new ServiceRefDataDefine.ServiceReference(); + serviceReference.setId(timeBucket + Const.ID_SPLIT + agg); + serviceReference.setAgg(agg); + serviceReference.setTimeBucket(timeBucket); + } + } + + class ServiceTemp { + private final String entryService; + private final String parentService; + + public ServiceTemp(String entryService, String parentService) { + this.entryService = entryService; + this.parentService = parentService; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefDataDefine.java new file mode 100644 index 000000000..00ff6b140 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefDataDefine.java @@ -0,0 +1,98 @@ +package org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define; + +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; +import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; + +/** + * @author pengys5 + */ +public class ServiceRefDataDefine extends DataDefine { + + public static final int DEFINE_ID = 501; + + @Override public int defineId() { + return DEFINE_ID; + } + + @Override protected int initialCapacity() { + return 4; + } + + @Override protected void attributeDefine() { + addAttribute(0, new Attribute(ServiceRefTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); + addAttribute(1, new Attribute(ServiceRefTable.COLUMN_ENTRY_SERVICE, AttributeType.STRING, new NonOperation())); + addAttribute(2, new Attribute(ServiceRefTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); + addAttribute(3, new Attribute(ServiceRefTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + } + + @Override public Object deserialize(RemoteData remoteData) { + String id = remoteData.getDataStrings(0); + String entryService = remoteData.getDataStrings(1); + String agg = remoteData.getDataStrings(2); + long timeBucket = remoteData.getDataLongs(0); + return new ServiceReference(id, entryService, agg, timeBucket); + } + + @Override public RemoteData serialize(Object object) { + ServiceReference serviceReference = (ServiceReference)object; + RemoteData.Builder builder = RemoteData.newBuilder(); + builder.addDataStrings(serviceReference.getId()); + builder.addDataStrings(serviceReference.getEntryService()); + builder.addDataStrings(serviceReference.getAgg()); + builder.addDataLongs(serviceReference.getTimeBucket()); + return builder.build(); + } + + public static class ServiceReference { + private String id; + private String entryService; + private String agg; + private long timeBucket; + + public ServiceReference(String id, String entryService, String agg, long timeBucket) { + this.id = id; + this.entryService = entryService; + this.agg = agg; + this.timeBucket = timeBucket; + } + + public ServiceReference() { + } + + public String getId() { + return id; + } + + public String getAgg() { + return agg; + } + + public long getTimeBucket() { + return timeBucket; + } + + public void setId(String id) { + this.id = id; + } + + public void setAgg(String agg) { + this.agg = agg; + } + + public void setTimeBucket(long timeBucket) { + this.timeBucket = timeBucket; + } + + public String getEntryService() { + return entryService; + } + + public void setEntryService(String entryService) { + this.entryService = entryService; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefEsTableDefine.java new file mode 100644 index 000000000..067624ed5 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefEsTableDefine.java @@ -0,0 +1,32 @@ +package org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define; + +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; + +/** + * @author pengys5 + */ +public class ServiceRefEsTableDefine extends ElasticSearchTableDefine { + + public ServiceRefEsTableDefine() { + super(ServiceRefTable.TABLE); + } + + @Override public int refreshInterval() { + return 0; + } + + @Override public int numberOfShards() { + return 2; + } + + @Override public int numberOfReplicas() { + return 0; + } + + @Override public void initialize() { + addColumn(new ElasticSearchColumnDefine(ServiceRefTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(ServiceRefTable.COLUMN_ENTRY_SERVICE, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(ServiceRefTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefH2TableDefine.java new file mode 100644 index 000000000..c6756cab5 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefH2TableDefine.java @@ -0,0 +1,20 @@ +package org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define; + +import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; +import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; + +/** + * @author pengys5 + */ +public class ServiceRefH2TableDefine extends H2TableDefine { + + public ServiceRefH2TableDefine() { + super(ServiceRefTable.TABLE); + } + + @Override public void initialize() { + addColumn(new H2ColumnDefine(ServiceRefTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(ServiceRefTable.COLUMN_ENTRY_SERVICE, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(ServiceRefTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefTable.java new file mode 100644 index 000000000..d145b93f5 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefTable.java @@ -0,0 +1,11 @@ +package org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define; + +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + +/** + * @author pengys5 + */ +public class ServiceRefTable extends CommonTable { + public static final String TABLE = "service_reference"; + public static final String COLUMN_ENTRY_SERVICE = "entry_service"; +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/util/TimeBucketUtils.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/util/TimeBucketUtils.java new file mode 100644 index 000000000..8f8b0a114 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/util/TimeBucketUtils.java @@ -0,0 +1,47 @@ +package org.skywalking.apm.collector.agentstream.worker.util; + +import java.text.SimpleDateFormat; +import java.util.Calendar; +import java.util.TimeZone; + +/** + * @author pengys5 + */ +public enum TimeBucketUtils { + INSTANCE; + + private final SimpleDateFormat DAY_DATE_FORMAT = new SimpleDateFormat("yyyyMMdd"); + private final SimpleDateFormat HOUR_DATE_FORMAT = new SimpleDateFormat("yyyyMMddHH"); + private final SimpleDateFormat MINUTE_DATE_FORMAT = new SimpleDateFormat("yyyyMMddHHmm"); + + public long getMinuteTimeBucket(long time) { + Calendar calendar = Calendar.getInstance(); + calendar.setTimeInMillis(time); + String timeStr = MINUTE_DATE_FORMAT.format(calendar.getTime()); + return Long.valueOf(timeStr); + } + + public long getHourTimeBucket(long time) { + Calendar calendar = Calendar.getInstance(); + calendar.setTimeInMillis(time); + String timeStr = HOUR_DATE_FORMAT.format(calendar.getTime()) + "00"; + return Long.valueOf(timeStr); + } + + public long getDayTimeBucket(long time) { + Calendar calendar = Calendar.getInstance(); + calendar.setTimeInMillis(time); + String timeStr = DAY_DATE_FORMAT.format(calendar.getTime()) + "0000"; + return Long.valueOf(timeStr); + } + + public long changeToUTCTimeBucket(long timeBucket) { + String timeBucketStr = String.valueOf(timeBucket); + + if (TimeZone.getDefault().getID().equals("GMT+08:00") || timeBucketStr.endsWith("0000")) { + return timeBucket; + } else { + return timeBucket - 800; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define index 205e6c74a..c8abe9d72 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/data.define @@ -1,4 +1,4 @@ -org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentDataDefine +org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationDataDefine org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceDataDefine org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameDataDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define index 0d0a2264d..d5017553e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define @@ -1,7 +1,7 @@ -org.skywalking.apm.collector.agentstream.worker.node.define.NodeComponentEsTableDefine -org.skywalking.apm.collector.agentstream.worker.node.define.NodeComponentH2TableDefine -org.skywalking.apm.collector.agentstream.worker.node.define.NodeMappingEsTableDefine -org.skywalking.apm.collector.agentstream.worker.node.define.NodeMappingH2TableDefine +org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentEsTableDefine +org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentH2TableDefine +org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingEsTableDefine +org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingH2TableDefine org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationEsTableDefine org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationH2TableDefine diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/Const.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/Const.java deleted file mode 100644 index c3acf882b..000000000 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/Const.java +++ /dev/null @@ -1,13 +0,0 @@ -package org.skywalking.apm.collector.stream.worker.impl; - -/** - * @author pengys5 - */ -public class Const { - public static final String ID_SPLIT = "..-.."; - public static final String IDS_SPLIT = "\\.\\.-\\.\\."; - public static final String PEERS_FRONT_SPLIT = "["; - public static final String PEERS_BEHIND_SPLIT = "]"; - public static final String USER_CODE = "User"; - public static final String RESULT = "result"; -}