1.修复没有移除失败的Sender

This commit is contained in:
zhangxin10 2015-12-07 23:30:03 +08:00
parent 437a7742a7
commit 3805cc6321
1 changed files with 168 additions and 172 deletions

View File

@ -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<InetSocketAddress> unusedServerAddresses = new ArrayList<InetSocketAddress>();
private static Logger logger = Logger
.getLogger(DataSenderFactoryWithBalance.class.getName());
// unUsedServerAddress存放没有使用的服务器地址
private static List<InetSocketAddress> unusedServerAddresses = new ArrayList<InetSocketAddress>();
private static List<DataSender> usingDataSender = new ArrayList<DataSender>();
private static int maxKeepConnectingSenderSize;
private static List<DataSender> usingDataSender = new ArrayList<DataSender>();
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<InetSocketAddress> tmpInetSocketAddress = new HashSet<InetSocketAddress>();
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<InetSocketAddress> tmpInetSocketAddress = new HashSet<InetSocketAddress>();
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<DataSender> 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);
}
}
}