From 49dc8232ae4aaea56dbfe00a183e89601db2b051 Mon Sep 17 00:00:00 2001 From: wusheng Date: Mon, 7 Dec 2015 21:36:11 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BC=98=E5=8C=96=E6=A3=80=E6=9F=A5=E6=9C=BA?= =?UTF-8?q?=E5=88=B6=EF=BC=8C=E7=A7=BB=E9=99=A4NEED=5FADD=5FSENDER=5FFLAG?= =?UTF-8?q?=EF=BC=8C=E6=AD=A4=E6=A0=87=E8=AE=B0=E4=BD=8D=E5=9C=A8=E5=A4=9A?= =?UTF-8?q?=E7=BA=BF=E7=A8=8B=E9=97=B4=E5=8D=8F=E5=90=8C=EF=BC=8C=E6=9C=89?= =?UTF-8?q?=E5=8F=AF=E8=83=BD=E9=80=A0=E6=88=90=E6=95=B0=E6=8D=AE=E4=B8=8D?= =?UTF-8?q?=E4=B8=80=E8=87=B4=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../sender/DataSenderFactoryWithBalance.java | 49 +++++++++---------- 1 file changed, 22 insertions(+), 27 deletions(-) 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 5979908f3..f22238ff3 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 @@ -29,7 +29,6 @@ public class DataSenderFactoryWithBalance { private static List usingDataSender = new ArrayList(); private static int maxKeepConnectingSenderSize; - private static boolean NEED_ADD_SENDER_FLAG = false; private static int calculateMaxKeeperConnectingSenderSize(int allAddressSize) { if (CONNECT_PERCENT <= 0 || CONNECT_PERCENT > 100) { @@ -122,33 +121,30 @@ public class DataSenderFactoryWithBalance { while (true) { // 检查是否需要新增 // NEED_ADD_SENDER_FLAG 将会在unRegister方法修改值 - if (NEED_ADD_SENDER_FLAG) { - DataSender newSender; - 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) { - NEED_ADD_SENDER_FLAG = false; - break; + DataSender newSender; + 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; + } } } @@ -211,6 +207,5 @@ public class DataSenderFactoryWithBalance { usingDataSender.get(index) .setStatus(DataSender.SenderStatus.FAILED); } - NEED_ADD_SENDER_FLAG = true; } }