From cba4dc2883daf4490d2c94459ac2bb2ac5aec475 Mon Sep 17 00:00:00 2001 From: wusheng Date: Tue, 18 Jul 2017 21:47:24 +0800 Subject: [PATCH] Make sure the JVM stop hook works #295 --- .../apm/commons/datacarrier/consumer/ConsumerPool.java | 1 + .../org/skywalking/apm/agent/core/jvm/JVMService.java | 9 ++++++--- .../agent/core/remote/AppAndServiceRegisterClient.java | 5 +++++ .../agent/core/remote/CollectorDiscoveryService.java | 10 +++++++++- .../apm/agent/core/remote/GRPCChannelManager.java | 6 ++++++ .../agent/core/remote/TraceSegmentServiceClient.java | 5 +++++ .../apm/agent/core/sampling/SamplingService.java | 5 +++++ 7 files changed, 37 insertions(+), 4 deletions(-) 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/jvm/JVMService.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/jvm/JVMService.java index 3df7c1519..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 @@ -1,14 +1,11 @@ package org.skywalking.apm.agent.core.jvm; import io.grpc.ManagedChannel; -import java.text.SimpleDateFormat; -import java.util.Date; import java.util.LinkedList; import java.util.concurrent.Executors; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; -import java.util.concurrent.locks.ReentrantLock; import org.skywalking.apm.agent.core.boot.BootService; import org.skywalking.apm.agent.core.boot.ServiceManager; import org.skywalking.apm.agent.core.conf.Config; @@ -65,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