From eda2d8141bff82fd5e4a0b049d535d08ed65436f Mon Sep 17 00:00:00 2001 From: peng-yongsheng <8082209@qq.com> Date: Sun, 11 Feb 2018 17:23:21 +0800 Subject: [PATCH] Update span layer and server type by address id when those not equals to the values in cache. --- .../service/INetworkAddressIDService.java | 4 +- .../NetworkAddressRegisterSerialWorker.java | 44 +++++----- .../provider/service/InstanceIDService.java | 4 +- .../service/NetworkAddressIDService.java | 22 ++++- .../provider/buffer/SegmentBufferManager.java | 2 +- .../standardization/ReferenceIdExchanger.java | 3 +- .../standardization/SpanIdExchanger.java | 7 +- .../service/NetworkAddressCacheService.java | 2 + .../cache/guava/CacheModuleGuavaProvider.java | 4 +- ...a => NetworkAddressCacheGuavaService.java} | 42 +++++++--- .../dao/cache/INetworkAddressCacheDAO.java | 5 +- .../register/INetworkAddressRegisterDAO.java | 2 + .../table/register/NetworkAddress.java | 12 ++- .../table/register/NetworkAddressTable.java | 1 + .../storage/table/register/ServerType.java | 46 +++++++++++ .../table/register/ServerTypeDefine.java | 82 +++++++++++++++++++ .../register/ServerTypeDefineTestCase.java | 42 ++++++++++ .../dao/cache/NetworkAddressEsCacheDAO.java | 24 +++++- .../register/NetworkAddressRegisterEsDAO.java | 10 +++ .../register/NetworkAddressEsTableDefine.java | 2 + .../dao/cache/NetworkAddressH2CacheDAO.java | 29 ++++++- .../register/NetworkAddressRegisterH2DAO.java | 19 +++++ .../register/NetworkAddressH2TableDefine.java | 2 + 23 files changed, 359 insertions(+), 51 deletions(-) rename apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/service/{NetworAddressCacheGuavaService.java => NetworkAddressCacheGuavaService.java} (70%) create mode 100644 apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerType.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefine.java create mode 100644 apm-collector/apm-collector-storage/collector-storage-define/src/test/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefineTestCase.java diff --git a/apm-collector/apm-collector-analysis/analysis-register/register-define/src/main/java/org/apache/skywalking/apm/collector/analysis/register/define/service/INetworkAddressIDService.java b/apm-collector/apm-collector-analysis/analysis-register/register-define/src/main/java/org/apache/skywalking/apm/collector/analysis/register/define/service/INetworkAddressIDService.java index faa1ea16a..2f062f567 100644 --- a/apm-collector/apm-collector-analysis/analysis-register/register-define/src/main/java/org/apache/skywalking/apm/collector/analysis/register/define/service/INetworkAddressIDService.java +++ b/apm-collector/apm-collector-analysis/analysis-register/register-define/src/main/java/org/apache/skywalking/apm/collector/analysis/register/define/service/INetworkAddressIDService.java @@ -24,7 +24,9 @@ import org.apache.skywalking.apm.collector.core.module.Service; * @author peng-yongsheng */ public interface INetworkAddressIDService extends Service { - int create(String networkAddress, int spanLayer); + int getOrCreate(String networkAddress); int get(String networkAddress); + + void update(int addressId, int spanLayer, int serverType); } diff --git a/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/register/NetworkAddressRegisterSerialWorker.java b/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/register/NetworkAddressRegisterSerialWorker.java index f11b7f7bf..fd54dd370 100644 --- a/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/register/NetworkAddressRegisterSerialWorker.java +++ b/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/register/NetworkAddressRegisterSerialWorker.java @@ -41,7 +41,7 @@ public class NetworkAddressRegisterSerialWorker extends AbstractLocalAsyncWorker private final INetworkAddressRegisterDAO networkAddressRegisterDAO; private final NetworkAddressCacheService networkAddressCacheService; - public NetworkAddressRegisterSerialWorker(ModuleManager moduleManager) { + NetworkAddressRegisterSerialWorker(ModuleManager moduleManager) { super(moduleManager); this.networkAddressRegisterDAO = getModuleManager().find(StorageModule.NAME).getService(INetworkAddressRegisterDAO.class); this.networkAddressCacheService = getModuleManager().find(CacheModule.NAME).getService(NetworkAddressCacheService.class); @@ -53,28 +53,32 @@ public class NetworkAddressRegisterSerialWorker extends AbstractLocalAsyncWorker @Override protected void onWork(NetworkAddress networkAddress) throws WorkerException { logger.debug("register network address, address: {}", networkAddress.getNetworkAddress()); - int addressId = networkAddressCacheService.getAddressId(networkAddress.getNetworkAddress()); + if (networkAddress.getAddressId() == 0) { + int addressId = networkAddressCacheService.getAddressId(networkAddress.getNetworkAddress()); - if (addressId == 0) { - NetworkAddress newNetworkAddress; - int min = networkAddressRegisterDAO.getMinNetworkAddressId(); - if (min == 0) { - newNetworkAddress = new NetworkAddress(); - newNetworkAddress.setId("-1"); - newNetworkAddress.setAddressId(-1); - newNetworkAddress.setSpanLayer(networkAddress.getSpanLayer()); - newNetworkAddress.setNetworkAddress(networkAddress.getNetworkAddress()); - } else { - int max = networkAddressRegisterDAO.getMaxNetworkAddressId(); - addressId = IdAutoIncrement.INSTANCE.increment(min, max); + if (addressId == 0) { + NetworkAddress newNetworkAddress; + int min = networkAddressRegisterDAO.getMinNetworkAddressId(); + if (min == 0) { + newNetworkAddress = new NetworkAddress(); + newNetworkAddress.setId("-1"); + newNetworkAddress.setAddressId(-1); + newNetworkAddress.setSpanLayer(networkAddress.getSpanLayer()); + newNetworkAddress.setNetworkAddress(networkAddress.getNetworkAddress()); + } else { + int max = networkAddressRegisterDAO.getMaxNetworkAddressId(); + addressId = IdAutoIncrement.INSTANCE.increment(min, max); - newNetworkAddress = new NetworkAddress(); - newNetworkAddress.setId(String.valueOf(addressId)); - newNetworkAddress.setAddressId(addressId); - newNetworkAddress.setSpanLayer(networkAddress.getSpanLayer()); - newNetworkAddress.setNetworkAddress(networkAddress.getNetworkAddress()); + newNetworkAddress = new NetworkAddress(); + newNetworkAddress.setId(String.valueOf(addressId)); + newNetworkAddress.setAddressId(addressId); + newNetworkAddress.setSpanLayer(networkAddress.getSpanLayer()); + newNetworkAddress.setNetworkAddress(networkAddress.getNetworkAddress()); + } + networkAddressRegisterDAO.save(newNetworkAddress); } - networkAddressRegisterDAO.save(newNetworkAddress); + } else { + networkAddressRegisterDAO.update(networkAddress.getId(), networkAddress.getSpanLayer(), networkAddress.getServerType()); } } diff --git a/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/InstanceIDService.java b/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/InstanceIDService.java index 83a1a360f..b5c2ee38a 100644 --- a/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/InstanceIDService.java +++ b/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/InstanceIDService.java @@ -72,7 +72,7 @@ public class InstanceIDService implements IInstanceIDService { } @Override public int getOrCreateByAgentUUID(int applicationId, String agentUUID, long registerTime, String osInfo) { - logger.debug("get or create instance id by agent UUID, application id: {}, agentUUID: {}, registerTime: {}, osInfo: {}", applicationId, agentUUID, registerTime, osInfo); + logger.debug("get or getOrCreate instance id by agent UUID, application id: {}, agentUUID: {}, registerTime: {}, osInfo: {}", applicationId, agentUUID, registerTime, osInfo); int instanceId = getInstanceCacheService().getInstanceIdByAgentUUID(applicationId, agentUUID); if (instanceId == 0) { @@ -93,7 +93,7 @@ public class InstanceIDService implements IInstanceIDService { } @Override public int getOrCreateByAddressId(int applicationId, int addressId, long registerTime) { - logger.debug("get or create instance id by address id, application id: {}, address id: {}, registerTime: {}", applicationId, addressId, registerTime); + logger.debug("get or getOrCreate instance id by address id, application id: {}, address id: {}, registerTime: {}", applicationId, addressId, registerTime); int instanceId = getInstanceCacheService().getInstanceIdByAddressId(applicationId, addressId); if (instanceId == 0) { diff --git a/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/NetworkAddressIDService.java b/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/NetworkAddressIDService.java index bdf37b1bb..1b524ecb4 100644 --- a/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/NetworkAddressIDService.java +++ b/apm-collector/apm-collector-analysis/analysis-register/register-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/register/provider/service/NetworkAddressIDService.java @@ -28,6 +28,7 @@ import org.apache.skywalking.apm.collector.cache.service.NetworkAddressCacheServ import org.apache.skywalking.apm.collector.core.graph.Graph; import org.apache.skywalking.apm.collector.core.graph.GraphManager; import org.apache.skywalking.apm.collector.core.module.ModuleManager; +import org.apache.skywalking.apm.collector.core.util.Const; import org.apache.skywalking.apm.collector.core.util.ObjectUtils; import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress; @@ -74,7 +75,7 @@ public class NetworkAddressIDService implements INetworkAddressIDService { return this.networkAddressGraph; } - @Override public int create(String networkAddress, int spanLayer) { + @Override public int getOrCreate(String networkAddress) { int addressId = getNetworkAddressCacheService().getAddressId(networkAddress); if (addressId != 0) { @@ -89,10 +90,11 @@ public class NetworkAddressIDService implements INetworkAddressIDService { } } else { NetworkAddress newNetworkAddress = new NetworkAddress(); - newNetworkAddress.setId("0"); + newNetworkAddress.setId(String.valueOf(Const.NONE)); newNetworkAddress.setNetworkAddress(networkAddress); - newNetworkAddress.setSpanLayer(spanLayer); - newNetworkAddress.setAddressId(0); + newNetworkAddress.setSpanLayer(Const.NONE); + newNetworkAddress.setServerType(Const.NONE); + newNetworkAddress.setAddressId(Const.NONE); getNetworkAddressGraph().start(newNetworkAddress); } @@ -103,4 +105,16 @@ public class NetworkAddressIDService implements INetworkAddressIDService { @Override public int get(String networkAddress) { return getNetworkAddressCacheService().getAddressId(networkAddress); } + + @Override public void update(int addressId, int spanLayer, int serverType) { + if (!networkAddressCacheService.compare(addressId, spanLayer, serverType)) { + NetworkAddress newNetworkAddress = new NetworkAddress(); + newNetworkAddress.setId(String.valueOf(addressId)); + newNetworkAddress.setSpanLayer(spanLayer); + newNetworkAddress.setServerType(serverType); + newNetworkAddress.setAddressId(addressId); + + getNetworkAddressGraph().start(newNetworkAddress); + } + } } diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferManager.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferManager.java index c146f6653..6fac316c3 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferManager.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/buffer/SegmentBufferManager.java @@ -81,7 +81,7 @@ public enum SegmentBufferManager { } private void newDataFile() throws IOException { - logger.debug("create new segment buffer file"); + logger.debug("getOrCreate new segment buffer file"); String timeBucket = String.valueOf(TimeBucketUtils.INSTANCE.getSecondTimeBucket(System.currentTimeMillis())); String writeFileName = DATA_FILE_PREFIX + "_" + timeBucket + "." + Const.FILE_SUFFIX; File dataFile = new File(BufferFileConfig.BUFFER_PATH + writeFileName); diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/ReferenceIdExchanger.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/ReferenceIdExchanger.java index 90bcbc364..1ce10b5af 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/ReferenceIdExchanger.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/ReferenceIdExchanger.java @@ -89,7 +89,8 @@ public class ReferenceIdExchanger implements IdExchanger { } if (standardBuilder.getNetworkAddressId() == 0 && StringUtils.isNotEmpty(standardBuilder.getNetworkAddress())) { - int networkAddressId = networkAddressIDService.get(standardBuilder.getNetworkAddress()); + int networkAddressId = networkAddressIDService.getOrCreate(standardBuilder.getNetworkAddress()); + if (networkAddressId == 0) { if (logger.isDebugEnabled()) { logger.debug("network address: {} from application id: {} exchange failed", standardBuilder.getNetworkAddress(), applicationId); diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SpanIdExchanger.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SpanIdExchanger.java index c5389a44f..9236c2d14 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SpanIdExchanger.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/main/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/parser/standardization/SpanIdExchanger.java @@ -25,6 +25,7 @@ import org.apache.skywalking.apm.collector.analysis.segment.parser.define.decora import org.apache.skywalking.apm.collector.core.module.ModuleManager; import org.apache.skywalking.apm.collector.core.util.Const; import org.apache.skywalking.apm.collector.core.util.StringUtils; +import org.apache.skywalking.apm.collector.storage.table.register.ServerTypeDefine; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -53,7 +54,7 @@ public class SpanIdExchanger implements IdExchanger { @Override public boolean exchange(SpanDecorator standardBuilder, int applicationId) { if (standardBuilder.getPeerId() == 0 && StringUtils.isNotEmpty(standardBuilder.getPeer())) { - int peerId = networkAddressIDService.create(standardBuilder.getPeer(), standardBuilder.getSpanLayer().getNumber()); + int peerId = networkAddressIDService.getOrCreate(standardBuilder.getPeer()); if (peerId == 0) { logger.debug("peer: {} in application: {} exchange failed", standardBuilder.getPeer(), applicationId); @@ -62,6 +63,10 @@ public class SpanIdExchanger implements IdExchanger { standardBuilder.toBuilder(); standardBuilder.setPeerId(peerId); standardBuilder.setPeer(Const.EMPTY_STRING); + + int spanLayer = standardBuilder.getSpanLayerValue(); + int serverType = ServerTypeDefine.getInstance().getServerTypeId(standardBuilder.getComponentId()); + networkAddressIDService.update(peerId, spanLayer, serverType); } } diff --git a/apm-collector/apm-collector-cache/collector-cache-define/src/main/java/org/apache/skywalking/apm/collector/cache/service/NetworkAddressCacheService.java b/apm-collector/apm-collector-cache/collector-cache-define/src/main/java/org/apache/skywalking/apm/collector/cache/service/NetworkAddressCacheService.java index f0483af11..e435301a6 100644 --- a/apm-collector/apm-collector-cache/collector-cache-define/src/main/java/org/apache/skywalking/apm/collector/cache/service/NetworkAddressCacheService.java +++ b/apm-collector/apm-collector-cache/collector-cache-define/src/main/java/org/apache/skywalking/apm/collector/cache/service/NetworkAddressCacheService.java @@ -27,4 +27,6 @@ public interface NetworkAddressCacheService extends Service { int getAddressId(String networkAddress); String getAddress(int addressId); + + boolean compare(int addressId, int spanLayer, int serverType); } diff --git a/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/CacheModuleGuavaProvider.java b/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/CacheModuleGuavaProvider.java index af2d54cb1..240b8fae9 100644 --- a/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/CacheModuleGuavaProvider.java +++ b/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/CacheModuleGuavaProvider.java @@ -22,7 +22,7 @@ import java.util.Properties; import org.apache.skywalking.apm.collector.cache.CacheModule; import org.apache.skywalking.apm.collector.cache.guava.service.ApplicationCacheGuavaService; import org.apache.skywalking.apm.collector.cache.guava.service.InstanceCacheGuavaService; -import org.apache.skywalking.apm.collector.cache.guava.service.NetworAddressCacheGuavaService; +import org.apache.skywalking.apm.collector.cache.guava.service.NetworkAddressCacheGuavaService; import org.apache.skywalking.apm.collector.cache.guava.service.ServiceIdCacheGuavaService; import org.apache.skywalking.apm.collector.cache.guava.service.ServiceNameCacheGuavaService; import org.apache.skywalking.apm.collector.cache.service.ApplicationCacheService; @@ -53,7 +53,7 @@ public class CacheModuleGuavaProvider extends ModuleProvider { this.registerServiceImplementation(InstanceCacheService.class, new InstanceCacheGuavaService(getManager())); this.registerServiceImplementation(ServiceIdCacheService.class, new ServiceIdCacheGuavaService(getManager())); this.registerServiceImplementation(ServiceNameCacheService.class, new ServiceNameCacheGuavaService(getManager())); - this.registerServiceImplementation(NetworkAddressCacheService.class, new NetworAddressCacheGuavaService(getManager())); + this.registerServiceImplementation(NetworkAddressCacheService.class, new NetworkAddressCacheGuavaService(getManager())); } @Override public void start(Properties config) throws ServiceNotProvidedException { diff --git a/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/service/NetworAddressCacheGuavaService.java b/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/service/NetworkAddressCacheGuavaService.java similarity index 70% rename from apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/service/NetworAddressCacheGuavaService.java rename to apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/service/NetworkAddressCacheGuavaService.java index 6bf6c8469..5de4c0399 100644 --- a/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/service/NetworAddressCacheGuavaService.java +++ b/apm-collector/apm-collector-cache/collector-cache-guava-provider/src/main/java/org/apache/skywalking/apm/collector/cache/guava/service/NetworkAddressCacheGuavaService.java @@ -27,22 +27,23 @@ import org.apache.skywalking.apm.collector.core.util.ObjectUtils; import org.apache.skywalking.apm.collector.core.util.StringUtils; import org.apache.skywalking.apm.collector.storage.StorageModule; import org.apache.skywalking.apm.collector.storage.dao.cache.INetworkAddressCacheDAO; +import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ -public class NetworAddressCacheGuavaService implements NetworkAddressCacheService { +public class NetworkAddressCacheGuavaService implements NetworkAddressCacheService { - private final Logger logger = LoggerFactory.getLogger(NetworAddressCacheGuavaService.class); + private final Logger logger = LoggerFactory.getLogger(NetworkAddressCacheGuavaService.class); private final Cache addressCache = CacheBuilder.newBuilder().initialCapacity(100).maximumSize(5000).build(); private final ModuleManager moduleManager; private INetworkAddressCacheDAO networkAddressCacheDAO; - public NetworAddressCacheGuavaService(ModuleManager moduleManager) { + public NetworkAddressCacheGuavaService(ModuleManager moduleManager) { this.moduleManager = moduleManager; } @@ -57,16 +58,17 @@ public class NetworAddressCacheGuavaService implements NetworkAddressCacheServic int addressId = 0; try { addressId = addressCache.get(networkAddress, () -> getNetworkAddressCacheDAO().getAddressId(networkAddress)); + + if (addressId == 0) { + addressId = getNetworkAddressCacheDAO().getAddressId(networkAddress); + if (addressId != 0) { + addressCache.put(networkAddress, addressId); + } + } } catch (Throwable e) { logger.error(e.getMessage(), e); } - if (addressId == 0) { - addressId = getNetworkAddressCacheDAO().getAddressId(networkAddress); - if (addressId != 0) { - addressCache.put(networkAddress, addressId); - } - } return addressId; } @@ -75,17 +77,33 @@ public class NetworAddressCacheGuavaService implements NetworkAddressCacheServic public String getAddress(int addressId) { String networkAddress = Const.EMPTY_STRING; try { - networkAddress = idCache.get(addressId, () -> getNetworkAddressCacheDAO().getAddress(addressId)); + networkAddress = idCache.get(addressId, () -> getNetworkAddressCacheDAO().getAddressById(addressId)); } catch (Throwable e) { logger.error(e.getMessage(), e); } if (StringUtils.isEmpty(networkAddress)) { - networkAddress = getNetworkAddressCacheDAO().getAddress(addressId); + networkAddress = getNetworkAddressCacheDAO().getAddressById(addressId); if (StringUtils.isNotEmpty(networkAddress)) { - addressCache.put(networkAddress, addressId); + idCache.put(addressId, networkAddress); } } return networkAddress; } + + private final Cache addressObjCache = CacheBuilder.newBuilder().maximumSize(5000).build(); + + @Override public boolean compare(int addressId, int spanLayer, int serverType) { + try { + NetworkAddress address = addressObjCache.get(addressId, () -> getNetworkAddressCacheDAO().getAddress(addressId)); + if (ObjectUtils.isNotEmpty(address)) { + if (spanLayer != address.getSpanLayer() || serverType != address.getServerType()) { + return false; + } + } + } catch (Throwable e) { + logger.error(e.getMessage(), e); + } + return true; + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/cache/INetworkAddressCacheDAO.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/cache/INetworkAddressCacheDAO.java index 50e758199..a4501b13c 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/cache/INetworkAddressCacheDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/cache/INetworkAddressCacheDAO.java @@ -19,6 +19,7 @@ package org.apache.skywalking.apm.collector.storage.dao.cache; import org.apache.skywalking.apm.collector.storage.base.dao.DAO; +import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress; /** * @author peng-yongsheng @@ -26,5 +27,7 @@ import org.apache.skywalking.apm.collector.storage.base.dao.DAO; public interface INetworkAddressCacheDAO extends DAO { int getAddressId(String networkAddress); - String getAddress(int addressId); + String getAddressById(int addressId); + + NetworkAddress getAddress(int addressId); } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/register/INetworkAddressRegisterDAO.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/register/INetworkAddressRegisterDAO.java index 262b45d0d..848f7adf3 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/register/INetworkAddressRegisterDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/dao/register/INetworkAddressRegisterDAO.java @@ -30,4 +30,6 @@ public interface INetworkAddressRegisterDAO extends DAO { int getMinNetworkAddressId(); void save(NetworkAddress networkAddress); + + void update(String id, int spanLayer, int serverType); } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java index b6e6c7129..74dd2930c 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddress.java @@ -20,6 +20,7 @@ package org.apache.skywalking.apm.collector.storage.table.register; import org.apache.skywalking.apm.collector.core.data.Column; import org.apache.skywalking.apm.collector.core.data.StreamData; +import org.apache.skywalking.apm.collector.core.data.operator.CoverOperation; import org.apache.skywalking.apm.collector.core.data.operator.NonOperation; /** @@ -39,7 +40,8 @@ public class NetworkAddress extends StreamData { private static final Column[] INTEGER_COLUMNS = { new Column(NetworkAddressTable.COLUMN_ADDRESS_ID, new NonOperation()), - new Column(NetworkAddressTable.COLUMN_SPAN_LAYER, new NonOperation()), + new Column(NetworkAddressTable.COLUMN_SPAN_LAYER, new CoverOperation()), + new Column(NetworkAddressTable.COLUMN_SERVER_TYPE, new CoverOperation()), }; private static final Column[] BYTE_COLUMNS = {}; @@ -87,4 +89,12 @@ public class NetworkAddress extends StreamData { public void setSpanLayer(Integer spanLayer) { setDataInteger(1, spanLayer); } + + public Integer getServerType() { + return getDataInteger(2); + } + + public void setServerType(Integer serverType) { + setDataInteger(2, serverType); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddressTable.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddressTable.java index 5374fb34b..d0e5d4735 100644 --- a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddressTable.java +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/NetworkAddressTable.java @@ -27,5 +27,6 @@ public class NetworkAddressTable extends CommonTable { public static final String TABLE = "network_address"; public static final String COLUMN_NETWORK_ADDRESS = "network_address"; public static final String COLUMN_SPAN_LAYER = "span_layer"; + public static final String COLUMN_SERVER_TYPE = "server_type"; public static final String COLUMN_ADDRESS_ID = "address_id"; } diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerType.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerType.java new file mode 100644 index 000000000..b7b71a22b --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerType.java @@ -0,0 +1,46 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.table.register; + +/** + * @author peng-yongsheng + */ +public class ServerType { + private int componentId; + private int id; + private String name; + + public ServerType(int componentId, int id, String name) { + this.componentId = componentId; + this.id = id; + this.name = name; + } + + public int getId() { + return id; + } + + public String getName() { + return name; + } + + public int getComponentId() { + return componentId; + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefine.java b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefine.java new file mode 100644 index 000000000..2810c2ab1 --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/main/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefine.java @@ -0,0 +1,82 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.table.register; + +import org.apache.skywalking.apm.collector.core.util.Const; +import org.apache.skywalking.apm.network.trace.component.ComponentsDefine; + +/** + * @author peng-yongsheng + */ +public class ServerTypeDefine { + + private static ServerTypeDefine INSTANCE = new ServerTypeDefine(); + + private String[] serverTypeNames; + private ServerType[] serverTypes; + + private ServerTypeDefine() { + this.serverTypes = new ServerType[28]; + this.serverTypeNames = new String[11]; + addServerType(new ServerType(ComponentsDefine.TOMCAT.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.HTTPCLIENT.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.DUBBO.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.H2.getId(), 1, ComponentsDefine.H2.getName())); + addServerType(new ServerType(ComponentsDefine.MYSQL.getId(), 2, ComponentsDefine.MYSQL.getName())); + addServerType(new ServerType(ComponentsDefine.ORACLE.getId(), 3, ComponentsDefine.ORACLE.getName())); + addServerType(new ServerType(ComponentsDefine.REDIS.getId(), 4, ComponentsDefine.REDIS.getName())); + addServerType(new ServerType(ComponentsDefine.MOTAN.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.MONGODB.getId(), 5, ComponentsDefine.MONGODB.getName())); + addServerType(new ServerType(ComponentsDefine.RESIN.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.FEIGN.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.OKHTTP.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.SPRING_REST_TEMPLATE.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.SPRING_MVC_ANNOTATION.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.STRUTS2.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.NUTZ_MVC_ANNOTATION.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.NUTZ_HTTP.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.JETTY_CLIENT.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.JETTY_SERVER.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.MEMCACHED.getId(), 6, ComponentsDefine.MEMCACHED.getName())); + addServerType(new ServerType(ComponentsDefine.SHARDING_JDBC.getId(), 7, ComponentsDefine.SHARDING_JDBC.getName())); + addServerType(new ServerType(ComponentsDefine.POSTGRESQL.getId(), 8, ComponentsDefine.POSTGRESQL.getName())); + addServerType(new ServerType(ComponentsDefine.GRPC.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.ELASTIC_JOB.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.ROCKET_MQ.getId(), 9, ComponentsDefine.ROCKET_MQ.getName())); + addServerType(new ServerType(ComponentsDefine.HTTP_ASYNC_CLIENT.getId(), Const.NONE, Const.EMPTY_STRING)); + addServerType(new ServerType(ComponentsDefine.KAFKA.getId(), 10, ComponentsDefine.KAFKA.getName())); + } + + public static ServerTypeDefine getInstance() { + return INSTANCE; + } + + private void addServerType(ServerType serverType) { + serverTypeNames[serverType.getId()] = serverType.getName(); + serverTypes[serverType.getComponentId()] = serverType; + } + + public int getServerTypeId(int componentId) { + return serverTypes[componentId].getId(); + } + + public String getServerType(int serverTypeId) { + return serverTypeNames[serverTypeId]; + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-define/src/test/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefineTestCase.java b/apm-collector/apm-collector-storage/collector-storage-define/src/test/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefineTestCase.java new file mode 100644 index 000000000..015eff79b --- /dev/null +++ b/apm-collector/apm-collector-storage/collector-storage-define/src/test/java/org/apache/skywalking/apm/collector/storage/table/register/ServerTypeDefineTestCase.java @@ -0,0 +1,42 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.apm.collector.storage.table.register; + +import java.lang.reflect.Field; +import org.apache.skywalking.apm.network.trace.component.ComponentsDefine; +import org.apache.skywalking.apm.network.trace.component.OfficialComponent; +import org.junit.Test; + +/** + * @author peng-yongsheng + */ +public class ServerTypeDefineTestCase { + + @Test + public void check() throws IllegalAccessException { + Field[] fields = ComponentsDefine.class.getDeclaredFields(); + + for (Field field : fields) { + if (field.getType().equals(OfficialComponent.class)) { + OfficialComponent component = (OfficialComponent)field.get(ComponentsDefine.getInstance()); + ServerTypeDefine.getInstance().getServerTypeId(component.getId()); + } + } + } +} diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cache/NetworkAddressEsCacheDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cache/NetworkAddressEsCacheDAO.java index b7315e4ba..80c500caf 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cache/NetworkAddressEsCacheDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/cache/NetworkAddressEsCacheDAO.java @@ -22,6 +22,7 @@ import org.apache.skywalking.apm.collector.client.elasticsearch.ElasticSearchCli import org.apache.skywalking.apm.collector.core.util.Const; import org.apache.skywalking.apm.collector.storage.dao.cache.INetworkAddressCacheDAO; import org.apache.skywalking.apm.collector.storage.es.base.dao.EsDAO; +import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress; import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddressTable; import org.elasticsearch.action.get.GetRequestBuilder; import org.elasticsearch.action.get.GetResponse; @@ -56,12 +57,12 @@ public class NetworkAddressEsCacheDAO extends EsDAO implements INetworkAddressCa SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); if (searchResponse.getHits().totalHits > 0) { SearchHit searchHit = searchResponse.getHits().iterator().next(); - return (int)searchHit.getSource().get(NetworkAddressTable.COLUMN_ADDRESS_ID); + return ((Number)searchHit.getSource().get(NetworkAddressTable.COLUMN_ADDRESS_ID)).intValue(); } - return 0; + return Const.NONE; } - @Override public String getAddress(int addressId) { + @Override public String getAddressById(int addressId) { logger.debug("get network address, address id: {}", addressId); ElasticSearchClient client = getClient(); GetRequestBuilder getRequestBuilder = client.prepareGet(NetworkAddressTable.TABLE, String.valueOf(addressId)); @@ -72,4 +73,21 @@ public class NetworkAddressEsCacheDAO extends EsDAO implements INetworkAddressCa } return Const.EMPTY_STRING; } + + @Override public NetworkAddress getAddress(int addressId) { + ElasticSearchClient client = getClient(); + GetRequestBuilder getRequestBuilder = client.prepareGet(NetworkAddressTable.TABLE, String.valueOf(addressId)); + + GetResponse getResponse = getRequestBuilder.get(); + if (getResponse.isExists()) { + NetworkAddress address = new NetworkAddress(); + address.setId((String)getResponse.getSource().get(NetworkAddressTable.COLUMN_ID)); + address.setAddressId(((Number)getResponse.getSource().get(NetworkAddressTable.COLUMN_ADDRESS_ID)).intValue()); + address.setSpanLayer(((Number)getResponse.getSource().get(NetworkAddressTable.COLUMN_SPAN_LAYER)).intValue()); + address.setServerType(((Number)getResponse.getSource().get(NetworkAddressTable.COLUMN_SERVER_TYPE)).intValue()); + address.setNetworkAddress((String)getResponse.getSource().get(NetworkAddressTable.COLUMN_NETWORK_ADDRESS)); + return address; + } + return null; + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/register/NetworkAddressRegisterEsDAO.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/register/NetworkAddressRegisterEsDAO.java index a59b6c66b..b979f1603 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/register/NetworkAddressRegisterEsDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/dao/register/NetworkAddressRegisterEsDAO.java @@ -56,8 +56,18 @@ public class NetworkAddressRegisterEsDAO extends EsDAO implements INetworkAddres source.put(NetworkAddressTable.COLUMN_NETWORK_ADDRESS, networkAddress.getNetworkAddress()); source.put(NetworkAddressTable.COLUMN_ADDRESS_ID, networkAddress.getAddressId()); source.put(NetworkAddressTable.COLUMN_SPAN_LAYER, networkAddress.getSpanLayer()); + source.put(NetworkAddressTable.COLUMN_SERVER_TYPE, networkAddress.getServerType()); IndexResponse response = client.prepareIndex(NetworkAddressTable.TABLE, networkAddress.getId()).setSource(source).setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE).get(); logger.debug("save network address register info, address getId: {}, network address code: {}, status: {}", networkAddress.getAddressId(), networkAddress.getNetworkAddress(), response.status().name()); } + + @Override public void update(String id, int spanLayer, int serverType) { + ElasticSearchClient client = getClient(); + + Map source = new HashMap<>(); + source.put(NetworkAddressTable.COLUMN_SPAN_LAYER, spanLayer); + source.put(NetworkAddressTable.COLUMN_SERVER_TYPE, serverType); + client.prepareUpdate(NetworkAddressTable.TABLE, id).setDoc(source).setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE).get(); + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/register/NetworkAddressEsTableDefine.java b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/register/NetworkAddressEsTableDefine.java index 59f1f4f41..ea9eaa28a 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/register/NetworkAddressEsTableDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/java/org/apache/skywalking/apm/collector/storage/es/define/register/NetworkAddressEsTableDefine.java @@ -38,5 +38,7 @@ public class NetworkAddressEsTableDefine extends ElasticSearchTableDefine { @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(NetworkAddressTable.COLUMN_NETWORK_ADDRESS, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(NetworkAddressTable.COLUMN_ADDRESS_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NetworkAddressTable.COLUMN_SPAN_LAYER, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(NetworkAddressTable.COLUMN_SERVER_TYPE, ElasticSearchColumnDefine.Type.Integer.name())); } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/cache/NetworkAddressH2CacheDAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/cache/NetworkAddressH2CacheDAO.java index d9aca53a0..99ffd1040 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/cache/NetworkAddressH2CacheDAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/cache/NetworkAddressH2CacheDAO.java @@ -26,6 +26,7 @@ import org.apache.skywalking.apm.collector.core.util.Const; import org.apache.skywalking.apm.collector.storage.base.sql.SqlBuilder; import org.apache.skywalking.apm.collector.storage.dao.cache.INetworkAddressCacheDAO; import org.apache.skywalking.apm.collector.storage.h2.base.dao.H2DAO; +import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress; import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddressTable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -47,6 +48,7 @@ public class NetworkAddressH2CacheDAO extends H2DAO implements INetworkAddressCa public int getAddressId(String networkAddress) { logger.info("get the address id with network address = {}", networkAddress); H2Client client = getClient(); + String sql = SqlBuilder.buildSql(GET_ADDRESS_ID_OR_CODE_SQL, NetworkAddressTable.COLUMN_ADDRESS_ID, NetworkAddressTable.TABLE, NetworkAddressTable.COLUMN_NETWORK_ADDRESS); Object[] params = new Object[] {networkAddress}; @@ -57,10 +59,10 @@ public class NetworkAddressH2CacheDAO extends H2DAO implements INetworkAddressCa } catch (SQLException | H2ClientException e) { logger.error(e.getMessage(), e); } - return 0; + return Const.NONE; } - @Override public String getAddress(int addressId) { + @Override public String getAddressById(int addressId) { logger.debug("get network address, address id: {}", addressId); H2Client client = getClient(); String sql = SqlBuilder.buildSql(GET_ADDRESS_ID_OR_CODE_SQL, NetworkAddressTable.COLUMN_NETWORK_ADDRESS, NetworkAddressTable.TABLE, NetworkAddressTable.COLUMN_ADDRESS_ID); @@ -74,4 +76,27 @@ public class NetworkAddressH2CacheDAO extends H2DAO implements INetworkAddressCa } return Const.EMPTY_STRING; } + + @Override public NetworkAddress getAddress(int addressId) { + logger.debug("get network address, address id: {}", addressId); + H2Client client = getClient(); + + String dynamicSql = "select * from {0} where {1} = ?"; + String sql = SqlBuilder.buildSql(dynamicSql, NetworkAddressTable.TABLE, NetworkAddressTable.COLUMN_ADDRESS_ID); + Object[] params = new Object[] {addressId}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + NetworkAddress networkAddress = new NetworkAddress(); + networkAddress.setId(rs.getString(NetworkAddressTable.COLUMN_ID)); + networkAddress.setAddressId(rs.getInt(NetworkAddressTable.COLUMN_ADDRESS_ID)); + networkAddress.setNetworkAddress(rs.getString(NetworkAddressTable.COLUMN_NETWORK_ADDRESS)); + networkAddress.setSpanLayer(rs.getInt(NetworkAddressTable.COLUMN_SPAN_LAYER)); + networkAddress.setServerType(rs.getInt(NetworkAddressTable.COLUMN_SERVER_TYPE)); + return networkAddress; + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return null; + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/register/NetworkAddressRegisterH2DAO.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/register/NetworkAddressRegisterH2DAO.java index b0b632fac..8c0afcb35 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/register/NetworkAddressRegisterH2DAO.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/dao/register/NetworkAddressRegisterH2DAO.java @@ -25,6 +25,7 @@ import org.apache.skywalking.apm.collector.client.h2.H2ClientException; import org.apache.skywalking.apm.collector.storage.base.sql.SqlBuilder; import org.apache.skywalking.apm.collector.storage.dao.register.INetworkAddressRegisterDAO; import org.apache.skywalking.apm.collector.storage.h2.base.dao.H2DAO; +import org.apache.skywalking.apm.collector.storage.table.register.InstanceTable; import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddress; import org.apache.skywalking.apm.collector.storage.table.register.NetworkAddressTable; import org.slf4j.Logger; @@ -60,6 +61,7 @@ public class NetworkAddressRegisterH2DAO extends H2DAO implements INetworkAddres source.put(NetworkAddressTable.COLUMN_NETWORK_ADDRESS, networkAddress.getNetworkAddress()); source.put(NetworkAddressTable.COLUMN_ADDRESS_ID, networkAddress.getAddressId()); source.put(NetworkAddressTable.COLUMN_SPAN_LAYER, networkAddress.getSpanLayer()); + source.put(NetworkAddressTable.COLUMN_SERVER_TYPE, networkAddress.getServerType()); String sql = SqlBuilder.buildBatchInsertSql(NetworkAddressTable.TABLE, source.keySet()); Object[] params = source.values().toArray(new Object[0]); @@ -69,4 +71,21 @@ public class NetworkAddressRegisterH2DAO extends H2DAO implements INetworkAddres logger.error(e.getMessage(), e); } } + + @Override public void update(String id, int spanLayer, int serverType) { + H2Client client = getClient(); + + Map source = new HashMap<>(); + source.put(NetworkAddressTable.COLUMN_SPAN_LAYER, spanLayer); + source.put(NetworkAddressTable.COLUMN_SERVER_TYPE, serverType); + + String sql = SqlBuilder.buildBatchUpdateSql(InstanceTable.TABLE, source.keySet(), InstanceTable.COLUMN_INSTANCE_ID); + Object[] params = source.values().toArray(new Object[] {id}); + + try { + client.execute(sql, params); + } catch (H2ClientException e) { + logger.error(e.getMessage(), e); + } + } } diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/define/register/NetworkAddressH2TableDefine.java b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/define/register/NetworkAddressH2TableDefine.java index d90c13cb3..0b21be38d 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/define/register/NetworkAddressH2TableDefine.java +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/java/org/apache/skywalking/apm/collector/storage/h2/define/register/NetworkAddressH2TableDefine.java @@ -34,6 +34,8 @@ public class NetworkAddressH2TableDefine extends H2TableDefine { @Override public void initialize() { addColumn(new H2ColumnDefine(NetworkAddressTable.COLUMN_ID, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(NetworkAddressTable.COLUMN_NETWORK_ADDRESS, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(NetworkAddressTable.COLUMN_SPAN_LAYER, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(NetworkAddressTable.COLUMN_SERVER_TYPE, H2ColumnDefine.Type.Int.name())); addColumn(new H2ColumnDefine(NetworkAddressTable.COLUMN_ADDRESS_ID, H2ColumnDefine.Type.Int.name())); } }