From 2db52c1fca53630c00a2ce174c42da6f4456821f Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Fri, 25 Aug 2017 16:05:56 +0800 Subject: [PATCH] Trace dag web ui test success. #354 --- .../component/NodeComponentSpanListener.java | 73 ++++--- .../component/dao/NodeComponentEsDAO.java | 19 +- .../define/NodeComponentEsTableDefine.java | 7 +- .../define/NodeComponentH2TableDefine.java | 7 +- .../node/mapping/NodeMappingSpanListener.java | 33 +-- .../node/mapping/dao/NodeMappingEsDAO.java | 18 +- .../define/NodeMappingEsTableDefine.java | 6 +- .../define/NodeMappingH2TableDefine.java | 6 +- .../define/NodeRefEsTableDefine.java | 12 +- .../define/NodeRefH2TableDefine.java | 12 +- .../summary/NodeRefSumSpanListener.java | 77 +++---- .../noderef/summary/dao/NodeRefSumEsDAO.java | 61 +++--- .../define/NodeRefSumEsTableDefine.java | 18 +- .../define/NodeRefSumH2TableDefine.java | 18 +- .../worker/segment/SegmentParse.java | 4 +- .../agentstream/mock/SegmentPost.java | 4 +- .../apm/collector/core/stream/Data.java | 25 +++ .../apm/collector/core/util/Const.java | 4 +- .../define/node/NodeComponentDataDefine.java | 80 ++++++-- .../define/node/NodeComponentTable.java | 4 + .../define/node/NodeMappingDataDefine.java | 66 ++++-- .../storage/define/node/NodeMappingTable.java | 3 + .../define/noderef/NodeRefDataDefine.java | 8 +- .../define/noderef/NodeRefSumDataDefine.java | 189 ++++++++++-------- .../define/noderef/NodeRefSumTable.java | 13 +- .../storage/define/noderef/NodeRefTable.java | 9 + .../stream/worker/impl/AggregationWorker.java | 2 + .../collector/ui/dao/NodeComponentEsDAO.java | 93 +++++++-- .../collector/ui/dao/NodeMappingEsDAO.java | 47 +++-- .../apm/collector/ui/dao/NodeRefSumEsDAO.java | 108 ++++++---- .../collector/ui/dao/NodeReferenceEsDAO.java | 15 +- .../ui/jetty/handler/TraceDagGetHandler.java | 25 ++- .../ui/service/TraceDagDataBuilder.java | 61 ++---- .../collector/ui/service/TraceDagService.java | 10 +- 34 files changed, 722 insertions(+), 415 deletions(-) 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 index 7866954f4..6423f460e 100644 --- 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 @@ -2,15 +2,13 @@ package org.skywalking.apm.collector.agentstream.worker.node.component; import java.util.ArrayList; import java.util.List; -import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.storage.define.node.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.segment.LocalSpanListener; -import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; -import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.collector.storage.define.node.NodeComponentDataDefine; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; @@ -22,48 +20,59 @@ import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanListener, FirstSpanListener, LocalSpanListener { +public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanListener, FirstSpanListener { private final Logger logger = LoggerFactory.getLogger(NodeComponentSpanListener.class); - private List nodeComponents = new ArrayList<>(); + private List nodeComponents = new ArrayList<>(); private long timeBucket; @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String componentName = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getComponentId()); + NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); + nodeComponent.setComponentId(spanObject.getComponentId()); + + String id; if (spanObject.getComponentId() == 0) { - componentName = spanObject.getComponent(); - } - String peer = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getPeerId()); - if (spanObject.getPeerId() == 0) { - peer = spanObject.getPeer(); + nodeComponent.setComponentName(spanObject.getComponent()); + id = nodeComponent.getComponentName(); + } else { + nodeComponent.setComponentName(Const.EMPTY_STRING); + id = String.valueOf(nodeComponent.getComponentId()); } - String agg = peer + Const.ID_SPLIT + componentName; - nodeComponents.add(agg); + nodeComponent.setPeerId(spanObject.getPeerId()); + if (spanObject.getPeerId() == 0) { + nodeComponent.setPeer(spanObject.getPeer()); + id = id + Const.ID_SPLIT + nodeComponent.getPeer(); + } else { + nodeComponent.setPeer(Const.EMPTY_STRING); + id = id + Const.ID_SPLIT + nodeComponent.getPeerId(); + } + nodeComponent.setId(id); + nodeComponents.add(nodeComponent); } @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - buildEntryOrLocal(spanObject, applicationId); - } - - @Override - public void parseLocal(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - buildEntryOrLocal(spanObject, applicationId); - } - - private void buildEntryOrLocal(SpanObject spanObject, int applicationId) { - String componentName = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getComponentId()); + NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); + nodeComponent.setComponentId(spanObject.getComponentId()); + String id; if (spanObject.getComponentId() == 0) { - componentName = spanObject.getComponent(); + nodeComponent.setComponentName(spanObject.getComponent()); + id = nodeComponent.getComponentName(); + } else { + id = String.valueOf(nodeComponent.getComponentId()); + nodeComponent.setComponentName(Const.EMPTY_STRING); } - String peer = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); - String agg = peer + Const.ID_SPLIT + componentName; - nodeComponents.add(agg); + nodeComponent.setPeerId(applicationId); + nodeComponent.setPeer(Const.EMPTY_STRING); + id = id + Const.ID_SPLIT + String.valueOf(applicationId); + nodeComponent.setId(id); + + nodeComponents.add(nodeComponent); } @Override @@ -74,10 +83,8 @@ public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanLis @Override public void build() { StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); - nodeComponents.forEach(agg -> { - NodeComponentDataDefine.NodeComponent nodeComponent = new NodeComponentDataDefine.NodeComponent(); - nodeComponent.setId(timeBucket + Const.ID_SPLIT + agg); - nodeComponent.setAgg(agg); + nodeComponents.forEach(nodeComponent -> { + nodeComponent.setId(timeBucket + Const.ID_SPLIT + nodeComponent.getId()); nodeComponent.setTimeBucket(timeBucket); try { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java index cd53b5537..25cf19d4b 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentEsDAO.java @@ -5,11 +5,11 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.skywalking.apm.collector.core.stream.Data; +import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; -import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -21,7 +21,10 @@ public class NodeComponentEsDAO extends EsDAO implements INodeComponentDAO, IPer if (getResponse.isExists()) { Data data = dataDefine.build(id); Map source = getResponse.getSource(); - data.setDataString(1, (String)source.get(NodeComponentTable.COLUMN_AGG)); + data.setDataInteger(0, ((Number)source.get(NodeComponentTable.COLUMN_COMPONENT_ID)).intValue()); + data.setDataString(1, (String)source.get(NodeComponentTable.COLUMN_COMPONENT_NAME)); + data.setDataInteger(1, ((Number)source.get(NodeComponentTable.COLUMN_PEER_ID)).intValue()); + data.setDataString(2, (String)source.get(NodeComponentTable.COLUMN_PEER)); data.setDataLong(0, (Long)source.get(NodeComponentTable.COLUMN_TIME_BUCKET)); return data; } else { @@ -31,7 +34,10 @@ public class NodeComponentEsDAO extends EsDAO implements INodeComponentDAO, IPer @Override public IndexRequestBuilder prepareBatchInsert(Data data) { Map source = new HashMap<>(); - source.put(NodeComponentTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_COMPONENT_ID, data.getDataInteger(0)); + source.put(NodeComponentTable.COLUMN_COMPONENT_NAME, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_PEER_ID, data.getDataInteger(1)); + source.put(NodeComponentTable.COLUMN_PEER, data.getDataString(2)); source.put(NodeComponentTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); return getClient().prepareIndex(NodeComponentTable.TABLE, data.getDataString(0)).setSource(source); @@ -39,7 +45,10 @@ public class NodeComponentEsDAO extends EsDAO implements INodeComponentDAO, IPer @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { Map source = new HashMap<>(); - source.put(NodeComponentTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_COMPONENT_ID, data.getDataInteger(0)); + source.put(NodeComponentTable.COLUMN_COMPONENT_NAME, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_PEER_ID, data.getDataInteger(1)); + source.put(NodeComponentTable.COLUMN_PEER, data.getDataString(2)); source.put(NodeComponentTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); return getClient().prepareUpdate(NodeComponentTable.TABLE, data.getDataString(0)).setDoc(source); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java index ebcc93646..1ded3d4fe 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.node.component.define; +import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; -import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; /** * @author pengys5 @@ -26,7 +26,10 @@ public class NodeComponentEsTableDefine extends ElasticSearchTableDefine { } @Override public void initialize() { - addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_COMPONENT_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_COMPONENT_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_PEER_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_PEER, 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/component/define/NodeComponentH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java index ea6bf0963..c38535ea2 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentH2TableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.node.component.define; +import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; -import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; /** * @author pengys5 @@ -15,7 +15,10 @@ public class NodeComponentH2TableDefine extends H2TableDefine { @Override public void initialize() { 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_COMPONENT_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_COMPONENT_NAME, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_PEER_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeComponentTable.COLUMN_PEER, 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/mapping/NodeMappingSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java index 608c0d267..c7d3a91aa 100644 --- 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 @@ -2,13 +2,12 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping; import java.util.ArrayList; import java.util.List; -import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.storage.define.node.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.ExchangeMarkUtils; -import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.collector.storage.define.node.NodeMappingDataDefine; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; @@ -25,19 +24,27 @@ public class NodeMappingSpanListener implements RefsListener, FirstSpanListener private final Logger logger = LoggerFactory.getLogger(NodeMappingSpanListener.class); - private List nodeMappings = new ArrayList<>(); + private List nodeMappings = new ArrayList<>(); private long timeBucket; @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId, String segmentId) { logger.debug("node mapping listener parse reference"); - String peers = reference.getNetworkAddress(); + NodeMappingDataDefine.NodeMapping nodeMapping = new NodeMappingDataDefine.NodeMapping(); + nodeMapping.setApplicationId(applicationId); + nodeMapping.setAddressId(reference.getNetworkAddressId()); + + String id = String.valueOf(applicationId); if (reference.getNetworkAddressId() != 0) { - peers = ExchangeMarkUtils.INSTANCE.buildMarkedID(reference.getNetworkAddressId()); + nodeMapping.setAddress(Const.EMPTY_STRING); + id = id + Const.ID_SPLIT + String.valueOf(nodeMapping.getAddressId()); + } else { + id = id + Const.ID_SPLIT + reference.getNetworkAddress(); + nodeMapping.setAddress(reference.getNetworkAddress()); } - String agg = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId) + Const.ID_SPLIT + peers; - nodeMappings.add(agg); + nodeMapping.setId(id); + nodeMappings.add(nodeMapping); } @Override @@ -48,13 +55,11 @@ public class NodeMappingSpanListener implements RefsListener, FirstSpanListener @Override public void build() { logger.debug("node mapping listener build"); StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); - for (String agg : nodeMappings) { - NodeMappingDataDefine.NodeMapping nodeMapping = new NodeMappingDataDefine.NodeMapping(); - nodeMapping.setId(timeBucket + Const.ID_SPLIT + agg); - nodeMapping.setAgg(agg); - nodeMapping.setTimeBucket(timeBucket); + for (NodeMappingDataDefine.NodeMapping nodeMapping : nodeMappings) { try { + nodeMapping.setId(timeBucket + Const.ID_SPLIT + nodeMapping.getId()); + nodeMapping.setTimeBucket(timeBucket); logger.debug("send to node mapping aggregation worker, id: {}", nodeMapping.getId()); context.getClusterWorkerContext().lookup(NodeMappingAggregationWorker.WorkerRole.INSTANCE).tell(nodeMapping.toData()); } catch (WorkerInvokeException | WorkerNotFoundException e) { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java index 907be8c8e..1afd946de 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingEsDAO.java @@ -5,11 +5,11 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.skywalking.apm.collector.core.stream.Data; +import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; -import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -21,8 +21,10 @@ public class NodeMappingEsDAO extends EsDAO implements INodeMappingDAO, IPersist if (getResponse.isExists()) { Data data = dataDefine.build(id); Map source = getResponse.getSource(); - data.setDataString(1, (String)source.get(NodeMappingTable.COLUMN_AGG)); - data.setDataLong(0, (Long)source.get(NodeMappingTable.COLUMN_TIME_BUCKET)); + data.setDataInteger(0, ((Number)source.get(NodeMappingTable.COLUMN_APPLICATION_ID)).intValue()); + data.setDataInteger(1, ((Number)source.get(NodeMappingTable.COLUMN_ADDRESS_ID)).intValue()); + data.setDataString(1, (String)source.get(NodeMappingTable.COLUMN_ADDRESS)); + data.setDataLong(0, ((Number)source.get(NodeMappingTable.COLUMN_TIME_BUCKET)).longValue()); return data; } else { return null; @@ -31,7 +33,9 @@ public class NodeMappingEsDAO extends EsDAO implements INodeMappingDAO, IPersist @Override public IndexRequestBuilder prepareBatchInsert(Data data) { Map source = new HashMap<>(); - source.put(NodeMappingTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeMappingTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(NodeMappingTable.COLUMN_ADDRESS_ID, data.getDataInteger(1)); + source.put(NodeMappingTable.COLUMN_ADDRESS, data.getDataString(1)); source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); return getClient().prepareIndex(NodeMappingTable.TABLE, data.getDataString(0)).setSource(source); @@ -39,7 +43,9 @@ public class NodeMappingEsDAO extends EsDAO implements INodeMappingDAO, IPersist @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { Map source = new HashMap<>(); - source.put(NodeMappingTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeMappingTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(NodeMappingTable.COLUMN_ADDRESS_ID, data.getDataInteger(1)); + source.put(NodeMappingTable.COLUMN_ADDRESS, data.getDataString(1)); source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); return getClient().prepareUpdate(NodeMappingTable.TABLE, data.getDataString(0)).setDoc(source); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java index 978f2e1eb..e1027b392 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.define; +import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; -import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; /** * @author pengys5 @@ -26,7 +26,9 @@ public class NodeMappingEsTableDefine extends ElasticSearchTableDefine { } @Override public void initialize() { - addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_ADDRESS_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_ADDRESS, 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/mapping/define/NodeMappingH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingH2TableDefine.java index ddd7a5b65..721b323b1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingH2TableDefine.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.define; +import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; -import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; /** * @author pengys5 @@ -15,7 +15,9 @@ public class NodeMappingH2TableDefine extends H2TableDefine { @Override public void initialize() { 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_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_ADDRESS_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeMappingTable.COLUMN_ADDRESS, 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/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 index 5cb48064b..04a269c7c 100644 --- 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 @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define; +import org.skywalking.apm.collector.storage.define.noderef.NodeRefTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; -import org.skywalking.apm.collector.storage.define.noderef.NodeRefTable; /** * @author pengys5 @@ -26,7 +26,15 @@ public class NodeRefEsTableDefine extends ElasticSearchTableDefine { } @Override public void initialize() { - addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_FRONT_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_BEHIND_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_BEHIND_PEER, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_S1_LTE, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_S3_LTE, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_S5_LTE, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_S5_GT, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_SUMMARY, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefTable.COLUMN_ERROR, ElasticSearchColumnDefine.Type.Integer.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 index 5a9173a1f..6a23f1e8c 100644 --- 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 @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define; +import org.skywalking.apm.collector.storage.define.noderef.NodeRefTable; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; -import org.skywalking.apm.collector.storage.define.noderef.NodeRefTable; /** * @author pengys5 @@ -14,7 +14,15 @@ public class NodeRefH2TableDefine extends H2TableDefine { } @Override public void initialize() { - addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_FRONT_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_BEHIND_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_BEHIND_PEER, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_S1_LTE, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_S3_LTE, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_S5_LTE, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_S5_GT, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_SUMMARY, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefTable.COLUMN_ERROR, H2ColumnDefine.Type.Int.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/summary/NodeRefSumSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java index 4742fbbed..262d926e0 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java @@ -2,16 +2,15 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary; import java.util.ArrayList; import java.util.List; -import org.skywalking.apm.collector.core.util.Const; import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache; -import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumDataDefine; 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.ExchangeMarkUtils; -import org.skywalking.apm.collector.core.util.TimeBucketUtils; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumDataDefine; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; @@ -30,7 +29,7 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen private List nodeExitReferences = new ArrayList<>(); private List nodeEntryReferences = new ArrayList<>(); - private List nodeReferences = new ArrayList<>(); + private List nodeReferences = new ArrayList<>(); private long timeBucket; private boolean hasReference = false; private long startTime; @@ -39,42 +38,49 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); - String behind = spanObject.getPeer(); - if (spanObject.getPeerId() != 0) { - behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getPeerId()); - } + NodeRefSumDataDefine.NodeReferenceSum referenceSum = new NodeRefSumDataDefine.NodeReferenceSum(); + referenceSum.setApplicationId(applicationId); + referenceSum.setBehindApplicationId(spanObject.getPeerId()); - String agg = front + Const.ID_SPLIT + behind; - nodeExitReferences.add(buildNodeRefSum(spanObject.getStartTime(), spanObject.getEndTime(), agg, spanObject.getIsError())); + String id = String.valueOf(applicationId); + if (spanObject.getPeerId() != 0) { + referenceSum.setBehindPeer(Const.EMPTY_STRING); + id = id + Const.ID_SPLIT + String.valueOf(spanObject.getPeerId()); + } else { + referenceSum.setBehindPeer(spanObject.getPeer()); + id = id + Const.ID_SPLIT + spanObject.getPeer(); + } + referenceSum.setId(id); + nodeExitReferences.add(buildNodeRefSum(referenceSum, spanObject.getStartTime(), spanObject.getEndTime(), spanObject.getIsError())); } @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); - String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(Const.USER_ID); - String agg = front + Const.ID_SPLIT + behind; - nodeEntryReferences.add(buildNodeRefSum(spanObject.getStartTime(), spanObject.getEndTime(), agg, spanObject.getIsError())); + NodeRefSumDataDefine.NodeReferenceSum referenceSum = new NodeRefSumDataDefine.NodeReferenceSum(); + referenceSum.setApplicationId(Const.USER_ID); + referenceSum.setBehindApplicationId(applicationId); + referenceSum.setBehindPeer(Const.EMPTY_STRING); + + String id = String.valueOf(Const.USER_ID) + Const.ID_SPLIT + String.valueOf(applicationId); + referenceSum.setId(id); + nodeEntryReferences.add(buildNodeRefSum(referenceSum, spanObject.getStartTime(), spanObject.getEndTime(), spanObject.getIsError())); } - private NodeRefSumDataDefine.NodeReferenceSum buildNodeRefSum(long startTime, long endTime, String agg, - boolean isError) { - NodeRefSumDataDefine.NodeReferenceSum referenceSum = new NodeRefSumDataDefine.NodeReferenceSum(); - referenceSum.setAgg(agg); - + private NodeRefSumDataDefine.NodeReferenceSum buildNodeRefSum(NodeRefSumDataDefine.NodeReferenceSum referenceSum, + long startTime, long endTime, boolean isError) { long cost = endTime - startTime; if (cost <= 1000 && !isError) { - referenceSum.setOneSecondLess(1L); + referenceSum.setS1LTE(1); } else if (1000 < cost && cost <= 3000 && !isError) { - referenceSum.setThreeSecondLess(1L); + referenceSum.setS3LTE(1); } else if (3000 < cost && cost <= 5000 && !isError) { - referenceSum.setFiveSecondLess(1L); + referenceSum.setS5LTE(1); } else if (5000 < cost && !isError) { - referenceSum.setFiveSecondGreater(1L); + referenceSum.setS5GT(1); } else { - referenceSum.setError(1L); + referenceSum.setError(1); } - referenceSum.setSummary(1L); + referenceSum.setSummary(1); return referenceSum; } @@ -90,13 +96,16 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen String segmentId) { int parentApplicationId = InstanceCache.get(reference.getParentApplicationInstanceId()); - String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(parentApplicationId); - String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); + NodeRefSumDataDefine.NodeReferenceSum referenceSum = new NodeRefSumDataDefine.NodeReferenceSum(); + referenceSum.setApplicationId(parentApplicationId); + referenceSum.setBehindApplicationId(applicationId); + referenceSum.setBehindPeer(Const.EMPTY_STRING); - String agg = front + Const.ID_SPLIT + behind; + String id = String.valueOf(parentApplicationId) + Const.ID_SPLIT + String.valueOf(applicationId); + referenceSum.setId(id); hasReference = true; - nodeReferences.add(agg); + nodeReferences.add(referenceSum); } @Override public void build() { @@ -105,13 +114,13 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen if (!hasReference) { nodeExitReferences.addAll(nodeEntryReferences); } else { - nodeReferences.forEach(agg -> { - nodeExitReferences.add(buildNodeRefSum(startTime, endTime, agg, isError)); + nodeReferences.forEach(referenceSum -> { + nodeExitReferences.add(buildNodeRefSum(referenceSum, startTime, endTime, isError)); }); } for (NodeRefSumDataDefine.NodeReferenceSum referenceSum : nodeExitReferences) { - referenceSum.setId(timeBucket + Const.ID_SPLIT + referenceSum.getAgg()); + referenceSum.setId(timeBucket + Const.ID_SPLIT + referenceSum.getId()); referenceSum.setTimeBucket(timeBucket); try { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java index 5632b5332..e0f73c036 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java @@ -5,12 +5,11 @@ import java.util.Map; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.update.UpdateRequestBuilder; -import org.skywalking.apm.collector.storage.define.noderef.NodeRefTable; +import org.skywalking.apm.collector.core.stream.Data; +import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; -import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -22,14 +21,16 @@ public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO, IPersisten if (getResponse.isExists()) { Data data = dataDefine.build(id); Map source = getResponse.getSource(); - data.setDataLong(0, ((Number)source.get(NodeRefSumTable.COLUMN_ONE_SECOND_LESS)).longValue()); - data.setDataLong(1, ((Number)source.get(NodeRefSumTable.COLUMN_THREE_SECOND_LESS)).longValue()); - data.setDataLong(2, ((Number)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS)).longValue()); - data.setDataLong(3, ((Number)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER)).longValue()); - data.setDataLong(4, ((Number)source.get(NodeRefSumTable.COLUMN_ERROR)).longValue()); - data.setDataLong(5, ((Number)source.get(NodeRefSumTable.COLUMN_SUMMARY)).longValue()); - data.setDataLong(6, ((Number)source.get(NodeRefSumTable.COLUMN_TIME_BUCKET)).longValue()); - data.setDataString(1, (String)source.get(NodeRefSumTable.COLUMN_AGG)); + data.setDataInteger(0, ((Number)source.get(NodeRefSumTable.COLUMN_FRONT_APPLICATION_ID)).intValue()); + data.setDataInteger(1, ((Number)source.get(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID)).intValue()); + data.setDataString(1, (String)source.get(NodeRefSumTable.COLUMN_BEHIND_PEER)); + data.setDataInteger(2, ((Number)source.get(NodeRefSumTable.COLUMN_S1_LTE)).intValue()); + data.setDataInteger(3, ((Number)source.get(NodeRefSumTable.COLUMN_S3_LTE)).intValue()); + data.setDataInteger(4, ((Number)source.get(NodeRefSumTable.COLUMN_S5_LTE)).intValue()); + data.setDataInteger(5, ((Number)source.get(NodeRefSumTable.COLUMN_S5_GT)).intValue()); + data.setDataInteger(6, ((Number)source.get(NodeRefSumTable.COLUMN_SUMMARY)).intValue()); + data.setDataInteger(7, ((Number)source.get(NodeRefSumTable.COLUMN_ERROR)).intValue()); + data.setDataLong(0, ((Number)source.get(NodeRefSumTable.COLUMN_TIME_BUCKET)).longValue()); return data; } else { return null; @@ -38,29 +39,33 @@ public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO, IPersisten @Override public IndexRequestBuilder prepareBatchInsert(Data data) { Map source = new HashMap<>(); - source.put(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, data.getDataLong(0)); - source.put(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, data.getDataLong(1)); - source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, data.getDataLong(2)); - source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, data.getDataLong(3)); - source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataLong(4)); - source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataLong(5)); - source.put(NodeRefSumTable.COLUMN_AGG, data.getDataString(1)); - source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(6)); + source.put(NodeRefSumTable.COLUMN_FRONT_APPLICATION_ID, data.getDataInteger(0)); + source.put(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID, data.getDataInteger(1)); + source.put(NodeRefSumTable.COLUMN_BEHIND_PEER, data.getDataString(1)); + source.put(NodeRefSumTable.COLUMN_S1_LTE, data.getDataInteger(2)); + source.put(NodeRefSumTable.COLUMN_S3_LTE, data.getDataInteger(3)); + source.put(NodeRefSumTable.COLUMN_S5_LTE, data.getDataInteger(4)); + source.put(NodeRefSumTable.COLUMN_S5_GT, data.getDataInteger(5)); + source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataInteger(6)); + source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataInteger(7)); + source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); return getClient().prepareIndex(NodeRefSumTable.TABLE, data.getDataString(0)).setSource(source); } @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { Map source = new HashMap<>(); - source.put(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, data.getDataLong(0)); - source.put(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, data.getDataLong(1)); - source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, data.getDataLong(2)); - source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, data.getDataLong(3)); - source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataLong(4)); - source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataLong(5)); - source.put(NodeRefSumTable.COLUMN_AGG, data.getDataString(1)); - source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(6)); + source.put(NodeRefSumTable.COLUMN_FRONT_APPLICATION_ID, data.getDataInteger(0)); + source.put(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID, data.getDataInteger(1)); + source.put(NodeRefSumTable.COLUMN_BEHIND_PEER, data.getDataString(1)); + source.put(NodeRefSumTable.COLUMN_S1_LTE, data.getDataInteger(2)); + source.put(NodeRefSumTable.COLUMN_S3_LTE, data.getDataInteger(3)); + source.put(NodeRefSumTable.COLUMN_S5_LTE, data.getDataInteger(4)); + source.put(NodeRefSumTable.COLUMN_S5_GT, data.getDataInteger(5)); + source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataInteger(6)); + source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataInteger(7)); + source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); - return getClient().prepareUpdate(NodeRefTable.TABLE, data.getDataString(0)).setDoc(source); + return getClient().prepareUpdate(NodeRefSumTable.TABLE, data.getDataString(0)).setDoc(source); } } 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 index 91e6946ed..7d9685661 100644 --- 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 @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define; +import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; -import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumTable; /** * @author pengys5 @@ -26,13 +26,15 @@ public class NodeRefSumEsTableDefine extends ElasticSearchTableDefine { } @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_FRONT_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_BEHIND_PEER, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_S1_LTE, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_S3_LTE, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_S5_LTE, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_S5_GT, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_SUMMARY, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NodeRefSumTable.COLUMN_ERROR, ElasticSearchColumnDefine.Type.Integer.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 index 978df4992..1e8484a6f 100644 --- 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 @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define; +import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumTable; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; -import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumTable; /** * @author pengys5 @@ -14,13 +14,15 @@ public class NodeRefSumH2TableDefine extends H2TableDefine { } @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_FRONT_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_BEHIND_PEER, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_S1_LTE, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_S3_LTE, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_S5_LTE, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_S5_GT, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_SUMMARY, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NodeRefSumTable.COLUMN_ERROR, H2ColumnDefine.Type.Int.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/segment/SegmentParse.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java index 5564127a4..155433e7f 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 @@ -6,15 +6,14 @@ import org.skywalking.apm.collector.agentstream.worker.global.GlobalTraceSpanLis import org.skywalking.apm.collector.agentstream.worker.instance.performance.InstPerformanceSpanListener; 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.agentstream.worker.noderef.summary.NodeRefSumSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.cost.SegmentCostSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.origin.SegmentPersistenceWorker; -import org.skywalking.apm.collector.storage.define.segment.SegmentDataDefine; import org.skywalking.apm.collector.agentstream.worker.service.entry.ServiceEntrySpanListener; import org.skywalking.apm.collector.agentstream.worker.serviceref.reference.ServiceRefSpanListener; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.skywalking.apm.collector.storage.define.segment.SegmentDataDefine; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; @@ -38,7 +37,6 @@ public class SegmentParse { public SegmentParse() { spanListeners = new ArrayList<>(); - spanListeners.add(new NodeRefSpanListener()); spanListeners.add(new NodeComponentSpanListener()); spanListeners.add(new NodeMappingSpanListener()); spanListeners.add(new NodeRefSumSpanListener()); diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java index 506b0d0ee..f6b877a70 100644 --- a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java @@ -32,6 +32,8 @@ public class SegmentPost { ApplicationEsDAO applicationEsDAO = new ApplicationEsDAO(); applicationEsDAO.setClient(client); + ApplicationDataDefine.Application userApplication = new ApplicationDataDefine.Application("1", "User", 1); + applicationEsDAO.save(userApplication); ApplicationDataDefine.Application consumerApplication = new ApplicationDataDefine.Application("2", "dubbox-consumer", 2); applicationEsDAO.save(consumerApplication); ApplicationDataDefine.Application providerApplication = new ApplicationDataDefine.Application("3", "dubbox-provider", 3); @@ -46,7 +48,7 @@ public class SegmentPost { modifyTime(provider); HttpClientTools.INSTANCE.post("http://localhost:12800/segments", provider.toString()); - Thread.sleep(200); + Thread.sleep(1000); } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java index 6789e26e2..2a82c75d5 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/stream/Data.java @@ -86,4 +86,29 @@ public class Data extends AbstractHashMessage { public String id() { return dataStrings[0]; } + + @Override public String toString() { + StringBuilder dataStr = new StringBuilder(); + dataStr.append("string: ["); + for (int i = 0; i < dataStrings.length; i++) { + dataStr.append(dataStrings[i]).append(","); + } + dataStr.append("], longs: ["); + for (int i = 0; i < dataLongs.length; i++) { + dataStr.append(dataLongs[i]).append(","); + } + dataStr.append("], double: ["); + for (int i = 0; i < dataDoubles.length; i++) { + dataStr.append(dataDoubles[i]).append(","); + } + dataStr.append("], integer: ["); + for (int i = 0; i < dataIntegers.length; i++) { + dataStr.append(dataIntegers[i]).append(","); + } + dataStr.append("], boolean: ["); + for (int i = 0; i < dataBooleans.length; i++) { + dataStr.append(dataBooleans[i]).append(","); + } + return dataStr.toString(); + } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/Const.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/Const.java index 942ef9f3b..0d0ecf643 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/Const.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/Const.java @@ -4,11 +4,11 @@ package org.skywalking.apm.collector.core.util; * @author pengys5 */ public class Const { - public static final String ID_SPLIT = "..-.."; - public static final String IDS_SPLIT = "\\.\\.-\\.\\."; + public static final String ID_SPLIT = "_"; public static final int USER_ID = 1; public static final String USER_CODE = "User"; public static final String SEGMENT_SPAN_SPLIT = "S"; public static final String UNKNOWN = "Unknown"; public static final String EXCEPTION = "Exception"; + public static final String EMPTY_STRING = ""; } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentDataDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentDataDefine.java index 3cc9bd5fb..fad5a335a 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentDataDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentDataDefine.java @@ -1,13 +1,13 @@ package org.skywalking.apm.collector.storage.define.node; -import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; -import org.skywalking.apm.collector.storage.define.Attribute; -import org.skywalking.apm.collector.storage.define.AttributeType; import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.core.stream.Transform; import org.skywalking.apm.collector.core.stream.operate.CoverOperation; import org.skywalking.apm.collector.core.stream.operate.NonOperation; +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.storage.define.Attribute; +import org.skywalking.apm.collector.storage.define.AttributeType; +import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -15,13 +15,16 @@ import org.skywalking.apm.collector.core.stream.operate.NonOperation; public class NodeComponentDataDefine extends DataDefine { @Override protected int initialCapacity() { - return 3; + return 6; } @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())); + addAttribute(1, new Attribute(NodeComponentTable.COLUMN_COMPONENT_ID, AttributeType.INTEGER, new CoverOperation())); + addAttribute(2, new Attribute(NodeComponentTable.COLUMN_COMPONENT_NAME, AttributeType.STRING, new CoverOperation())); + addAttribute(3, new Attribute(NodeComponentTable.COLUMN_PEER_ID, AttributeType.INTEGER, new CoverOperation())); + addAttribute(4, new Attribute(NodeComponentTable.COLUMN_PEER, AttributeType.STRING, new CoverOperation())); + addAttribute(5, new Attribute(NodeComponentTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); } @Override public Object deserialize(RemoteData remoteData) { @@ -34,12 +37,19 @@ public class NodeComponentDataDefine extends DataDefine { public static class NodeComponent implements Transform { private String id; - private String agg; + private Integer componentId; + private String componentName; + private Integer peerId; + private String peer; private long timeBucket; - NodeComponent(String id, String agg, long timeBucket) { + public NodeComponent(String id, Integer componentId, String componentName, Integer peerId, String peer, + long timeBucket) { this.id = id; - this.agg = agg; + this.componentId = componentId; + this.componentName = componentName; + this.peerId = peerId; + this.peer = peer; this.timeBucket = timeBucket; } @@ -50,14 +60,20 @@ public class NodeComponentDataDefine extends DataDefine { NodeComponentDataDefine define = new NodeComponentDataDefine(); Data data = define.build(id); data.setDataString(0, this.id); - data.setDataString(1, this.agg); + data.setDataInteger(0, this.componentId); + data.setDataString(1, this.componentName); + data.setDataInteger(1, this.peerId); + data.setDataString(2, this.peer); data.setDataLong(0, this.timeBucket); return data; } @Override public NodeComponent toSelf(Data data) { this.id = data.getDataString(0); - this.agg = data.getDataString(1); + this.componentId = data.getDataInteger(0); + this.componentName = data.getDataString(1); + this.peerId = data.getDataInteger(1); + this.peer = data.getDataString(2); this.timeBucket = data.getDataLong(0); return this; } @@ -66,10 +82,6 @@ public class NodeComponentDataDefine extends DataDefine { return id; } - public String getAgg() { - return agg; - } - public long getTimeBucket() { return timeBucket; } @@ -78,12 +90,40 @@ public class NodeComponentDataDefine extends DataDefine { this.id = id; } - public void setAgg(String agg) { - this.agg = agg; - } - public void setTimeBucket(long timeBucket) { this.timeBucket = timeBucket; } + + public Integer getComponentId() { + return componentId; + } + + public void setComponentId(Integer componentId) { + this.componentId = componentId; + } + + public String getComponentName() { + return componentName; + } + + public void setComponentName(String componentName) { + this.componentName = componentName; + } + + public Integer getPeerId() { + return peerId; + } + + public void setPeerId(Integer peerId) { + this.peerId = peerId; + } + + public String getPeer() { + return peer; + } + + public void setPeer(String peer) { + this.peer = peer; + } } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentTable.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentTable.java index 0354e194b..e4ab07892 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentTable.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeComponentTable.java @@ -7,4 +7,8 @@ import org.skywalking.apm.collector.storage.define.CommonTable; */ public class NodeComponentTable extends CommonTable { public static final String TABLE = "node_component"; + public static final String COLUMN_COMPONENT_ID = "component_id"; + public static final String COLUMN_COMPONENT_NAME = "component_name"; + public static final String COLUMN_PEER = "peer"; + public static final String COLUMN_PEER_ID = "peer_id"; } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingDataDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingDataDefine.java index a978fe51c..aea8a7d40 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingDataDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingDataDefine.java @@ -1,13 +1,13 @@ package org.skywalking.apm.collector.storage.define.node; -import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; -import org.skywalking.apm.collector.storage.define.Attribute; -import org.skywalking.apm.collector.storage.define.AttributeType; import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.core.stream.Transform; import org.skywalking.apm.collector.core.stream.operate.CoverOperation; import org.skywalking.apm.collector.core.stream.operate.NonOperation; +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.storage.define.Attribute; +import org.skywalking.apm.collector.storage.define.AttributeType; +import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -15,13 +15,15 @@ import org.skywalking.apm.collector.core.stream.operate.NonOperation; public class NodeMappingDataDefine extends DataDefine { @Override protected int initialCapacity() { - return 3; + return 5; } @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())); + addAttribute(1, new Attribute(NodeMappingTable.COLUMN_APPLICATION_ID, AttributeType.INTEGER, new CoverOperation())); + addAttribute(2, new Attribute(NodeMappingTable.COLUMN_ADDRESS_ID, AttributeType.INTEGER, new CoverOperation())); + addAttribute(3, new Attribute(NodeMappingTable.COLUMN_ADDRESS, AttributeType.STRING, new CoverOperation())); + addAttribute(4, new Attribute(NodeMappingTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); } @Override public Object deserialize(RemoteData remoteData) { @@ -34,12 +36,16 @@ public class NodeMappingDataDefine extends DataDefine { public static class NodeMapping implements Transform { private String id; - private String agg; + private int applicationId; + private int addressId; + private String address; private long timeBucket; - NodeMapping(String id, String agg, long timeBucket) { + public NodeMapping(String id, int applicationId, int addressId, String address, long timeBucket) { this.id = id; - this.agg = agg; + this.applicationId = applicationId; + this.addressId = addressId; + this.address = address; this.timeBucket = timeBucket; } @@ -50,14 +56,18 @@ public class NodeMappingDataDefine extends DataDefine { NodeMappingDataDefine define = new NodeMappingDataDefine(); Data data = define.build(id); data.setDataString(0, this.id); - data.setDataString(1, this.agg); + data.setDataInteger(0, this.applicationId); + data.setDataInteger(1, this.addressId); + data.setDataString(1, this.address); data.setDataLong(0, this.timeBucket); return data; } @Override public NodeMapping toSelf(Data data) { this.id = data.getDataString(0); - this.agg = data.getDataString(1); + this.applicationId = data.getDataInteger(0); + this.addressId = data.getDataInteger(1); + this.address = data.getDataString(1); this.timeBucket = data.getDataLong(0); return this; } @@ -66,10 +76,6 @@ public class NodeMappingDataDefine extends DataDefine { return id; } - public String getAgg() { - return agg; - } - public long getTimeBucket() { return timeBucket; } @@ -78,12 +84,32 @@ public class NodeMappingDataDefine extends DataDefine { this.id = id; } - public void setAgg(String agg) { - this.agg = agg; - } - public void setTimeBucket(long timeBucket) { this.timeBucket = timeBucket; } + + public int getApplicationId() { + return applicationId; + } + + public void setApplicationId(int applicationId) { + this.applicationId = applicationId; + } + + public int getAddressId() { + return addressId; + } + + public void setAddressId(int addressId) { + this.addressId = addressId; + } + + public String getAddress() { + return address; + } + + public void setAddress(String address) { + this.address = address; + } } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingTable.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingTable.java index 53fbd98e6..3c775989f 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingTable.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/node/NodeMappingTable.java @@ -7,4 +7,7 @@ import org.skywalking.apm.collector.storage.define.CommonTable; */ public class NodeMappingTable extends CommonTable { public static final String TABLE = "node_mapping"; + public static final String COLUMN_APPLICATION_ID = "application_id"; + public static final String COLUMN_ADDRESS_ID = "address_id"; + public static final String COLUMN_ADDRESS = "address"; } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefDataDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefDataDefine.java index e55efe0c9..b0ba846a3 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefDataDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefDataDefine.java @@ -1,13 +1,13 @@ package org.skywalking.apm.collector.storage.define.noderef; -import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; -import org.skywalking.apm.collector.storage.define.Attribute; -import org.skywalking.apm.collector.storage.define.AttributeType; import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.core.stream.Transform; import org.skywalking.apm.collector.core.stream.operate.CoverOperation; import org.skywalking.apm.collector.core.stream.operate.NonOperation; +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.storage.define.Attribute; +import org.skywalking.apm.collector.storage.define.AttributeType; +import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumDataDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumDataDefine.java index 9190d724d..445f5c00d 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumDataDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumDataDefine.java @@ -1,13 +1,13 @@ package org.skywalking.apm.collector.storage.define.noderef; -import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; -import org.skywalking.apm.collector.storage.define.Attribute; -import org.skywalking.apm.collector.storage.define.AttributeType; import org.skywalking.apm.collector.core.stream.Data; -import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.core.stream.Transform; import org.skywalking.apm.collector.core.stream.operate.AddOperation; import org.skywalking.apm.collector.core.stream.operate.NonOperation; +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.storage.define.Attribute; +import org.skywalking.apm.collector.storage.define.AttributeType; +import org.skywalking.apm.collector.storage.define.DataDefine; /** * @author pengys5 @@ -15,70 +15,81 @@ import org.skywalking.apm.collector.core.stream.operate.NonOperation; public class NodeRefSumDataDefine extends DataDefine { @Override protected int initialCapacity() { - return 9; + return 11; } @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 AddOperation())); - addAttribute(2, new Attribute(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, AttributeType.LONG, new AddOperation())); - addAttribute(3, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, AttributeType.LONG, new AddOperation())); - addAttribute(4, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, AttributeType.LONG, new AddOperation())); - addAttribute(5, new Attribute(NodeRefSumTable.COLUMN_ERROR, AttributeType.LONG, new AddOperation())); - addAttribute(6, new Attribute(NodeRefSumTable.COLUMN_SUMMARY, AttributeType.LONG, new AddOperation())); - addAttribute(7, new Attribute(NodeRefSumTable.COLUMN_AGG, AttributeType.STRING, new NonOperation())); - addAttribute(8, new Attribute(NodeRefSumTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new NonOperation())); + addAttribute(1, new Attribute(NodeRefSumTable.COLUMN_FRONT_APPLICATION_ID, AttributeType.INTEGER, new NonOperation())); + addAttribute(2, new Attribute(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID, AttributeType.INTEGER, new NonOperation())); + addAttribute(3, new Attribute(NodeRefSumTable.COLUMN_BEHIND_PEER, AttributeType.STRING, new NonOperation())); + addAttribute(4, new Attribute(NodeRefSumTable.COLUMN_S1_LTE, AttributeType.INTEGER, new AddOperation())); + addAttribute(5, new Attribute(NodeRefSumTable.COLUMN_S3_LTE, AttributeType.INTEGER, new AddOperation())); + addAttribute(6, new Attribute(NodeRefSumTable.COLUMN_S5_LTE, AttributeType.INTEGER, new AddOperation())); + addAttribute(7, new Attribute(NodeRefSumTable.COLUMN_S5_GT, AttributeType.INTEGER, new AddOperation())); + addAttribute(8, new Attribute(NodeRefSumTable.COLUMN_SUMMARY, AttributeType.INTEGER, new AddOperation())); + addAttribute(9, new Attribute(NodeRefSumTable.COLUMN_ERROR, AttributeType.INTEGER, new AddOperation())); + addAttribute(10, new Attribute(NodeRefSumTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new NonOperation())); } @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); + int applicationId = remoteData.getDataIntegers(0); + int behindApplicationId = remoteData.getDataIntegers(1); + String behindPeer = remoteData.getDataStrings(1); + int s1LTE = remoteData.getDataIntegers(2); + int s3LTE = remoteData.getDataIntegers(3); + int s5LTE = remoteData.getDataIntegers(4); + int s5GT = remoteData.getDataIntegers(5); + int summary = remoteData.getDataIntegers(6); + int error = remoteData.getDataIntegers(7); + long timeBucket = remoteData.getDataLongs(0); + return new NodeReferenceSum(id, applicationId, behindApplicationId, behindPeer, s1LTE, s3LTE, s5LTE, s5GT, summary, error, 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.addDataIntegers(nodeReferenceSum.getApplicationId()); + builder.addDataIntegers(nodeReferenceSum.getBehindApplicationId()); + builder.addDataStrings(nodeReferenceSum.getBehindPeer()); + builder.addDataIntegers(nodeReferenceSum.getS1LTE()); + builder.addDataIntegers(nodeReferenceSum.getS3LTE()); + builder.addDataIntegers(nodeReferenceSum.getS5LTE()); + builder.addDataIntegers(nodeReferenceSum.getS5GT()); + builder.addDataIntegers(nodeReferenceSum.getSummary()); + builder.addDataIntegers(nodeReferenceSum.getError()); builder.addDataLongs(nodeReferenceSum.getTimeBucket()); return builder.build(); } public static class NodeReferenceSum implements Transform { private String id; - private Long oneSecondLess = 0L; - private Long threeSecondLess = 0L; - private Long fiveSecondLess = 0L; - private Long fiveSecondGreater = 0L; - private Long error = 0L; - private Long summary = 0L; - private String agg; + private int applicationId; + private int behindApplicationId; + private String behindPeer; + private int s1LTE = 0; + private int s3LTE = 0; + private int s5LTE = 0; + private int s5GT = 0; + private int summary = 0; + private int error = 0; private long timeBucket; - public NodeReferenceSum(String id, Long oneSecondLess, Long threeSecondLess, Long fiveSecondLess, - Long fiveSecondGreater, Long error, Long summary, String agg, long timeBucket) { + public NodeReferenceSum(String id, int applicationId, int behindApplicationId, String behindPeer, int s1LTE, + int s3LTE, + int s5LTE, int s5GT, int summary, int error, long timeBucket) { this.id = id; - this.oneSecondLess = oneSecondLess; - this.threeSecondLess = threeSecondLess; - this.fiveSecondLess = fiveSecondLess; - this.fiveSecondGreater = fiveSecondGreater; - this.error = error; + this.applicationId = applicationId; + this.behindApplicationId = behindApplicationId; + this.behindPeer = behindPeer; + this.s1LTE = s1LTE; + this.s3LTE = s3LTE; + this.s5LTE = s5LTE; + this.s5GT = s5GT; this.summary = summary; - this.agg = agg; + this.error = error; this.timeBucket = timeBucket; } @@ -89,14 +100,16 @@ public class NodeRefSumDataDefine extends DataDefine { NodeRefSumDataDefine define = new NodeRefSumDataDefine(); Data data = define.build(id); data.setDataString(0, this.id); - data.setDataString(1, this.agg); - data.setDataLong(0, this.oneSecondLess); - data.setDataLong(1, this.threeSecondLess); - data.setDataLong(2, this.fiveSecondLess); - data.setDataLong(3, this.fiveSecondGreater); - data.setDataLong(4, this.error); - data.setDataLong(5, this.summary); - data.setDataLong(6, this.timeBucket); + data.setDataInteger(0, this.applicationId); + data.setDataInteger(1, this.behindApplicationId); + data.setDataString(1, this.behindPeer); + data.setDataInteger(2, this.s1LTE); + data.setDataInteger(3, this.s3LTE); + data.setDataInteger(4, this.s5LTE); + data.setDataInteger(5, this.s5GT); + data.setDataInteger(6, this.summary); + data.setDataInteger(7, this.error); + data.setDataLong(0, this.timeBucket); return data; } @@ -112,62 +125,78 @@ public class NodeRefSumDataDefine extends DataDefine { this.id = id; } - public Long getOneSecondLess() { - return oneSecondLess; + public int getApplicationId() { + return applicationId; } - public void setOneSecondLess(Long oneSecondLess) { - this.oneSecondLess = oneSecondLess; + public void setApplicationId(int applicationId) { + this.applicationId = applicationId; } - public Long getThreeSecondLess() { - return threeSecondLess; + public int getBehindApplicationId() { + return behindApplicationId; } - public void setThreeSecondLess(Long threeSecondLess) { - this.threeSecondLess = threeSecondLess; + public void setBehindApplicationId(int behindApplicationId) { + this.behindApplicationId = behindApplicationId; } - public Long getFiveSecondLess() { - return fiveSecondLess; + public String getBehindPeer() { + return behindPeer; } - public void setFiveSecondLess(Long fiveSecondLess) { - this.fiveSecondLess = fiveSecondLess; + public void setBehindPeer(String behindPeer) { + this.behindPeer = behindPeer; } - public Long getFiveSecondGreater() { - return fiveSecondGreater; + public int getS1LTE() { + return s1LTE; } - public void setFiveSecondGreater(Long fiveSecondGreater) { - this.fiveSecondGreater = fiveSecondGreater; + public void setS1LTE(int s1LTE) { + this.s1LTE = s1LTE; } - public Long getError() { + public int getS3LTE() { + return s3LTE; + } + + public void setS3LTE(int s3LTE) { + this.s3LTE = s3LTE; + } + + public int getS5LTE() { + return s5LTE; + } + + public void setS5LTE(int s5LTE) { + this.s5LTE = s5LTE; + } + + public int getS5GT() { + return s5GT; + } + + public void setS5GT(int s5GT) { + this.s5GT = s5GT; + } + + public int getError() { return error; } - public void setError(Long error) { + public void setError(int error) { this.error = error; } - public Long getSummary() { + public int getSummary() { return summary; } - public void setSummary(Long summary) { + public void setSummary(int summary) { this.summary = summary; } - public String getAgg() { - return agg; - } - - public void setAgg(String agg) { - this.agg = agg; - } - public long getTimeBucket() { return timeBucket; } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumTable.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumTable.java index de9cedb1d..49fe4bc7b 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumTable.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefSumTable.java @@ -7,10 +7,13 @@ import org.skywalking.apm.collector.storage.define.CommonTable; */ 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_FRONT_APPLICATION_ID = "front_application_id"; + public static final String COLUMN_BEHIND_APPLICATION_ID = "behind_application_id"; + public static final String COLUMN_BEHIND_PEER = "behind_peer"; + public static final String COLUMN_S1_LTE = "s1_lte"; + public static final String COLUMN_S3_LTE = "s3_lte"; + public static final String COLUMN_S5_LTE = "s5_lte"; + public static final String COLUMN_S5_GT = "s5_gt"; public static final String COLUMN_SUMMARY = "summary"; + public static final String COLUMN_ERROR = "error"; } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefTable.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefTable.java index 873d55e0d..51e0302d9 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefTable.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/define/noderef/NodeRefTable.java @@ -7,4 +7,13 @@ import org.skywalking.apm.collector.storage.define.CommonTable; */ public class NodeRefTable extends CommonTable { public static final String TABLE = "node_reference"; + public static final String COLUMN_FRONT_APPLICATION_ID = "front_application_id"; + public static final String COLUMN_BEHIND_APPLICATION_ID = "behind_application_id"; + public static final String COLUMN_BEHIND_PEER = "behind_peer"; + public static final String COLUMN_S1_LTE = "s1_lte"; + public static final String COLUMN_S3_LTE = "s3_lte"; + public static final String COLUMN_S5_LTE = "s5_lte"; + public static final String COLUMN_S5_GT = "s5_gt"; + public static final String COLUMN_SUMMARY = "summary"; + public static final String COLUMN_ERROR = "error"; } diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java index 08498bc98..4cc4676d6 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/AggregationWorker.java @@ -60,11 +60,13 @@ public abstract class AggregationWorker extends AbstractLocalAsyncWorker { } dataCache.getLast().asMap().forEach((id, data) -> { try { + logger.debug(data.toString()); nextWorkRef(id).tell(data); } catch (WorkerNotFoundException | WorkerInvokeException e) { logger.error(e.getMessage(), e); } }); + dataCache.releaseLast(); } protected final void aggregate(Object message) { diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentEsDAO.java index 8d6b1f238..a35192f7d 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentEsDAO.java @@ -8,9 +8,11 @@ import org.elasticsearch.action.search.SearchType; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.bucket.terms.Terms; -import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.core.util.StringUtils; import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.ui.cache.ApplicationCache; +import org.skywalking.apm.network.trace.component.ComponentsDefine; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -22,30 +24,89 @@ public class NodeComponentEsDAO extends EsDAO implements INodeComponentDAO { private final Logger logger = LoggerFactory.getLogger(NodeComponentEsDAO.class); @Override public JsonArray load(long startTime, long endTime) { + logger.debug("node component load, start time: {}, end time: {}", startTime, endTime); + JsonArray nodeComponentArray = new JsonArray(); + nodeComponentArray.addAll(aggregationByComponentId(startTime, endTime)); + nodeComponentArray.addAll(aggregationByComponentName(startTime, endTime)); + return nodeComponentArray; + } + + private JsonArray aggregationByComponentId(long startTime, long endTime) { SearchRequestBuilder searchRequestBuilder = getClient().prepareSearch(NodeComponentTable.TABLE); searchRequestBuilder.setTypes(NodeComponentTable.TABLE_TYPE); searchRequestBuilder.setSearchType(SearchType.DFS_QUERY_THEN_FETCH); searchRequestBuilder.setQuery(QueryBuilders.rangeQuery(NodeComponentTable.COLUMN_TIME_BUCKET).gte(startTime).lte(endTime)); searchRequestBuilder.setSize(0); - searchRequestBuilder.addAggregation(AggregationBuilders.terms(NodeComponentTable.COLUMN_AGG).field(NodeComponentTable.COLUMN_AGG).size(100)); + searchRequestBuilder.addAggregation(AggregationBuilders.terms(NodeComponentTable.COLUMN_COMPONENT_ID).field(NodeComponentTable.COLUMN_COMPONENT_ID).size(100) + .subAggregation(AggregationBuilders.terms(NodeComponentTable.COLUMN_PEER).field(NodeComponentTable.COLUMN_PEER).size(100)) + .subAggregation(AggregationBuilders.terms(NodeComponentTable.COLUMN_PEER_ID).field(NodeComponentTable.COLUMN_PEER_ID).size(100))); + SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); - Terms genders = searchResponse.getAggregations().get(NodeComponentTable.COLUMN_AGG); - + Terms componentIdTerms = searchResponse.getAggregations().get(NodeComponentTable.COLUMN_COMPONENT_ID); JsonArray nodeComponentArray = new JsonArray(); - for (Terms.Bucket entry : genders.getBuckets()) { - String aggId = entry.getKeyAsString(); - String[] aggIds = aggId.split(Const.IDS_SPLIT); - String code = aggIds[0]; - String peers = aggIds[1]; - - JsonObject nodeComponentObj = new JsonObject(); - nodeComponentObj.addProperty("name", code); - nodeComponentObj.addProperty("peers", peers); - nodeComponentArray.add(nodeComponentObj); + for (Terms.Bucket componentIdBucket : componentIdTerms.getBuckets()) { + int componentId = componentIdBucket.getKeyAsNumber().intValue(); + String componentName = ComponentsDefine.getInstance().getComponentName(componentId); + if (componentId != 0) { + buildComponentArray(componentIdBucket, componentName, nodeComponentArray); + } } - logger.debug("node component data: {}", nodeComponentArray.toString()); + return nodeComponentArray; } + + private JsonArray aggregationByComponentName(long startTime, long endTime) { + SearchRequestBuilder searchRequestBuilder = getClient().prepareSearch(NodeComponentTable.TABLE); + searchRequestBuilder.setTypes(NodeComponentTable.TABLE_TYPE); + searchRequestBuilder.setSearchType(SearchType.DFS_QUERY_THEN_FETCH); + searchRequestBuilder.setQuery(QueryBuilders.rangeQuery(NodeComponentTable.COLUMN_TIME_BUCKET).gte(startTime).lte(endTime)); + searchRequestBuilder.setSize(0); + + searchRequestBuilder.addAggregation(AggregationBuilders.terms(NodeComponentTable.COLUMN_COMPONENT_NAME).field(NodeComponentTable.COLUMN_COMPONENT_NAME).size(100) + .subAggregation(AggregationBuilders.terms(NodeComponentTable.COLUMN_PEER).field(NodeComponentTable.COLUMN_PEER).size(100)) + .subAggregation(AggregationBuilders.terms(NodeComponentTable.COLUMN_PEER_ID).field(NodeComponentTable.COLUMN_PEER_ID).size(100))); + + SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); + + Terms componentNameTerms = searchResponse.getAggregations().get(NodeComponentTable.COLUMN_COMPONENT_NAME); + JsonArray nodeComponentArray = new JsonArray(); + for (Terms.Bucket componentNameBucket : componentNameTerms.getBuckets()) { + String componentName = componentNameBucket.getKeyAsString(); + if (StringUtils.isNotEmpty(componentName)) { + buildComponentArray(componentNameBucket, componentName, nodeComponentArray); + } + } + + return nodeComponentArray; + } + + private void buildComponentArray(Terms.Bucket componentBucket, String componentName, JsonArray nodeComponentArray) { + Terms peerIdTerms = componentBucket.getAggregations().get(NodeComponentTable.COLUMN_PEER_ID); + for (Terms.Bucket peerIdBucket : peerIdTerms.getBuckets()) { + int peerId = peerIdBucket.getKeyAsNumber().intValue(); + + if (peerId != 0) { + String peer = ApplicationCache.getForUI(peerId); + + JsonObject nodeComponentObj = new JsonObject(); + nodeComponentObj.addProperty("componentName", componentName); + nodeComponentObj.addProperty("peer", peer); + nodeComponentArray.add(nodeComponentObj); + } + } + + Terms peerTerms = componentBucket.getAggregations().get(NodeComponentTable.COLUMN_PEER); + for (Terms.Bucket peerBucket : peerTerms.getBuckets()) { + String peer = peerBucket.getKeyAsString(); + + if (StringUtils.isNotEmpty(peer)) { + JsonObject nodeComponentObj = new JsonObject(); + nodeComponentObj.addProperty("componentName", componentName); + nodeComponentObj.addProperty("peer", peer); + nodeComponentArray.add(nodeComponentObj); + } + } + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingEsDAO.java index 55d07837d..fafb6a448 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingEsDAO.java @@ -8,9 +8,10 @@ import org.elasticsearch.action.search.SearchType; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.bucket.terms.Terms; -import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.core.util.StringUtils; import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.ui.cache.ApplicationCache; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -28,22 +29,42 @@ public class NodeMappingEsDAO extends EsDAO implements INodeMappingDAO { searchRequestBuilder.setQuery(QueryBuilders.rangeQuery(NodeMappingTable.COLUMN_TIME_BUCKET).gte(startTime).lte(endTime)); searchRequestBuilder.setSize(0); - searchRequestBuilder.addAggregation(AggregationBuilders.terms(NodeMappingTable.COLUMN_AGG).field(NodeMappingTable.COLUMN_AGG).size(100)); + searchRequestBuilder.addAggregation( + AggregationBuilders.terms(NodeMappingTable.COLUMN_APPLICATION_ID).field(NodeMappingTable.COLUMN_APPLICATION_ID).size(100) + .subAggregation(AggregationBuilders.terms(NodeMappingTable.COLUMN_ADDRESS_ID).field(NodeMappingTable.COLUMN_ADDRESS_ID).size(100)) + .subAggregation(AggregationBuilders.terms(NodeMappingTable.COLUMN_ADDRESS).field(NodeMappingTable.COLUMN_ADDRESS).size(100))); SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); - Terms genders = searchResponse.getAggregations().get(NodeMappingTable.COLUMN_AGG); + Terms applicationIdTerms = searchResponse.getAggregations().get(NodeMappingTable.COLUMN_APPLICATION_ID); JsonArray nodeMappingArray = new JsonArray(); - for (Terms.Bucket entry : genders.getBuckets()) { - String aggId = entry.getKeyAsString(); - String[] aggIds = aggId.split(Const.IDS_SPLIT); - String code = aggIds[0]; - String peers = aggIds[1]; + for (Terms.Bucket applicationIdBucket : applicationIdTerms.getBuckets()) { + int applicationId = applicationIdBucket.getKeyAsNumber().intValue(); + String applicationCode = ApplicationCache.getForUI(applicationId); + Terms addressIdTerms = applicationIdBucket.getAggregations().get(NodeMappingTable.COLUMN_ADDRESS_ID); + for (Terms.Bucket addressIdBucket : addressIdTerms.getBuckets()) { + int addressId = addressIdBucket.getKeyAsNumber().intValue(); + String address = ApplicationCache.getForUI(addressId); - JsonObject nodeMappingObj = new JsonObject(); - nodeMappingObj.addProperty("code", code); - nodeMappingObj.addProperty("peers", peers); - nodeMappingArray.add(nodeMappingObj); + if (addressId != 0) { + JsonObject nodeMappingObj = new JsonObject(); + nodeMappingObj.addProperty("applicationCode", applicationCode); + nodeMappingObj.addProperty("address", address); + nodeMappingArray.add(nodeMappingObj); + } + } + + Terms addressTerms = applicationIdBucket.getAggregations().get(NodeMappingTable.COLUMN_ADDRESS); + for (Terms.Bucket addressBucket : addressTerms.getBuckets()) { + String address = addressBucket.getKeyAsString(); + + if (StringUtils.isNotEmpty(address)) { + JsonObject nodeMappingObj = new JsonObject(); + nodeMappingObj.addProperty("applicationCode", applicationCode); + nodeMappingObj.addProperty("address", address); + nodeMappingArray.add(nodeMappingObj); + } + } } logger.debug("node mapping data: {}", nodeMappingArray.toString()); return nodeMappingArray; diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeRefSumEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeRefSumEsDAO.java index 520c30618..f1c465608 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeRefSumEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeRefSumEsDAO.java @@ -10,9 +10,10 @@ import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.bucket.terms.TermsAggregationBuilder; import org.elasticsearch.search.aggregations.metrics.sum.Sum; -import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.core.util.StringUtils; import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.ui.cache.ApplicationCache; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -30,46 +31,85 @@ public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO { searchRequestBuilder.setQuery(QueryBuilders.rangeQuery(NodeRefSumTable.COLUMN_TIME_BUCKET).gte(startTime).lte(endTime)); searchRequestBuilder.setSize(0); - TermsAggregationBuilder aggregationBuilder = AggregationBuilders.terms(NodeRefSumTable.COLUMN_AGG).field(NodeRefSumTable.COLUMN_AGG); - aggregationBuilder.subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_ONE_SECOND_LESS).field(NodeRefSumTable.COLUMN_ONE_SECOND_LESS)); - aggregationBuilder.subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_THREE_SECOND_LESS).field(NodeRefSumTable.COLUMN_THREE_SECOND_LESS)); - aggregationBuilder.subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS).field(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS)); - aggregationBuilder.subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER).field(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER)); - aggregationBuilder.subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_ERROR).field(NodeRefSumTable.COLUMN_ERROR)); - aggregationBuilder.subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_SUMMARY).field(NodeRefSumTable.COLUMN_SUMMARY)); + TermsAggregationBuilder aggregationBuilder = AggregationBuilders.terms(NodeRefSumTable.COLUMN_FRONT_APPLICATION_ID).field(NodeRefSumTable.COLUMN_FRONT_APPLICATION_ID).size(100); + aggregationBuilder.subAggregation(AggregationBuilders.terms(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID).field(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID).size(100) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S1_LTE).field(NodeRefSumTable.COLUMN_S1_LTE)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S3_LTE).field(NodeRefSumTable.COLUMN_S3_LTE)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S5_LTE).field(NodeRefSumTable.COLUMN_S5_LTE)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S5_GT).field(NodeRefSumTable.COLUMN_S5_GT)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_SUMMARY).field(NodeRefSumTable.COLUMN_SUMMARY)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_ERROR).field(NodeRefSumTable.COLUMN_ERROR))); + aggregationBuilder.subAggregation(AggregationBuilders.terms(NodeRefSumTable.COLUMN_BEHIND_PEER).field(NodeRefSumTable.COLUMN_BEHIND_PEER).size(100) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S1_LTE).field(NodeRefSumTable.COLUMN_S1_LTE)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S3_LTE).field(NodeRefSumTable.COLUMN_S3_LTE)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S5_LTE).field(NodeRefSumTable.COLUMN_S5_LTE)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_S5_GT).field(NodeRefSumTable.COLUMN_S5_GT)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_SUMMARY).field(NodeRefSumTable.COLUMN_SUMMARY)) + .subAggregation(AggregationBuilders.sum(NodeRefSumTable.COLUMN_ERROR).field(NodeRefSumTable.COLUMN_ERROR))); searchRequestBuilder.addAggregation(aggregationBuilder); - SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); JsonArray nodeRefResSumArray = new JsonArray(); - Terms aggTerms = searchResponse.getAggregations().get(NodeRefSumTable.COLUMN_AGG); - for (Terms.Bucket bucket : aggTerms.getBuckets()) { - String aggId = String.valueOf(bucket.getKey()); - Sum oneSecondLess = bucket.getAggregations().get(NodeRefSumTable.COLUMN_ONE_SECOND_LESS); - Sum threeSecondLess = bucket.getAggregations().get(NodeRefSumTable.COLUMN_THREE_SECOND_LESS); - Sum fiveSecondLess = bucket.getAggregations().get(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS); - Sum fiveSecondGreater = bucket.getAggregations().get(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER); - Sum error = bucket.getAggregations().get(NodeRefSumTable.COLUMN_ERROR); - Sum summary = bucket.getAggregations().get(NodeRefSumTable.COLUMN_SUMMARY); - logger.debug("aggId: {}, oneSecondLess: {}, threeSecondLess: {}, fiveSecondLess: {}, fiveSecondGreater: {}, error: {}, summary: {}", aggId, - oneSecondLess.getValue(), threeSecondLess.getValue(), fiveSecondLess.getValue(), fiveSecondGreater.getValue(), error.getValue(), summary.getValue()); + Terms frontApplicationIdTerms = searchResponse.getAggregations().get(NodeRefSumTable.COLUMN_FRONT_APPLICATION_ID); + for (Terms.Bucket frontApplicationIdBucket : frontApplicationIdTerms.getBuckets()) { + int applicationId = frontApplicationIdBucket.getKeyAsNumber().intValue(); + String applicationCode = ApplicationCache.getForUI(applicationId); + Terms behindApplicationIdTerms = frontApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_BEHIND_APPLICATION_ID); + for (Terms.Bucket behindApplicationIdBucket : behindApplicationIdTerms.getBuckets()) { + int behindApplicationId = behindApplicationIdBucket.getKeyAsNumber().intValue(); - JsonObject nodeRefResSumObj = new JsonObject(); - String[] ids = aggId.split(Const.IDS_SPLIT); - String front = ids[0]; - String behind = ids[1]; + if (behindApplicationId != 0) { + String behindApplicationCode = ApplicationCache.getForUI(behindApplicationId); - nodeRefResSumObj.addProperty("front", front); - nodeRefResSumObj.addProperty("behind", behind); + Sum s1LTE = behindApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_S1_LTE); + Sum s3LTE = behindApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_S3_LTE); + Sum s5LTE = behindApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_S5_LTE); + Sum s5GT = behindApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_S5_GT); + Sum summary = behindApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_SUMMARY); + Sum error = behindApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_ERROR); + logger.debug("applicationId: {}, behindApplicationId: {}, s1LTE: {}, s3LTE: {}, s5LTE: {}, s5GT: {}, error: {}, summary: {}", applicationId, + behindApplicationId, s1LTE.getValue(), s3LTE.getValue(), s5LTE.getValue(), s5GT.getValue(), error.getValue(), summary.getValue()); - nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, oneSecondLess.getValue()); - nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, threeSecondLess.getValue()); - nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, fiveSecondLess.getValue()); - nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, fiveSecondGreater.getValue()); - nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_ERROR, error.getValue()); - nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_SUMMARY, summary.getValue()); - nodeRefResSumArray.add(nodeRefResSumObj); + JsonObject nodeRefResSumObj = new JsonObject(); + nodeRefResSumObj.addProperty("front", applicationCode); + nodeRefResSumObj.addProperty("behind", behindApplicationCode); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S1_LTE, s1LTE.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S3_LTE, s3LTE.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S5_LTE, s5LTE.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S5_GT, s5GT.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_ERROR, error.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_SUMMARY, summary.getValue()); + nodeRefResSumArray.add(nodeRefResSumObj); + } + } + + Terms behindPeerTerms = frontApplicationIdBucket.getAggregations().get(NodeRefSumTable.COLUMN_BEHIND_PEER); + for (Terms.Bucket behindPeerBucket : behindPeerTerms.getBuckets()) { + String behindPeer = behindPeerBucket.getKeyAsString(); + + if (StringUtils.isNotEmpty(behindPeer)) { + Sum s1LTE = behindPeerBucket.getAggregations().get(NodeRefSumTable.COLUMN_S1_LTE); + Sum s3LTE = behindPeerBucket.getAggregations().get(NodeRefSumTable.COLUMN_S3_LTE); + Sum s5LTE = behindPeerBucket.getAggregations().get(NodeRefSumTable.COLUMN_S5_LTE); + Sum s5GT = behindPeerBucket.getAggregations().get(NodeRefSumTable.COLUMN_S5_GT); + Sum summary = behindPeerBucket.getAggregations().get(NodeRefSumTable.COLUMN_SUMMARY); + Sum error = behindPeerBucket.getAggregations().get(NodeRefSumTable.COLUMN_ERROR); + logger.debug("applicationId: {}, behindPeer: {}, s1LTE: {}, s3LTE: {}, s5LTE: {}, s5GT: {}, error: {}, summary: {}", applicationId, + behindPeer, s1LTE.getValue(), s3LTE.getValue(), s5LTE.getValue(), s5GT.getValue(), error.getValue(), summary.getValue()); + + JsonObject nodeRefResSumObj = new JsonObject(); + nodeRefResSumObj.addProperty("front", applicationCode); + nodeRefResSumObj.addProperty("behind", behindPeer); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S1_LTE, s1LTE.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S3_LTE, s3LTE.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S5_LTE, s5LTE.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_S5_GT, s5GT.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_ERROR, error.getValue()); + nodeRefResSumObj.addProperty(NodeRefSumTable.COLUMN_SUMMARY, summary.getValue()); + nodeRefResSumArray.add(nodeRefResSumObj); + } + } } return nodeRefResSumArray; diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceEsDAO.java index b189c339b..6ac547ee7 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceEsDAO.java @@ -1,16 +1,14 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonArray; -import com.google.gson.JsonObject; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.action.search.SearchType; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.bucket.terms.Terms; -import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.storage.define.noderef.NodeRefTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -35,17 +33,6 @@ public class NodeReferenceEsDAO extends EsDAO implements INodeReferenceDAO { Terms genders = searchResponse.getAggregations().get(NodeRefTable.COLUMN_AGG); JsonArray nodeRefArray = new JsonArray(); - for (Terms.Bucket entry : genders.getBuckets()) { - String aggId = entry.getKeyAsString(); - String[] aggIds = aggId.split(Const.IDS_SPLIT); - String front = aggIds[0]; - String behind = aggIds[1]; - - JsonObject nodeRefObj = new JsonObject(); - nodeRefObj.addProperty("front", front); - nodeRefObj.addProperty("behind", behind); - nodeRefArray.add(nodeRefObj); - } logger.debug("node ref data: {}", nodeRefArray.toString()); return nodeRefArray; } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java index a11540d83..687edd156 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java @@ -22,14 +22,29 @@ public class TraceDagGetHandler extends JettyHandler { private TraceDagService service = new TraceDagService(); @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + if (!req.getParameterMap().containsKey("startTime") || !req.getParameterMap().containsKey("endTime")) { + throw new ArgumentsParseException("the request parameter must contains startTime, endTime"); + } + String startTimeStr = req.getParameter("startTime"); String endTimeStr = req.getParameter("endTime"); - String timeBucketType = req.getParameter("timeBucketType"); - logger.debug("startTime: {}, endTimeStr: {}, timeBucketType: {}", startTimeStr, endTimeStr, timeBucketType); + logger.debug("startTime: {}, endTimeStr: {}", startTimeStr, endTimeStr); - long startTime = Long.valueOf(startTimeStr); - long endTime = Long.valueOf(endTimeStr); - return service.load(startTime, endTime, timeBucketType); + long startTime; + try { + startTime = Long.valueOf(req.getParameter("startTime")); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("the request parameter startTime must be a long"); + } + + long endTime; + try { + endTime = Long.valueOf(req.getParameter("endTime")); + } catch (NumberFormatException e) { + throw new ArgumentsParseException("the request parameter endTime must be a long"); + } + + return service.load(startTime, endTime); } @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagDataBuilder.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagDataBuilder.java index 0c9d03d5d..73d842a7e 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagDataBuilder.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagDataBuilder.java @@ -5,6 +5,7 @@ import com.google.gson.JsonObject; import java.util.HashMap; import java.util.Map; import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.storage.define.noderef.NodeRefSumTable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -17,29 +18,27 @@ public class TraceDagDataBuilder { private Integer nodeId = new Integer(-1); private Map mappingMap = new HashMap<>(); private Map nodeCompMap = new HashMap<>(); - private Map resSumMap = new HashMap<>(); private Map nodeIdMap = new HashMap<>(); private JsonArray pointArray = new JsonArray(); private JsonArray lineArray = new JsonArray(); - public JsonObject build(JsonArray nodeCompArray, JsonArray nodesMappingArray, JsonArray nodeRefsArray, - JsonArray resSumArray) { - changeMapping2Map(nodesMappingArray); + public JsonObject build(JsonArray nodeCompArray, JsonArray nodesMappingArray, JsonArray resSumArray) { changeNodeComp2Map(nodeCompArray); - resSumMerge(resSumArray); + changeMapping2Map(nodesMappingArray); - for (int i = 0; i < nodeRefsArray.size(); i++) { - JsonObject nodeRefJsonObj = nodeRefsArray.get(i).getAsJsonObject(); + for (int i = 0; i < resSumArray.size(); i++) { + JsonObject nodeRefJsonObj = resSumArray.get(i).getAsJsonObject(); String front = nodeRefJsonObj.get("front").getAsString(); String behind = nodeRefJsonObj.get("behind").getAsString(); - String behindCode = findRealCode(behind); - logger.debug("behind: %s, behindCode: {}", behind, behindCode); + if (hasMapping(behind)) { + continue; + } JsonObject lineJsonObj = new JsonObject(); lineJsonObj.addProperty("from", findOrCreateNode(front)); - lineJsonObj.addProperty("to", findOrCreateNode(behindCode)); - lineJsonObj.addProperty("resSum", resSumMap.get(front + Const.ID_SPLIT + behindCode)); + lineJsonObj.addProperty("to", findOrCreateNode(behind)); + lineJsonObj.addProperty("resSum", nodeRefJsonObj.get(NodeRefSumTable.COLUMN_SUMMARY).getAsInt()); lineArray.add(lineJsonObj); logger.debug("line: {}", lineJsonObj); @@ -75,9 +74,9 @@ public class TraceDagDataBuilder { private void changeMapping2Map(JsonArray nodesMappingArray) { for (int i = 0; i < nodesMappingArray.size(); i++) { JsonObject nodesMappingJsonObj = nodesMappingArray.get(i).getAsJsonObject(); - String code = nodesMappingJsonObj.get("code").getAsString(); - String peers = nodesMappingJsonObj.get("peers").getAsString(); - mappingMap.put(peers, code); + String applicationCode = nodesMappingJsonObj.get("applicationCode").getAsString(); + String address = nodesMappingJsonObj.get("address").getAsString(); + mappingMap.put(address, applicationCode); } } @@ -85,37 +84,13 @@ public class TraceDagDataBuilder { for (int i = 0; i < nodeCompArray.size(); i++) { JsonObject nodesJsonObj = nodeCompArray.get(i).getAsJsonObject(); logger.debug(nodesJsonObj.toString()); - String component = nodesJsonObj.get("name").getAsString(); - String peers = nodesJsonObj.get("peers").getAsString(); - nodeCompMap.put(peers, component); + String componentName = nodesJsonObj.get("componentName").getAsString(); + String peer = nodesJsonObj.get("peer").getAsString(); + nodeCompMap.put(peer, componentName); } } - private String findRealCode(String peers) { - if (mappingMap.containsKey(peers)) { - return mappingMap.get(peers); - } else { - return peers; - } - } - - private void resSumMerge(JsonArray resSumArray) { - for (int i = 0; i < resSumArray.size(); i++) { - JsonObject resSumJsonObj = resSumArray.get(i).getAsJsonObject(); - String front = resSumJsonObj.get("front").getAsString(); - String behind = resSumJsonObj.get("behind").getAsString(); - Long summary = resSumJsonObj.get("summary").getAsLong(); - - if (mappingMap.containsKey(behind)) { - behind = mappingMap.get(behind); - } - - String id = front + Const.ID_SPLIT + behind; - if (resSumMap.containsKey(id)) { - resSumMap.put(id, summary + resSumMap.get(id)); - } else { - resSumMap.put(id, summary); - } - } + private boolean hasMapping(String peers) { + return mappingMap.containsKey(peers); } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagService.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagService.java index 1dbc43fff..5bcc0ddc3 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagService.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/TraceDagService.java @@ -6,7 +6,6 @@ import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.ui.dao.INodeComponentDAO; import org.skywalking.apm.collector.ui.dao.INodeMappingDAO; import org.skywalking.apm.collector.ui.dao.INodeRefSumDAO; -import org.skywalking.apm.collector.ui.dao.INodeReferenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -17,22 +16,19 @@ public class TraceDagService { private final Logger logger = LoggerFactory.getLogger(TraceDagService.class); - public JsonObject load(long startTime, long endTime, String timeBucketType) { - logger.debug("startTime: {}, endTime: {}, timeSliceType: {}", startTime, endTime, timeBucketType); + public JsonObject load(long startTime, long endTime) { + logger.debug("startTime: {}, endTime: {}", startTime, endTime); INodeComponentDAO nodeComponentDAO = (INodeComponentDAO)DAOContainer.INSTANCE.get(INodeComponentDAO.class.getName()); JsonArray nodeComponentArray = nodeComponentDAO.load(startTime, endTime); INodeMappingDAO nodeMappingDAO = (INodeMappingDAO)DAOContainer.INSTANCE.get(INodeMappingDAO.class.getName()); JsonArray nodeMappingArray = nodeMappingDAO.load(startTime, endTime); - INodeReferenceDAO nodeReferenceDAO = (INodeReferenceDAO)DAOContainer.INSTANCE.get(INodeReferenceDAO.class.getName()); - JsonArray nodeRefArray = nodeReferenceDAO.load(startTime, endTime); - INodeRefSumDAO nodeRefSumDAO = (INodeRefSumDAO)DAOContainer.INSTANCE.get(INodeRefSumDAO.class.getName()); JsonArray nodeRefSumArray = nodeRefSumDAO.load(startTime, endTime); TraceDagDataBuilder builder = new TraceDagDataBuilder(); - JsonObject traceDag = builder.build(nodeComponentArray, nodeMappingArray, nodeRefArray, nodeRefSumArray); + JsonObject traceDag = builder.build(nodeComponentArray, nodeMappingArray, nodeRefSumArray); return traceDag; }