1.DataSenderFactoryWithBalance使用新的检查和发送逻辑

This commit is contained in:
wusheng 2015-12-07 22:27:10 +08:00
parent fe5e742d41
commit 437a7742a7
3 changed files with 75 additions and 99 deletions

View File

@ -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;
}
}

View File

@ -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 {

View File

@ -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() {