diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java index 2b48ecbec..358224105 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java @@ -96,7 +96,7 @@ public abstract class ModuleProvider { Service service) throws ServiceNotProvidedException { if (serviceType.isInstance(service)) { if (manager.isServiceInstrument()) { - service = ServiceInstrumentation.INSTANCE.buildServiceUnderMonitor(service); + service = ServiceInstrumentation.INSTANCE.buildServiceUnderMonitor(module.name(), name(), service); } this.services.put(serviceType, service); } else { diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/MetricCollector.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/MetricCollector.java index f0053ee27..00c8b1617 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/MetricCollector.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/MetricCollector.java @@ -18,11 +18,154 @@ package org.skywalking.apm.collector.core.module.instrument; +import java.lang.reflect.Method; +import java.util.HashMap; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + /** * The MetricCollector collects the service metrics by Module/Provider/Service structure. */ -public enum MetricCollector { +public enum MetricCollector implements Runnable { INSTANCE; + private final Logger logger = LoggerFactory.getLogger(MetricCollector.class); + private HashMap modules = new HashMap<>(); + MetricCollector() { + ScheduledExecutorService service = Executors + .newSingleThreadScheduledExecutor(); + service.scheduleAtFixedRate(this, 10, 60, TimeUnit.SECONDS); + } + + @Override + public void run() { + if (!logger.isDebugEnabled()) { + return; + } + StringBuilder report = new StringBuilder(); + report.append("\n"); + report.append("##################################################################################################################\n"); + report.append("# Collector Service Report #\n"); + report.append("##################################################################################################################\n"); + modules.forEach((moduleName, moduleMetric) -> { + report.append(moduleName).append(":\n"); + moduleMetric.providers.forEach((providerName, providerMetric) -> { + report.append("\t").append(providerName).append(":\n"); + providerMetric.services.forEach((serviceName, serviceMetric) -> { + serviceMetric.methodMetrics.forEach((method, metric) -> { + report.append("\t\t").append(method).append(":\n"); + report.append("\t\t\t").append(metric).append("\n"); + serviceMetric.methodMetrics.put(method, new ServiceMethodMetric()); + }); + }); + }); + }); + + logger.debug(report.toString()); + + } + + ServiceMetric registerService(String module, String provider, String service) { + return initIfAbsent(module).initIfAbsent(provider).initIfAbsent(service); + } + + private ModuleMetric initIfAbsent(String moduleName) { + if (!modules.containsKey(moduleName)) { + ModuleMetric metric = new ModuleMetric(moduleName); + modules.put(moduleName, metric); + return metric; + } + return modules.get(moduleName); + } + + private class ModuleMetric { + private String moduleName; + private HashMap providers = new HashMap<>(); + + public ModuleMetric(String moduleName) { + this.moduleName = moduleName; + } + + private ProviderMetric initIfAbsent(String providerName) { + if (!providers.containsKey(providerName)) { + ProviderMetric metric = new ProviderMetric(providerName); + providers.put(providerName, metric); + return metric; + } + return providers.get(providerName); + } + } + + private class ProviderMetric { + private String providerName; + private HashMap services = new HashMap<>(); + + public ProviderMetric(String providerName) { + this.providerName = providerName; + } + + private ServiceMetric initIfAbsent(String serviceName) { + if (!services.containsKey(serviceName)) { + ServiceMetric metric = new ServiceMetric(serviceName); + services.put(serviceName, metric); + return metric; + } + return services.get(serviceName); + } + } + + class ServiceMetric { + private String serviceName; + private ConcurrentHashMap methodMetrics = new ConcurrentHashMap<>(); + + public ServiceMetric(String serviceName) { + this.serviceName = serviceName; + } + + void trace(Method method, long nano, boolean occurException) { + if (logger.isDebugEnabled()) { + ServiceMethodMetric metric = methodMetrics.get(method); + if (metric == null) { + ServiceMethodMetric methodMetric = new ServiceMethodMetric(); + methodMetrics.putIfAbsent(method, methodMetric); + metric = methodMetrics.get(method); + } + metric.add(nano, occurException); + } + } + } + + private class ServiceMethodMetric { + private AtomicLong totalTimeNano; + private AtomicLong counter; + private AtomicLong errorCounter; + + public ServiceMethodMetric() { + totalTimeNano = new AtomicLong(0); + counter = new AtomicLong(0); + errorCounter = new AtomicLong(0); + } + + private void add(long nano, boolean occurException) { + totalTimeNano.addAndGet(nano); + counter.incrementAndGet(); + if (occurException) + errorCounter.incrementAndGet(); + } + + @Override public String toString() { + if (counter.longValue() == 0) { + return "Avg=N/A"; + } + return "Avg=" + (totalTimeNano.longValue() / counter.longValue()) + " (nano)" + + ", Success Rate=" + (counter.longValue() - errorCounter.longValue()) * 100 / counter.longValue() + + "%"; + } + } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceInstrumentation.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceInstrumentation.java index a33b1b695..a41a6232d 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceInstrumentation.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceInstrumentation.java @@ -42,7 +42,7 @@ public enum ServiceInstrumentation { private final Logger logger = LoggerFactory.getLogger(ServiceInstrumentation.class); private ElementMatcher excludeObjectMethodsMatcher; - public Service buildServiceUnderMonitor(Service implementation) { + public Service buildServiceUnderMonitor(String moduleName, String providerName, Service implementation) { if (implementation instanceof TracedService) { // Duplicate service instrument, ignore. return implementation; @@ -51,7 +51,7 @@ public enum ServiceInstrumentation { return new ByteBuddy().subclass(implementation.getClass()) .implement(TracedService.class) .method(getDefaultMatcher()).intercept( - MethodDelegation.withDefaultConfiguration().to(new ServiceMetricCollector()) + MethodDelegation.withDefaultConfiguration().to(new ServiceMetricTracing(moduleName, providerName, implementation.getClass().getName())) ).make().load(getClass().getClassLoader() ).getLoaded().newInstance(); } catch (InstantiationException e) { diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceMetricCollector.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceMetricTracing.java similarity index 68% rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceMetricCollector.java rename to apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceMetricTracing.java index 2da824857..1f99ffc94 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceMetricCollector.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/instrument/ServiceMetricTracing.java @@ -29,7 +29,12 @@ import net.bytebuddy.implementation.bind.annotation.This; /** * @author wu-sheng */ -public class ServiceMetricCollector { +public class ServiceMetricTracing { + private MetricCollector.ServiceMetric serviceMetric; + + public ServiceMetricTracing(String module, String provider, String service) { + serviceMetric = MetricCollector.INSTANCE.registerService(module, provider, service); + } @RuntimeType public Object intercept(@This Object obj, @@ -37,6 +42,17 @@ public class ServiceMetricCollector { @SuperCall Callable zuper, @Origin Method method ) throws Throwable { - return zuper.call(); + boolean occurError = false; + long startNano = System.nanoTime(); + long endNano; + try { + return zuper.call(); + } catch (Throwable t) { + occurError = true; + throw t; + } finally { + endNano = System.nanoTime(); + serviceMetric.trace(method, endNano - startNano, occurError); + } } } diff --git a/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/module/ModuleManagerTest.java b/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/module/ModuleManagerTest.java index d31164d86..5965e1025 100644 --- a/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/module/ModuleManagerTest.java +++ b/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/module/ModuleManagerTest.java @@ -54,5 +54,23 @@ public class ModuleManagerTest { BaseModuleA.ServiceABusiness1 serviceABusiness1 = manager.find("BaseA").getService(BaseModuleA.ServiceABusiness1.class); Assert.assertTrue(serviceABusiness1 instanceof TracedService); + +// for (int i = 0; i < 10000; i++) +// serviceABusiness1.print(); +// +// try { +// Thread.sleep(60 * 1000L); +// } catch (InterruptedException e) { +// e.printStackTrace(); +// } +// +// for (int i = 0; i < 10000; i++) +// serviceABusiness1.print(); +// +// try { +// Thread.sleep(120 * 1000L); +// } catch (InterruptedException e) { +// e.printStackTrace(); +// } } }