From 1b83b34401be7a26aee6e7300ed74dec589909d0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=90=B4=E6=99=9F=20Wu=20Sheng?= Date: Sun, 10 Jun 2018 22:00:49 +0800 Subject: [PATCH] Relied on grpc response for certain, to avoid OOM when collector is overloaded. (#1323) --- .../core/remote/GRPCStreamServiceStatus.java | 24 ++++++++++++++++++- .../remote/TraceSegmentServiceClient.java | 9 +++---- 2 files changed, 26 insertions(+), 7 deletions(-) diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCStreamServiceStatus.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCStreamServiceStatus.java index c6e90d5c8..56a0d798f 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCStreamServiceStatus.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCStreamServiceStatus.java @@ -16,13 +16,16 @@ * */ - package org.apache.skywalking.apm.agent.core.remote; +import org.apache.skywalking.apm.agent.core.logging.api.ILog; +import org.apache.skywalking.apm.agent.core.logging.api.LogManager; + /** * @author wusheng */ public class GRPCStreamServiceStatus { + private static final ILog logger = LogManager.getLogger(GRPCStreamServiceStatus.class); private volatile boolean status; public GRPCStreamServiceStatus(boolean status) { @@ -52,6 +55,25 @@ public class GRPCStreamServiceStatus { return status; } + /** + * Wait until success status reported. + */ + public void wait4Finish() { + long recheckCycle = 5; + long hasWaited = 0L; + long maxCycle = 30 * 1000L;// 30 seconds max. + while (!status) { + try2Sleep(recheckCycle); + hasWaited += recheckCycle; + + if (recheckCycle >= maxCycle) { + logger.warn("Collector traceSegment service doesn't response in {} seconds.", hasWaited); + } else { + recheckCycle = recheckCycle * 2 > maxCycle ? maxCycle : recheckCycle * 2; + } + } + } + /** * Try to sleep, and ignore the {@link InterruptedException} * diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java index de8db2d46..8fe211ec4 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java @@ -16,11 +16,11 @@ * */ - package org.apache.skywalking.apm.agent.core.remote; import io.grpc.Channel; import io.grpc.stub.StreamObserver; +import java.util.List; import org.apache.skywalking.apm.agent.core.boot.BootService; import org.apache.skywalking.apm.agent.core.boot.DefaultImplementor; import org.apache.skywalking.apm.agent.core.boot.ServiceManager; @@ -36,8 +36,6 @@ import org.apache.skywalking.apm.network.proto.Downstream; import org.apache.skywalking.apm.network.proto.TraceSegmentServiceGrpc; import org.apache.skywalking.apm.network.proto.UpstreamSegment; -import java.util.List; - import static org.apache.skywalking.apm.agent.core.conf.Config.Buffer.BUFFER_SIZE; import static org.apache.skywalking.apm.agent.core.conf.Config.Buffer.CHANNEL_SIZE; import static org.apache.skywalking.apm.agent.core.remote.GRPCChannelStatus.CONNECTED; @@ -122,9 +120,8 @@ public class TraceSegmentServiceClient implements BootService, IConsumer