diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/NetworkAddressInventory.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/NetworkAddressInventory.java index c68ba637f..fadb00cb6 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/NetworkAddressInventory.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/NetworkAddressInventory.java @@ -39,8 +39,10 @@ public class NetworkAddressInventory extends RegisterSource { public static final String MODEL_NAME = "network_address_inventory"; private static final String NAME = "name"; + public static final String SRC_LAYER = "src_layer"; @Setter @Getter @Column(columnName = NAME, matchQuery = true) private String name = Const.EMPTY_STRING; + @Setter @Getter @Column(columnName = SRC_LAYER) private int srcLayer; public static String buildId(String networkAddress) { return networkAddress; @@ -71,9 +73,16 @@ public class NetworkAddressInventory extends RegisterSource { return true; } + @Override public void combine(RegisterSource registerSource) { + super.combine(registerSource); + NetworkAddressInventory inventory = (NetworkAddressInventory)registerSource; + setSrcLayer(inventory.srcLayer); + } + @Override public RemoteData.Builder serialize() { RemoteData.Builder remoteBuilder = RemoteData.newBuilder(); remoteBuilder.setDataIntegers(0, getSequence()); + remoteBuilder.setDataIntegers(1, getSrcLayer()); remoteBuilder.setDataLongs(0, getRegisterTime()); remoteBuilder.setDataLongs(1, getHeartbeatTime()); @@ -84,6 +93,7 @@ public class NetworkAddressInventory extends RegisterSource { @Override public void deserialize(RemoteData remoteData) { setSequence(remoteData.getDataIntegers(0)); + setSrcLayer(remoteData.getDataIntegers(1)); setRegisterTime(remoteData.getDataLongs(0)); setHeartbeatTime(remoteData.getDataLongs(1)); @@ -101,6 +111,7 @@ public class NetworkAddressInventory extends RegisterSource { NetworkAddressInventory inventory = new NetworkAddressInventory(); inventory.setSequence((Integer)dbMap.get(SEQUENCE)); inventory.setName((String)dbMap.get(NAME)); + inventory.setSrcLayer((Integer)dbMap.get(SRC_LAYER)); inventory.setRegisterTime((Long)dbMap.get(REGISTER_TIME)); inventory.setHeartbeatTime((Long)dbMap.get(HEARTBEAT_TIME)); return inventory; @@ -110,6 +121,7 @@ public class NetworkAddressInventory extends RegisterSource { Map map = new HashMap<>(); map.put(SEQUENCE, storageData.getSequence()); map.put(NAME, storageData.getName()); + map.put(SRC_LAYER, storageData.getSrcLayer()); map.put(REGISTER_TIME, storageData.getRegisterTime()); map.put(HEARTBEAT_TIME, storageData.getHeartbeatTime()); return map; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/INetworkAddressInventoryRegister.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/INetworkAddressInventoryRegister.java index 49c88483b..02c93b1b6 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/INetworkAddressInventoryRegister.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/INetworkAddressInventoryRegister.java @@ -29,4 +29,6 @@ public interface INetworkAddressInventoryRegister extends Service { int get(String networkAddress); void heartbeat(int addressId, long heartBeatTime); + + void update(int addressId, int srcLayer); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/NetworkAddressInventoryRegister.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/NetworkAddressInventoryRegister.java index 3cf548198..db72c87ff 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/NetworkAddressInventoryRegister.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/service/NetworkAddressInventoryRegister.java @@ -106,4 +106,23 @@ public class NetworkAddressInventoryRegister implements INetworkAddressInventory logger.warn("Network address {} heartbeat, but not found in storage."); } } + + @Override public void update(int addressId, int srcLayer) { + if (!this.compare(addressId, srcLayer)) { + NetworkAddressInventory newNetworkAddress = getNetworkAddressInventoryCache().get(addressId); + newNetworkAddress.setSrcLayer(srcLayer); + newNetworkAddress.setHeartbeatTime(System.currentTimeMillis()); + + InventoryProcess.INSTANCE.in(newNetworkAddress); + } + } + + private boolean compare(int addressId, int srcLayer) { + NetworkAddressInventory networkAddress = getNetworkAddressInventoryCache().get(addressId); + + if (Objects.nonNull(networkAddress)) { + return srcLayer == networkAddress.getSrcLayer(); + } + return true; + } } diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/standardization/SpanIdExchanger.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/standardization/SpanIdExchanger.java index 98e9c9eb1..b19d9e637 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/standardization/SpanIdExchanger.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/standardization/SpanIdExchanger.java @@ -80,6 +80,9 @@ public class SpanIdExchanger implements IdExchanger { standardBuilder.toBuilder(); standardBuilder.setPeerId(peerId); standardBuilder.setPeer(Const.EMPTY_STRING); + + int spanLayer = standardBuilder.getSpanLayerValue(); + networkAddressInventoryRegister.update(peerId, spanLayer); } } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java index 14dac1620..c334431f2 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java @@ -77,6 +77,7 @@ public class MetadataQueryEsDAO extends EsDAO implements IMetadataQueryDAO { BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery(); boolQueryBuilder.must().add(timeRangeQueryBuild(startTimestamp, endTimestamp)); + boolQueryBuilder.must().add(QueryBuilders.termQuery(NetworkAddressInventory.SRC_LAYER, srcLayer)); sourceBuilder.query(boolQueryBuilder); sourceBuilder.size(0);