From 4399299bd56f39b5991109e5bd1f74383688b896 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=90=B4=E6=99=9F=20Wu=20Sheng?= Date: Thu, 24 Jan 2019 00:09:09 +0800 Subject: [PATCH] Optimize the agent grpc channel management. (#2198) * Optimize the grpc channel management. * Rename old version concepts * Notify when rebuild connection. * Remove notify in new channel created. * Remove println. --- .../apm/agent/core/conf/Config.java | 2 +- .../agent/core/remote/GRPCChannelManager.java | 29 ++++--- .../core/remote/GRPCChannelManagerTest.java | 79 +++++++++++++++++++ 3 files changed, 97 insertions(+), 13 deletions(-) create mode 100644 apm-sniffer/apm-agent-core/src/test/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManagerTest.java diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/conf/Config.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/conf/Config.java index a3a049e18..e7c3fb968 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/conf/Config.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/conf/Config.java @@ -88,7 +88,7 @@ public class Config { */ public static long GRPC_CHANNEL_CHECK_INTERVAL = 30; /** - * application and service registry check interval + * service and endpoint registry check interval */ public static long APP_AND_SERVICE_REGISTER_CHECK_INTERVAL = 3; /** diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java index 4334f4527..b5ddd84f9 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManager.java @@ -39,6 +39,7 @@ public class GRPCChannelManager implements BootService, Runnable { private Random random = new Random(); private List listeners = Collections.synchronizedList(new LinkedList()); private volatile List grpcServers; + private volatile int selectedIdx = -1; @Override public void prepare() throws Throwable { @@ -87,25 +88,29 @@ public class GRPCChannelManager implements BootService, Runnable { String server = ""; try { int index = Math.abs(random.nextInt()) % grpcServers.size(); - server = grpcServers.get(index); - String[] ipAndPort = server.split(":"); + if (index != selectedIdx) { + selectedIdx = index; - managedChannel = GRPCChannel.newBuilder(ipAndPort[0], Integer.parseInt(ipAndPort[1])) - .addManagedChannelBuilder(new StandardChannelBuilder()) - .addManagedChannelBuilder(new TLSChannelBuilder()) - .addChannelDecorator(new AuthenticationDecorator()) - .build(); + server = grpcServers.get(index); + String[] ipAndPort = server.split(":"); + + if (managedChannel != null) { + managedChannel.shutdownNow(); + } + + managedChannel = GRPCChannel.newBuilder(ipAndPort[0], Integer.parseInt(ipAndPort[1])) + .addManagedChannelBuilder(new StandardChannelBuilder()) + .addManagedChannelBuilder(new TLSChannelBuilder()) + .addChannelDecorator(new AuthenticationDecorator()) + .build(); - if (!managedChannel.isShutdown() && !managedChannel.isTerminated()) { - reconnect = false; notify(GRPCChannelStatus.CONNECTED); - } else { - notify(GRPCChannelStatus.DISCONNECT); } + + reconnect = false; return; } catch (Throwable t) { logger.error(t, "Create channel to {} fail.", server); - notify(GRPCChannelStatus.DISCONNECT); } } diff --git a/apm-sniffer/apm-agent-core/src/test/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManagerTest.java b/apm-sniffer/apm-agent-core/src/test/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManagerTest.java new file mode 100644 index 000000000..374def0bf --- /dev/null +++ b/apm-sniffer/apm-agent-core/src/test/java/org/apache/skywalking/apm/agent/core/remote/GRPCChannelManagerTest.java @@ -0,0 +1,79 @@ +/* + * 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.agent.core.remote; + +import io.grpc.Server; +import io.grpc.netty.NettyServerBuilder; +import io.grpc.stub.StreamObserver; +import java.net.InetSocketAddress; +import org.apache.skywalking.apm.agent.core.conf.Config; +import org.apache.skywalking.apm.network.register.v2.*; +import org.junit.*; + +public class GRPCChannelManagerTest { + @BeforeClass + public static void setup() { + Config.Collector.BACKEND_SERVICE = "127.0.0.1:8080"; + } + + @AfterClass + public static void clear() { + Config.Collector.BACKEND_SERVICE = ""; + } + + //@Test + public void testConnected() throws Throwable { + GRPCChannelManager manager = new GRPCChannelManager(); + manager.addChannelListener(new GRPCChannelListener() { + @Override public void statusChanged(GRPCChannelStatus status) { + } + }); + + manager.boot(); + Thread.sleep(1000); + + RegisterGrpc.RegisterBlockingStub stub = RegisterGrpc.newBlockingStub(manager.getChannel()); + try { + stub.doServiceRegister(Services.newBuilder().addServices(Service.newBuilder().setServiceName("abc")).build()); + } catch (Exception e) { + e.printStackTrace(); + } + + NettyServerBuilder nettyServerBuilder = NettyServerBuilder.forAddress(new InetSocketAddress("127.0.0.1", 8080)); + nettyServerBuilder.addService(new RegisterGrpc.RegisterImplBase() { + @Override + public void doServiceRegister(Services request, StreamObserver responseObserver) { + } + }); + Server server = nettyServerBuilder.build(); + server.start(); + Thread.sleep(1000); + + boolean registerSuccess = false; + try { + stub.doServiceRegister(Services.newBuilder().addServices(Service.newBuilder().setServiceName("abc")).build()); + registerSuccess = true; + } catch (Exception e) { + e.printStackTrace(); + } + + Assert.assertTrue(registerSuccess); + server.shutdownNow(); + } +}