From b367c36db9af7186cf37a97c015d83e33a6d58c7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=AD=E5=8B=87=E5=8D=87=20pengys?= <8082209@qq.com> Date: Wed, 17 Oct 2018 11:28:16 +0800 Subject: [PATCH] Fixed the endpoint topology bug. (#1778) * Delete the client side endpoint relation indicator. Fixed the endpoint topology bug. * Fixed package failure bug. --- .../EndpointCallRelationDispatcher.java | 13 +- .../EndpointRelationClientSideIndicator.java | 151 ------------------ .../core/query/TopologyQueryService.java | 32 +++- .../core/storage/query/ITopologyQueryDAO.java | 3 - .../endpoint/MultiScopesSpanListener.java | 1 - .../query/TopologyQueryEsDAO.java | 25 +-- 6 files changed, 34 insertions(+), 191 deletions(-) delete mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpointrelation/EndpointRelationClientSideIndicator.java diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpointrelation/EndpointCallRelationDispatcher.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpointrelation/EndpointCallRelationDispatcher.java index 8b011694c..6dc5ea0dd 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpointrelation/EndpointCallRelationDispatcher.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpointrelation/EndpointCallRelationDispatcher.java @@ -23,15 +23,12 @@ import org.apache.skywalking.oap.server.core.analysis.worker.IndicatorProcess; import org.apache.skywalking.oap.server.core.source.EndpointRelation; /** - * @author wusheng + * @author wusheng, peng-yongsheng */ public class EndpointCallRelationDispatcher implements SourceDispatcher { @Override public void dispatch(EndpointRelation source) { switch (source.getDetectPoint()) { - case CLIENT: - clientSide(source); - break; case SERVER: serverSide(source); break; @@ -45,12 +42,4 @@ public class EndpointCallRelationDispatcher implements SourceDispatcher { - - @Override public EndpointRelationClientSideIndicator map2Data(Map dbMap) { - EndpointRelationClientSideIndicator indicator = new EndpointRelationClientSideIndicator(); - indicator.setSourceEndpointId(((Number)dbMap.get(SOURCE_ENDPOINT_ID)).intValue()); - indicator.setDestEndpointId(((Number)dbMap.get(DEST_ENDPOINT_ID)).intValue()); - indicator.setTimeBucket(((Number)dbMap.get(TIME_BUCKET)).longValue()); - return indicator; - } - - @Override public Map data2Map(EndpointRelationClientSideIndicator storageData) { - Map map = new HashMap<>(); - map.put(TIME_BUCKET, storageData.getTimeBucket()); - map.put(SOURCE_ENDPOINT_ID, storageData.getSourceEndpointId()); - map.put(DEST_ENDPOINT_ID, storageData.getDestEndpointId()); - return map; - } - } -} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TopologyQueryService.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TopologyQueryService.java index 74b80fda2..766fa7d79 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TopologyQueryService.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TopologyQueryService.java @@ -109,15 +109,35 @@ public class TopologyQueryService implements Service { Map components = new HashMap<>(); serviceComponents.forEach(component -> components.put(component.getServiceId(), getComponentLibraryCatalogService().getComponentName(component.getComponentId()))); - List calls = getTopologyQueryDAO().loadSpecifiedDestOfServerSideEndpointRelations(step, startTB, endTB, endpointId); - calls.addAll(getTopologyQueryDAO().loadSpecifiedSourceOfClientSideEndpointRelations(step, startTB, endTB, endpointId)); + List serverSideCalls = getTopologyQueryDAO().loadSpecifiedDestOfServerSideEndpointRelations(step, startTB, endTB, endpointId); + serverSideCalls.forEach(call -> call.setDetectPoint(DetectPoint.SERVER)); - calls.forEach(call -> { - call.setCallType(components.getOrDefault(getEndpointInventoryCache().get(call.getTarget()).getServiceId(), Const.UNKNOWN)); - }); + serverSideCalls.forEach(call -> call.setCallType(components.getOrDefault(getEndpointInventoryCache().get(call.getTarget()).getServiceId(), Const.UNKNOWN))); Topology topology = new Topology(); - topology.getCalls().addAll(calls); + topology.getCalls().addAll(serverSideCalls); + + Set nodeIds = new HashSet<>(); + serverSideCalls.forEach(call -> { + if (!nodeIds.contains(call.getSource())) { + topology.getNodes().add(buildEndpointNode(components, call.getSource())); + nodeIds.add(call.getSource()); + } + if (!nodeIds.contains(call.getTarget())) { + topology.getNodes().add(buildEndpointNode(components, call.getTarget())); + nodeIds.add(call.getTarget()); + } + }); + return topology; } + + private Node buildEndpointNode(Map components, int endpointId) { + Node node = new Node(); + node.setId(endpointId); + node.setName(getEndpointInventoryCache().get(endpointId).getName()); + node.setType(components.getOrDefault(getEndpointInventoryCache().get(endpointId).getServiceId(), Const.UNKNOWN)); + node.setReal(true); + return node; + } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITopologyQueryDAO.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITopologyQueryDAO.java index 514fb75b8..440b8b3a8 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITopologyQueryDAO.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITopologyQueryDAO.java @@ -45,7 +45,4 @@ public interface ITopologyQueryDAO extends Service { List loadSpecifiedDestOfServerSideEndpointRelations(Step step, long startTB, long endTB, int destEndpointId) throws IOException; - - List loadSpecifiedSourceOfClientSideEndpointRelations(Step step, long startTB, long endTB, - int sourceEndpointId) throws IOException; } diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java index dd7797c56..4e849fa83 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/endpoint/MultiScopesSpanListener.java @@ -182,7 +182,6 @@ public class MultiScopesSpanListener implements EntrySpanListener, ExitSpanListe exitSourceBuilder.setTimeBucket(minuteTimeBucket); sourceReceiver.receive(exitSourceBuilder.toServiceRelation()); sourceReceiver.receive(exitSourceBuilder.toServiceInstanceRelation()); - sourceReceiver.receive(exitSourceBuilder.toEndpointRelation()); }); } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TopologyQueryEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TopologyQueryEsDAO.java index b9d6093d3..05689b665 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TopologyQueryEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TopologyQueryEsDAO.java @@ -21,7 +21,7 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.query; import java.io.IOException; import java.util.*; import org.apache.skywalking.oap.server.core.UnexpectedException; -import org.apache.skywalking.oap.server.core.analysis.manual.endpointrelation.*; +import org.apache.skywalking.oap.server.core.analysis.manual.endpointrelation.EndpointRelationServerSideIndicator; import org.apache.skywalking.oap.server.core.analysis.manual.service.*; import org.apache.skywalking.oap.server.core.analysis.manual.servicerelation.*; import org.apache.skywalking.oap.server.core.query.entity.*; @@ -175,28 +175,17 @@ public class TopologyQueryEsDAO extends EsDAO implements ITopologyQueryDAO { BoolQueryBuilder boolQuery = QueryBuilders.boolQuery(); boolQuery.must().add(QueryBuilders.rangeQuery(EndpointRelationServerSideIndicator.TIME_BUCKET).gte(startTB).lte(endTB)); - boolQuery.must().add(QueryBuilders.termQuery(EndpointRelationServerSideIndicator.DEST_ENDPOINT_ID, destEndpointId)); + + BoolQueryBuilder serviceIdBoolQuery = QueryBuilders.boolQuery(); + boolQuery.must().add(serviceIdBoolQuery); + serviceIdBoolQuery.should().add(QueryBuilders.termQuery(EndpointRelationServerSideIndicator.SOURCE_ENDPOINT_ID, destEndpointId)); + serviceIdBoolQuery.should().add(QueryBuilders.termQuery(EndpointRelationServerSideIndicator.DEST_ENDPOINT_ID, destEndpointId)); + sourceBuilder.query(boolQuery); return load(sourceBuilder, indexName, EndpointRelationServerSideIndicator.SOURCE_ENDPOINT_ID, EndpointRelationServerSideIndicator.DEST_ENDPOINT_ID, Source.Endpoint); } - @Override - public List loadSpecifiedSourceOfClientSideEndpointRelations(Step step, long startTB, long endTB, - int sourceEndpointId) throws IOException { - String indexName = DownsampleingModelNameBuilder.build(step, EndpointRelationClientSideIndicator.INDEX_NAME); - - SearchSourceBuilder sourceBuilder = SearchSourceBuilder.searchSource(); - sourceBuilder.size(0); - - BoolQueryBuilder boolQuery = QueryBuilders.boolQuery(); - boolQuery.must().add(QueryBuilders.rangeQuery(EndpointRelationClientSideIndicator.TIME_BUCKET).gte(startTB).lte(endTB)); - boolQuery.must().add(QueryBuilders.termQuery(EndpointRelationClientSideIndicator.SOURCE_ENDPOINT_ID, sourceEndpointId)); - sourceBuilder.query(boolQuery); - - return load(sourceBuilder, indexName, EndpointRelationClientSideIndicator.SOURCE_ENDPOINT_ID, EndpointRelationClientSideIndicator.DEST_ENDPOINT_ID, Source.Endpoint); - } - private List load(SearchSourceBuilder sourceBuilder, String indexName, String sourceCName, String destCName, Source source) throws IOException { TermsAggregationBuilder sourceAggregation = AggregationBuilders.terms(sourceCName).field(sourceCName).size(1000);