diff --git a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java index 75a2ad651..e6d4e7602 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/consumer/ConsumerPool.java @@ -20,6 +20,7 @@ public class ConsumerPool { this(channels, num); for (int i = 0; i < num; i++) { consumerThreads[i] = new ConsumerThread("DataCarrier.Consumser." + i + ".Thread", getNewConsumerInstance(consumerClass)); + consumerThreads[i].setDaemon(true); } } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/BootService.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/BootService.java index d3538f915..be80b2c86 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/BootService.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/BootService.java @@ -13,4 +13,6 @@ public interface BootService { void boot() throws Throwable; void afterBoot() throws Throwable; + + void shutdown() throws Throwable; } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/ServiceManager.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/ServiceManager.java index 2d7338994..2c5337ef2 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/ServiceManager.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/boot/ServiceManager.java @@ -27,6 +27,16 @@ public enum ServiceManager { afterBoot(); } + public void shutdown() { + for (BootService service : bootedServices.values()) { + try { + service.shutdown(); + } catch (Throwable e) { + logger.error(e, "ServiceManager try to shutdown [{}] fail.", service.getClass().getName()); + } + } + } + private Map loadAllServices() { HashMap bootedServices = new HashMap(); Iterator serviceIterator = load().iterator(); diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/ContextManager.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/ContextManager.java index a34458464..f6347dca0 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/ContextManager.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/ContextManager.java @@ -162,6 +162,10 @@ public class ContextManager implements TracingContextListener, BootService, Igno } + @Override public void shutdown() throws Throwable { + + } + @Override public void afterFinished(TraceSegment traceSegment) { CONTEXT.remove(); diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/jvm/JVMService.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/jvm/JVMService.java index db80d7969..6ac9a0a2b 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/jvm/JVMService.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/jvm/JVMService.java @@ -62,6 +62,12 @@ public class JVMService implements BootService, Runnable { } + @Override + public void shutdown() throws Throwable { + collectMetricFuture.cancel(true); + sendMetricFuture.cancel(true); + } + @Override public void run() { if (RemoteDownstreamConfig.Agent.APPLICATION_ID != DictionaryUtil.nullValue() diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/AppAndServiceRegisterClient.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/AppAndServiceRegisterClient.java index af24fc93d..c32ed0ddb 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/AppAndServiceRegisterClient.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/AppAndServiceRegisterClient.java @@ -77,6 +77,11 @@ public class AppAndServiceRegisterClient implements BootService, GRPCChannelList TracingContext.ListenerManager.add(this); } + @Override + public void shutdown() throws Throwable { + applicationRegisterFuture.cancel(true); + } + @Override public void run() { if (CONNECTED.equals(status)) { diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java index 76461b985..0b73added 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java @@ -1,6 +1,7 @@ package org.skywalking.apm.agent.core.remote; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import org.skywalking.apm.agent.core.boot.BootService; import org.skywalking.apm.agent.core.conf.Config; @@ -11,6 +12,8 @@ import org.skywalking.apm.agent.core.conf.Config; * @author wusheng */ public class CollectorDiscoveryService implements BootService { + private ScheduledFuture future; + @Override public void beforeBoot() throws Throwable { @@ -18,7 +21,7 @@ public class CollectorDiscoveryService implements BootService { @Override public void boot() throws Throwable { - Executors.newSingleThreadScheduledExecutor() + future = Executors.newSingleThreadScheduledExecutor() .scheduleAtFixedRate(new DiscoveryRestServiceClient(), 0, Config.Collector.DISCOVERY_CHECK_INTERVAL, TimeUnit.SECONDS); } @@ -27,4 +30,9 @@ public class CollectorDiscoveryService implements BootService { public void afterBoot() throws Throwable { } + + @Override + public void shutdown() throws Throwable { + future.cancel(true); + } } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/GRPCChannelManager.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/GRPCChannelManager.java index 6042c55bf..a56525e3a 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/GRPCChannelManager.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/GRPCChannelManager.java @@ -49,6 +49,12 @@ public class GRPCChannelManager implements BootService, Runnable { } + @Override + public void shutdown() throws Throwable { + connectCheckFuture.cancel(true); + managedChannel.shutdownNow(); + } + @Override public void run() { if (reconnect) { diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java index d04a625f1..288fa7e4f 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java @@ -49,6 +49,11 @@ public class TraceSegmentServiceClient implements BootService, IConsumer transform(DynamicType.Builder builder, TypeDescription typeDescription,