diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java index 092ed6920..a00570e9d 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java @@ -1,202 +1,198 @@ package com.ai.cloud.skywalking.sender; -import static com.ai.cloud.skywalking.conf.Config.Sender.CHECKER_THREAD_WAIT_INTERVAL; -import static com.ai.cloud.skywalking.conf.Config.Sender.CLOSE_SENDER_COUNTDOWN; -import static com.ai.cloud.skywalking.conf.Config.Sender.CONNECT_PERCENT; -import static com.ai.cloud.skywalking.conf.Config.Sender.RETRY_GET_SENDER_WAIT_INTERVAL; -import static com.ai.cloud.skywalking.conf.Config.Sender.SWITCH_SENDER_INTERVAL; - -import java.io.IOException; -import java.net.InetSocketAddress; -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Set; -import java.util.concurrent.ThreadLocalRandom; -import java.util.logging.Level; -import java.util.logging.Logger; - import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.util.StringUtil; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.util.*; +import java.util.concurrent.ThreadLocalRandom; +import java.util.logging.Level; +import java.util.logging.Logger; + +import static com.ai.cloud.skywalking.conf.Config.Sender.*; + public class DataSenderFactoryWithBalance { - private static Logger logger = Logger - .getLogger(DataSenderFactoryWithBalance.class.getName()); - // unUsedServerAddress存放没有使用的服务器地址, - private static List unusedServerAddresses = new ArrayList(); + private static Logger logger = Logger + .getLogger(DataSenderFactoryWithBalance.class.getName()); + // unUsedServerAddress存放没有使用的服务器地址, + private static List unusedServerAddresses = new ArrayList(); - private static List usingDataSender = new ArrayList(); - private static int maxKeepConnectingSenderSize; + private static List usingDataSender = new ArrayList(); + private static int maxKeepConnectingSenderSize; - private static int calculateMaxKeeperConnectingSenderSize(int allAddressSize) { - if (CONNECT_PERCENT <= 0 || CONNECT_PERCENT > 100) { - logger.log(Level.ALL, "CONNECT_PERCENT must between 1 and 100"); - System.exit(-1); - } - return (int) Math.ceil(allAddressSize - * ((1.0 * CONNECT_PERCENT / 100) % 100)); - } + private static int calculateMaxKeeperConnectingSenderSize(int allAddressSize) { + if (CONNECT_PERCENT <= 0 || CONNECT_PERCENT > 100) { + logger.log(Level.ALL, "CONNECT_PERCENT must between 1 and 100"); + System.exit(-1); + } + return (int) Math.ceil(allAddressSize + * ((1.0 * CONNECT_PERCENT / 100) % 100)); + } - // 初始化服务端的地址数据 - static { - // 获取数据 - if (StringUtil.isEmpty(Config.Sender.SERVERS_ADDR)) { - throw new IllegalArgumentException( - "Collection service configuration error."); - } + // 初始化服务端的地址数据 + static { + // 获取数据 + if (StringUtil.isEmpty(Config.Sender.SERVERS_ADDR)) { + throw new IllegalArgumentException( + "Collection service configuration error."); + } - // 初始化地址 - Set tmpInetSocketAddress = new HashSet(); - for (String serverConfig : Config.Sender.SERVERS_ADDR.split(";")) { - String[] server = serverConfig.split(":"); - if (server.length != 2) - throw new IllegalArgumentException( - "Collection service configuration error."); - tmpInetSocketAddress.add(new InetSocketAddress(server[0], Integer - .valueOf(server[1]))); - } + // 初始化地址 + Set tmpInetSocketAddress = new HashSet(); + for (String serverConfig : Config.Sender.SERVERS_ADDR.split(";")) { + String[] server = serverConfig.split(":"); + if (server.length != 2) + throw new IllegalArgumentException( + "Collection service configuration error."); + tmpInetSocketAddress.add(new InetSocketAddress(server[0], Integer + .valueOf(server[1]))); + } - unusedServerAddresses.addAll(tmpInetSocketAddress); + unusedServerAddresses.addAll(tmpInetSocketAddress); - // 根据配置的服务器集群的地址,来计算保持连接的Sender的数量 - maxKeepConnectingSenderSize = calculateMaxKeeperConnectingSenderSize(tmpInetSocketAddress - .size()); + // 根据配置的服务器集群的地址,来计算保持连接的Sender的数量 + maxKeepConnectingSenderSize = calculateMaxKeeperConnectingSenderSize(tmpInetSocketAddress + .size()); - // 初始化的发送程序 - int index = 0; - while (usingDataSender.size() < maxKeepConnectingSenderSize) { - index = ThreadLocalRandom.current().nextInt(0, - unusedServerAddresses.size()); - try { - usingDataSender.add(new DataSender(unusedServerAddresses - .get(index))); - unusedServerAddresses.remove(index); - } catch (IOException e) { - // 服务器连接不上 - logger.log(Level.SEVERE, "Failed to connect server[" - + unusedServerAddresses.get(index).getHostName() + "]"); - continue; - } - } + // 初始化的发送程序 + int index = 0; + while (usingDataSender.size() < maxKeepConnectingSenderSize) { + index = ThreadLocalRandom.current().nextInt(0, + unusedServerAddresses.size()); + try { + usingDataSender.add(new DataSender(unusedServerAddresses + .get(index))); + unusedServerAddresses.remove(index); + } catch (IOException e) { + // 服务器连接不上 + logger.log(Level.SEVERE, "Failed to connect server[" + + unusedServerAddresses.get(index).getHostName() + "]"); + continue; + } + } - new DataSenderChecker().start(); - } + new DataSenderChecker().start(); + } - // 获取连接 + // 获取连接 - public static DataSender getSender() { - DataSender readySender = null; - while (true) { - int index = ThreadLocalRandom.current().nextInt(0, - usingDataSender.size()); - if (usingDataSender.get(index).getStatus() == DataSender.SenderStatus.READY) { - readySender = usingDataSender.get(index); - break; - } + public static DataSender getSender() { + DataSender readySender = null; + while (true) { + int index = ThreadLocalRandom.current().nextInt(0, + usingDataSender.size()); + if (usingDataSender.get(index).getStatus() == DataSender.SenderStatus.READY) { + readySender = usingDataSender.get(index); + break; + } - if (readySender == null) { - try { - Thread.sleep(RETRY_GET_SENDER_WAIT_INTERVAL); - } catch (InterruptedException e) { - logger.log(Level.ALL, "Sleep failed"); - } - } + if (readySender == null) { + try { + Thread.sleep(RETRY_GET_SENDER_WAIT_INTERVAL); + } catch (InterruptedException e) { + logger.log(Level.ALL, "Sleep failed"); + } + } - } + } - return readySender; - } + return readySender; + } - // 定时Sender状态检查 - public static class DataSenderChecker extends Thread { - public DataSenderChecker() { - super("Data-Sender-Checker"); - } + // 定时Sender状态检查 + public static class DataSenderChecker extends Thread { + public DataSenderChecker() { + super("Data-Sender-Checker"); + } - @Override - public void run() { - long sleepTime = 0; - while (true) { - DataSender newSender; - // removing failed sender - for (int i = 0; i < usingDataSender.size(); i++) { - if (usingDataSender.get(i).getStatus() == DataSender.SenderStatus.FAILED) { - usingDataSender.get(i).close(); - unusedServerAddresses.add(usingDataSender.get(i) - .getServerIp()); - } - } + @Override + public void run() { + long sleepTime = 0; + while (true) { + DataSender newSender; + // removing failed sender + Iterator senderIterator = usingDataSender.iterator(); + DataSender tmpDataSender; + while (senderIterator.hasNext()) { + tmpDataSender = senderIterator.next(); + if (tmpDataSender.getStatus() == DataSender.SenderStatus.FAILED) { + tmpDataSender.close(); + unusedServerAddresses.add(tmpDataSender.getServerIp()); + senderIterator.remove(); + } + } - // try to fill up senders. if size is not enough. - while (usingDataSender.size() < maxKeepConnectingSenderSize) { - if ((newSender = findReadySender()) == null) { - // no available sender. ignore. - break; - } - usingDataSender.add(newSender); + // try to fill up senders. if size is not enough. + while (usingDataSender.size() < maxKeepConnectingSenderSize) { + if ((newSender = findReadySender()) == null) { + // no available sender. ignore. + break; + } + usingDataSender.add(newSender); - } + } - // try to switch. - if (sleepTime >= SWITCH_SENDER_INTERVAL) { - // if sender is enough, go to switch for balancing. - if (usingDataSender.size() >= maxKeepConnectingSenderSize) { - DataSender toBeSwitchSender; - DataSender tmpSender; - int toBeSwitchIndex = ThreadLocalRandom.current() - .nextInt(0, usingDataSender.size() - 1); - toBeSwitchSender = usingDataSender.get(toBeSwitchIndex); - tmpSender = findReadySender(); - if (tmpSender != null) { - usingDataSender.set(toBeSwitchIndex, tmpSender); - try { - Thread.sleep(CLOSE_SENDER_COUNTDOWN); - } catch (InterruptedException e) { - logger.log(Level.ALL, "Sleep Failed"); - } - toBeSwitchSender.close(); - unusedServerAddresses.remove(tmpSender - .getServerIp()); - unusedServerAddresses.add(toBeSwitchSender - .getServerIp()); - } - } - sleepTime = 0; - } + // try to switch. + if (sleepTime >= SWITCH_SENDER_INTERVAL) { + // if sender is enough, go to switch for balancing. + if (usingDataSender.size() >= maxKeepConnectingSenderSize) { + DataSender toBeSwitchSender; + DataSender tmpSender; + int toBeSwitchIndex = ThreadLocalRandom.current() + .nextInt(0, usingDataSender.size() - 1); + toBeSwitchSender = usingDataSender.get(toBeSwitchIndex); + tmpSender = findReadySender(); + if (tmpSender != null) { + usingDataSender.set(toBeSwitchIndex, tmpSender); + try { + Thread.sleep(CLOSE_SENDER_COUNTDOWN); + } catch (InterruptedException e) { + logger.log(Level.ALL, "Sleep Failed"); + } + toBeSwitchSender.close(); + unusedServerAddresses.remove(tmpSender + .getServerIp()); + unusedServerAddresses.add(toBeSwitchSender + .getServerIp()); + } + } + sleepTime = 0; + } - sleepTime += CHECKER_THREAD_WAIT_INTERVAL; - try { - Thread.sleep(CHECKER_THREAD_WAIT_INTERVAL); - } catch (InterruptedException e) { - logger.log(Level.ALL, "Sleep failed"); - } + sleepTime += CHECKER_THREAD_WAIT_INTERVAL; + try { + Thread.sleep(CHECKER_THREAD_WAIT_INTERVAL); + } catch (InterruptedException e) { + logger.log(Level.ALL, "Sleep failed"); + } - } - } - } + } + } + } - private static DataSender findReadySender() { - DataSender result = null; - for (InetSocketAddress serverAddress : unusedServerAddresses) { - try { - result = new DataSender(serverAddress); - break; - } catch (IOException e) { - if (result != null) { - result.close(); - } - continue; - } - } - return result; - } + private static DataSender findReadySender() { + DataSender result = null; + for (InetSocketAddress serverAddress : unusedServerAddresses) { + try { + result = new DataSender(serverAddress); + break; + } catch (IOException e) { + if (result != null) { + result.close(); + } + continue; + } + } + return result; + } - public static void unRegister(DataSender socket) { - int index = usingDataSender.indexOf(socket); - if (index != -1) { - usingDataSender.get(index) - .setStatus(DataSender.SenderStatus.FAILED); - } - } + public static void unRegister(DataSender socket) { + int index = usingDataSender.indexOf(socket); + if (index != -1) { + usingDataSender.get(index) + .setStatus(DataSender.SenderStatus.FAILED); + } + } }