From 437a7742a754ad6632585d3c88400f95bbbda8df Mon Sep 17 00:00:00 2001 From: wusheng Date: Mon, 7 Dec 2015 22:27:10 +0800 Subject: [PATCH] =?UTF-8?q?1.DataSenderFactoryWithBalance=E4=BD=BF?= =?UTF-8?q?=E7=94=A8=E6=96=B0=E7=9A=84=E6=A3=80=E6=9F=A5=E5=92=8C=E5=8F=91?= =?UTF-8?q?=E9=80=81=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../com/ai/cloud/skywalking/conf/Config.java | 93 +++++++++---------- .../cloud/skywalking/sender/DataSender.java | 11 +-- .../sender/DataSenderFactoryWithBalance.java | 70 +++++++------- 3 files changed, 75 insertions(+), 99 deletions(-) diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java index f293ce2b9..efe946620 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java @@ -2,71 +2,64 @@ package com.ai.cloud.skywalking.conf; public class Config { - public static class SkyWalking { - public static String USER_ID; + public static class SkyWalking { + public static String USER_ID; - public static String APPLICATION_ID; - } + public static String APPLICATION_ID; + } - public static class BuriedPoint { - // 是否打印埋点信息 - public static boolean PRINTF = false; + public static class BuriedPoint { + // 是否打印埋点信息 + public static boolean PRINTF = false; - public static int MAX_EXCEPTION_STACK_LENGTH = 4000; + public static int MAX_EXCEPTION_STACK_LENGTH = 4000; - // Business Key 最大长度 - public static int BUSINESSKEY_MAX_LENGTH = 300; - } + // Business Key 最大长度 + public static int BUSINESSKEY_MAX_LENGTH = 300; + } - public static class Consumer { - // 最大消费线程数 - public static int MAX_CONSUMER = 2; - // 消费者最大等待时间 - public static long MAX_WAIT_TIME = 5L; + public static class Consumer { + // 最大消费线程数 + public static int MAX_CONSUMER = 2; + // 消费者最大等待时间 + public static long MAX_WAIT_TIME = 5L; - // - public static long CONSUMER_FAIL_RETRY_WAIT_INTERVAL = 50L; - } + // + public static long CONSUMER_FAIL_RETRY_WAIT_INTERVAL = 50L; + } - public static class Buffer { - // 每个Buffer的最大个数 - public static int BUFFER_MAX_SIZE = 20000; + public static class Buffer { + // 每个Buffer的最大个数 + public static int BUFFER_MAX_SIZE = 20000; - // Buffer池的最大长度 - public static int POOL_SIZE = 5; - } + // Buffer池的最大长度 + public static int POOL_SIZE = 5; + } - public static class Sender { - // 最大发送者的连接数阀比例 - public static int CONNECT_PERCENT = 50; + public static class Sender { + // 最大发送者的连接数阀比例 + public static int CONNECT_PERCENT = 50; - // 发送服务端配置 - public static String SERVERS_ADDR = "127.0.0.1:34000;127.0.0.1:34001;127.0.0.1:34002"; + // 发送服务端配置 + public static String SERVERS_ADDR = "127.0.0.1:34000;127.0.0.1:34001;127.0.0.1:34002"; - // 是否开启发送 - public static boolean IS_OFF = false; + // 是否开启发送 + public static boolean IS_OFF = false; - // 发送的最大长度 - public static int MAX_SEND_LENGTH = 18500; + // 发送的最大长度 + public static int MAX_SEND_LENGTH = 18500; - public static long RETRY_GET_SENDER_WAIT_INTERVAL = 2000L; + public static long RETRY_GET_SENDER_WAIT_INTERVAL = 2000L; - //切换Sender的周期 - public static long SWITCH_SENDER_INTERVAL = 10 * 60 * 1000; - //public static long SWITCH_SENDER_INTERVAL = 10 * 1000; + // 切换Sender的周期 + public static long SWITCH_SENDER_INTERVAL = 10 * 60 * 1000; - // 切换Sender之后,关闭Sender的倒计时 - public static long CLOSE_SENDER_COUNTDOWN = 3 * 1000; + // 切换Sender之后,关闭Sender的倒计时 + public static long CLOSE_SENDER_COUNTDOWN = 10 * 1000; - // Checker线程处理完成等待周期 - public static long CHECKER_THREAD_WAIT_INTERVAL = 1000; + // Checker线程处理完成等待周期 + public static long CHECKER_THREAD_WAIT_INTERVAL = 1000; - public static long RETRY_FIND_CONNECTION_SENDER = 1000; - } - - public static class SenderChecker { - - // 检查周期时间 - public static long CHECK_POLLING_TIME = 200L; - } + public static long RETRY_FIND_CONNECTION_SENDER = 1000; + } } \ No newline at end of file diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java index 63f2f55ba..bb3a9dc69 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java @@ -21,16 +21,7 @@ public class DataSender { private SenderStatus status = SenderStatus.FAILED; public DataSender(String ip, int port) throws IOException { - selector = Selector.open(); - InetSocketAddress isa = new InetSocketAddress(ip, port); - //调用open的静态方法创建连接指定的主机的SocketChannel - socketChannel = SocketChannel.open(isa); - //设置该sc已非阻塞的方式工作 - socketChannel.configureBlocking(false); - socketChannel.register(selector, SelectionKey.OP_CONNECT); - socketChannel.setOption(StandardSocketOptions.SO_KEEPALIVE, true); - this.socketAddress = isa; - status = SenderStatus.READY; + this(new InetSocketAddress(ip, port)); } public DataSender(InetSocketAddress address) throws IOException { 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 8fe559def..092ed6920 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 @@ -3,7 +3,6 @@ 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_FIND_CONNECTION_SENDER; 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; @@ -119,59 +118,53 @@ public class DataSenderFactoryWithBalance { 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(); - // 正在使用的Sender的数量 <= maxKeepConnectingSenderSize - // 剩余的服务器地址数量 = 总得服务器地址数量 - 正在使用的Sender的数量 - // 可替换的服务器数量 = 剩余服务器地址数量 - // 当剩余服务器地址数量 <= 0 - // 时,可以替换的地址也不存在,替换操作就可以不执行,所以这里的while是这样的意思 - // 当剩余服务器地址数量 > 0 时, - // 就可以找到可以替换的地址,替换操作也就可以执行了,这里的就会跳出while循环 - while ((newSender = findReadySender()) == null) { - try { - Thread.sleep(RETRY_FIND_CONNECTION_SENDER); - } catch (InterruptedException e) { - logger.log(Level.ALL, "Sleep failed."); - } - } - usingDataSender.set(i, newSender); unusedServerAddresses.add(usingDataSender.get(i) .getServerIp()); - if (usingDataSender.size() >= maxKeepConnectingSenderSize) { - break; - } } + } + + // 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) { - 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"); + // 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()); } - toBeSwitchSender.close(); - unusedServerAddresses.remove(tmpSender.getServerIp()); - unusedServerAddresses.add(toBeSwitchSender - .getServerIp()); } sleepTime = 0; } - // sleepTime += CHECKER_THREAD_WAIT_INTERVAL; try { Thread.sleep(CHECKER_THREAD_WAIT_INTERVAL); @@ -181,7 +174,6 @@ public class DataSenderFactoryWithBalance { } } - } private static DataSender findReadySender() {