1.修改SpanType为字符串,修改webui,提高类型的扩展性
2.为后续增加DataSender支持副本发送,增加部分代码。功能暂未完成。 3.为DataSenderFactoryWithBalance.getSender方法增加异常机制,方法异常退出。 4.修复客户端大量线程缺少异常保障能力的问题
This commit is contained in:
parent
7f904eee32
commit
cd60b5adf7
|
|
@ -10,16 +10,19 @@ buriedpoint.max_exception_stack_length=4000
|
|||
#业务字段的最大长度
|
||||
buriedpoint.businesskey_max_length=300
|
||||
|
||||
#发送的最大长度
|
||||
sender.max_send_length=20000
|
||||
#最大发送者的连接数阀比例
|
||||
sender.connect_percent=100
|
||||
#发送服务端配置
|
||||
sender.servers_addr=127.0.0.1:34000
|
||||
#最大发送的副本数量
|
||||
sender.max_copy_num=2
|
||||
#发送的最大长度
|
||||
sender.max_send_length=20000
|
||||
#当没有Sender时,尝试获取sender的等待周期 +++++
|
||||
sender.retry_get_sender_wait_interval=2000
|
||||
#是否开启发送消息
|
||||
sender.is_off=false
|
||||
#发送服务端配置 +++++++
|
||||
sender.servers_addr=127.0.0.1:34000
|
||||
|
||||
|
||||
#最大消费线程数
|
||||
consumer.max_consumer=2
|
||||
|
|
|
|||
|
|
@ -1,6 +1,5 @@
|
|||
package com.ai.cloud.skywalking.buffer;
|
||||
|
||||
|
||||
import static com.ai.cloud.skywalking.conf.Config.Buffer.BUFFER_MAX_SIZE;
|
||||
import static com.ai.cloud.skywalking.conf.Config.Consumer.CONSUMER_FAIL_RETRY_WAIT_INTERVAL;
|
||||
import static com.ai.cloud.skywalking.conf.Config.Consumer.MAX_CONSUMER;
|
||||
|
|
@ -18,95 +17,107 @@ import com.ai.cloud.skywalking.selfexamination.HeathReading;
|
|||
import com.ai.cloud.skywalking.sender.DataSenderFactoryWithBalance;
|
||||
|
||||
public class BufferGroup {
|
||||
private static Logger logger = Logger.getLogger(BufferGroup.class.getName());
|
||||
private String groupName;
|
||||
private Span[] dataBuffer = new Span[BUFFER_MAX_SIZE];
|
||||
AtomicInteger index = new AtomicInteger(0);
|
||||
private static Logger logger = Logger
|
||||
.getLogger(BufferGroup.class.getName());
|
||||
private String groupName;
|
||||
private Span[] dataBuffer = new Span[BUFFER_MAX_SIZE];
|
||||
AtomicInteger index = new AtomicInteger(0);
|
||||
|
||||
public BufferGroup(String groupName) {
|
||||
this.groupName = groupName;
|
||||
public BufferGroup(String groupName) {
|
||||
this.groupName = groupName;
|
||||
|
||||
int step = (int) Math.ceil(BUFFER_MAX_SIZE * 1.0 / MAX_CONSUMER);
|
||||
int start = 0, end = 0;
|
||||
while (true) {
|
||||
if (end + step >= BUFFER_MAX_SIZE) {
|
||||
new ConsumerWorker(start, BUFFER_MAX_SIZE).start();
|
||||
break;
|
||||
}
|
||||
end += step;
|
||||
new ConsumerWorker(start, end).start();
|
||||
start = end;
|
||||
}
|
||||
}
|
||||
int step = (int) Math.ceil(BUFFER_MAX_SIZE * 1.0 / MAX_CONSUMER);
|
||||
int start = 0, end = 0;
|
||||
while (true) {
|
||||
if (end + step >= BUFFER_MAX_SIZE) {
|
||||
new ConsumerWorker(start, BUFFER_MAX_SIZE).start();
|
||||
break;
|
||||
}
|
||||
end += step;
|
||||
new ConsumerWorker(start, end).start();
|
||||
start = end;
|
||||
}
|
||||
}
|
||||
|
||||
public void save(Span span) {
|
||||
int i = Math.abs(index.getAndIncrement() % BUFFER_MAX_SIZE);
|
||||
if (dataBuffer[i] != null) {
|
||||
HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.WARNING, "Group[" + groupName + "] index[" + i + "] data collision, discard old data.");
|
||||
}
|
||||
dataBuffer[i] = span;
|
||||
}
|
||||
public void save(Span span) {
|
||||
int i = Math.abs(index.getAndIncrement() % BUFFER_MAX_SIZE);
|
||||
if (dataBuffer[i] != null) {
|
||||
HealthCollector.getCurrentHeathReading(null).updateData(
|
||||
HeathReading.WARNING,
|
||||
"Group[" + groupName + "] index[" + i
|
||||
+ "] data collision, discard old data.");
|
||||
}
|
||||
dataBuffer[i] = span;
|
||||
}
|
||||
|
||||
class ConsumerWorker extends Thread {
|
||||
private int start = 0;
|
||||
private int end = BUFFER_MAX_SIZE;
|
||||
class ConsumerWorker extends Thread {
|
||||
private int start = 0;
|
||||
private int end = BUFFER_MAX_SIZE;
|
||||
|
||||
private ConsumerWorker(int start, int end) {
|
||||
super("ConsumerWorker");
|
||||
this.start = start;
|
||||
this.end = end;
|
||||
}
|
||||
private ConsumerWorker(int start, int end) {
|
||||
super("ConsumerWorker");
|
||||
this.start = start;
|
||||
this.end = end;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
StringBuilder data = new StringBuilder();
|
||||
while (true) {
|
||||
boolean bool = false;
|
||||
for (int i = start; i < end; i++) {
|
||||
if (dataBuffer[i] == null) {
|
||||
continue;
|
||||
}
|
||||
bool = true;
|
||||
if (data.length() + dataBuffer[i].toString().length() >= Config.Sender.MAX_SEND_LENGTH) {
|
||||
while (!DataSenderFactoryWithBalance.getSender().send(data.toString())) {
|
||||
try {
|
||||
Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.INFO, "send buried-point data.");
|
||||
data = new StringBuilder();
|
||||
}
|
||||
@Override
|
||||
public void run() {
|
||||
StringBuilder data = new StringBuilder();
|
||||
while (true) {
|
||||
boolean bool = false;
|
||||
try {
|
||||
for (int i = start; i < end; i++) {
|
||||
if (dataBuffer[i] == null) {
|
||||
continue;
|
||||
}
|
||||
bool = true;
|
||||
if (data.length() + dataBuffer[i].toString().length() >= Config.Sender.MAX_SEND_LENGTH) {
|
||||
while (!DataSenderFactoryWithBalance.getSender()
|
||||
.send(data.toString())) {
|
||||
try {
|
||||
Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
HealthCollector.getCurrentHeathReading(null)
|
||||
.updateData(HeathReading.INFO,
|
||||
"send buried-point data.");
|
||||
data = new StringBuilder();
|
||||
}
|
||||
|
||||
data.append(dataBuffer[i] + Constants.DATA_SPILT);
|
||||
dataBuffer[i] = null;
|
||||
}
|
||||
data.append(dataBuffer[i] + Constants.DATA_SPILT);
|
||||
dataBuffer[i] = null;
|
||||
}
|
||||
|
||||
if (data != null && data.length() > 0) {
|
||||
while (!DataSenderFactoryWithBalance.getSender().send(data.toString())) {
|
||||
try {
|
||||
Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
data = new StringBuilder();
|
||||
}
|
||||
if (data != null && data.length() > 0) {
|
||||
while (!DataSenderFactoryWithBalance.getSender().send(
|
||||
data.toString())) {
|
||||
try {
|
||||
Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
data = new StringBuilder();
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
logger.log(Level.ALL, "buffer group running failed", e);
|
||||
}
|
||||
|
||||
if (!bool) {
|
||||
try {
|
||||
Thread.sleep(MAX_WAIT_TIME);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!bool) {
|
||||
try {
|
||||
Thread.sleep(MAX_WAIT_TIME);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep Failure");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public String getGroupName() {
|
||||
return groupName;
|
||||
}
|
||||
public String getGroupName() {
|
||||
return groupName;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -45,6 +45,9 @@ public class Config {
|
|||
|
||||
// 是否开启发送
|
||||
public static boolean IS_OFF = false;
|
||||
|
||||
// 最大发送副本数量
|
||||
public static int MAX_COPY_NUM = 2;
|
||||
|
||||
// 发送的最大长度
|
||||
public static int MAX_SEND_LENGTH = 18500;
|
||||
|
|
|
|||
|
|
@ -7,7 +7,7 @@ public class ContextData {
|
|||
private String traceId;
|
||||
private String parentLevel;
|
||||
private int levelId;
|
||||
private char spanType;
|
||||
private String spanType;
|
||||
|
||||
ContextData() {
|
||||
|
||||
|
|
@ -23,10 +23,13 @@ public class ContextData {
|
|||
public ContextData(String contextDataStr) {
|
||||
// 反序列化参数
|
||||
String[] value = contextDataStr.split("-");
|
||||
if(value == null || value.length != 4){
|
||||
throw new IllegalArgumentException("illegal context data.");
|
||||
}
|
||||
this.traceId = value[0];
|
||||
this.parentLevel = value[1];
|
||||
this.levelId = Integer.valueOf(value[2]);
|
||||
this.spanType = value[3].charAt(0);
|
||||
this.spanType = value[3];
|
||||
|
||||
}
|
||||
|
||||
|
|
@ -42,7 +45,7 @@ public class ContextData {
|
|||
return levelId;
|
||||
}
|
||||
|
||||
public char getSpanType() {
|
||||
public String getSpanType() {
|
||||
return spanType;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -3,7 +3,7 @@ package com.ai.cloud.skywalking.model;
|
|||
public class Identification {
|
||||
private String viewPoint;
|
||||
private String businessKey;
|
||||
private char spanType;
|
||||
private String spanType;
|
||||
|
||||
public Identification() {
|
||||
//Non
|
||||
|
|
@ -17,7 +17,7 @@ public class Identification {
|
|||
return businessKey;
|
||||
}
|
||||
|
||||
public char getSpanType(){
|
||||
public String getSpanType(){
|
||||
return spanType;
|
||||
}
|
||||
|
||||
|
|
@ -46,7 +46,7 @@ public class Identification {
|
|||
return this;
|
||||
}
|
||||
|
||||
public IdentificationBuilder spanType(char spanType) {
|
||||
public IdentificationBuilder spanType(String spanType) {
|
||||
sendData.spanType = spanType;
|
||||
return this;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,7 +1,5 @@
|
|||
package com.ai.cloud.skywalking.sender;
|
||||
|
||||
import com.ai.cloud.skywalking.sender.protocol.ProtocolBuilder;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.net.StandardSocketOptions;
|
||||
|
|
@ -12,7 +10,9 @@ import java.nio.channels.SocketChannel;
|
|||
import java.util.logging.Level;
|
||||
import java.util.logging.Logger;
|
||||
|
||||
public class DataSender {
|
||||
import com.ai.cloud.skywalking.sender.protocol.ProtocolBuilder;
|
||||
|
||||
public class DataSender implements IDataSender{
|
||||
private static Logger logger = Logger.getLogger(DataSender.class.getName());
|
||||
|
||||
private SocketChannel socketChannel;
|
||||
|
|
@ -41,6 +41,7 @@ public class DataSender {
|
|||
* @param data
|
||||
* @return
|
||||
*/
|
||||
@Override
|
||||
public boolean send(String data) {
|
||||
// 发送报文
|
||||
try {
|
||||
|
|
|
|||
|
|
@ -14,197 +14,217 @@ 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的数量,就不需要保持那么多的保持连接的Sender的数量
|
||||
if (maxKeepConnectingSenderSize > Config.Consumer.MAX_CONSUMER * Config.Buffer.POOL_SIZE) {
|
||||
maxKeepConnectingSenderSize = Config.Consumer.MAX_CONSUMER * Config.Buffer.POOL_SIZE;
|
||||
}
|
||||
// 根据配置的服务器集群的地址,来计算保持连接的Sender的数量
|
||||
maxKeepConnectingSenderSize = calculateMaxKeeperConnectingSenderSize(tmpInetSocketAddress
|
||||
.size());
|
||||
// 最大连接消费线程小于保持连接的Sender的数量,就不需要保持那么多的保持连接的Sender的数量
|
||||
if (maxKeepConnectingSenderSize > Config.Consumer.MAX_CONSUMER
|
||||
* Config.Buffer.POOL_SIZE) {
|
||||
maxKeepConnectingSenderSize = Config.Consumer.MAX_CONSUMER
|
||||
* Config.Buffer.POOL_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 IDataSender getSender() {
|
||||
DataSender readySender = null;
|
||||
while (true) {
|
||||
try {
|
||||
if (usingDataSender.size() == 0) {
|
||||
try {
|
||||
Thread.sleep(RETRY_GET_SENDER_WAIT_INTERVAL);
|
||||
} catch (InterruptedException e) {
|
||||
logger.log(Level.ALL, "Sleep failed");
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
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", e);
|
||||
}
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
logger.log(Level.ALL, "get sender failed", e);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
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
|
||||
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();
|
||||
}
|
||||
}
|
||||
@Override
|
||||
public void run() {
|
||||
long sleepTime = 0;
|
||||
while (true) {
|
||||
try {
|
||||
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;
|
||||
// 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;
|
||||
int toBeSwitchIndex;
|
||||
|
||||
if (usingDataSender.size() - 1 > 0) {
|
||||
toBeSwitchIndex = ThreadLocalRandom.current()
|
||||
.nextInt(0, usingDataSender.size() - 1);
|
||||
} else {
|
||||
toBeSwitchIndex = 0;
|
||||
}
|
||||
if (usingDataSender.size() - 1 > 0) {
|
||||
toBeSwitchIndex = ThreadLocalRandom.current()
|
||||
.nextInt(0, usingDataSender.size() - 1);
|
||||
} else {
|
||||
toBeSwitchIndex = 0;
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
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;
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
logger.log(Level.ALL, "DataSenderChecker running failed", e);
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,40 @@
|
|||
package com.ai.cloud.skywalking.sender;
|
||||
|
||||
import static com.ai.cloud.skywalking.conf.Config.Sender.MAX_COPY_NUM;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* 带副本的数据发送器
|
||||
* @author wusheng
|
||||
*
|
||||
*/
|
||||
public class DataSenderWithCopies implements IDataSender{
|
||||
private int maxCopyNum;
|
||||
|
||||
private List<IDataSender> senders = new ArrayList<IDataSender>();
|
||||
|
||||
public DataSenderWithCopies(int maxKeepConnectingSenderSize){
|
||||
//最大副本数量,不能大于可用最大连接数
|
||||
maxCopyNum = maxKeepConnectingSenderSize > MAX_COPY_NUM ? MAX_COPY_NUM: maxKeepConnectingSenderSize;
|
||||
}
|
||||
|
||||
/**
|
||||
* 尝试增加到最大可用副本数,极端情况可能不足
|
||||
* @param dataSender
|
||||
* @return
|
||||
*/
|
||||
public boolean append(IDataSender dataSender){
|
||||
senders.add(dataSender);
|
||||
return maxCopyNum == senders.size();
|
||||
}
|
||||
|
||||
/**
|
||||
* 尝试向所有副本发送
|
||||
*/
|
||||
public boolean send(String data) {
|
||||
return false;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -0,0 +1,5 @@
|
|||
package com.ai.cloud.skywalking.sender;
|
||||
|
||||
public interface IDataSender {
|
||||
public boolean send(String data);
|
||||
}
|
||||
|
|
@ -42,7 +42,7 @@ public class Span extends SpanData {
|
|||
exceptionStack = fieldValues[8].trim().replaceAll(SPAN_ATTR_SPILT_CHARACTER,
|
||||
NEW_LINE_CHARACTER_PATTERN);
|
||||
}
|
||||
spanType = fieldValues[9].charAt(0);
|
||||
spanType = fieldValues[9];
|
||||
isReceiver = Boolean.valueOf(fieldValues[10]);
|
||||
|
||||
businessKey = fieldValues[11].trim().replaceAll(SPAN_ATTR_SPILT_CHARACTER,
|
||||
|
|
|
|||
|
|
@ -16,7 +16,7 @@ public abstract class SpanData {
|
|||
protected String address = "";
|
||||
protected byte statusCode = 0;
|
||||
protected String exceptionStack;
|
||||
protected char spanType = 'M';
|
||||
protected String spanType = "";
|
||||
protected boolean isReceiver = false;
|
||||
protected String businessKey = "";
|
||||
protected String processNo = "";
|
||||
|
|
@ -69,11 +69,11 @@ public abstract class SpanData {
|
|||
this.address = address;
|
||||
}
|
||||
|
||||
public char getSpanType() {
|
||||
public String getSpanType() {
|
||||
return spanType;
|
||||
}
|
||||
|
||||
public void setSpanType(char spanType) {
|
||||
public void setSpanType(String spanType) {
|
||||
this.spanType = spanType;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -87,7 +87,7 @@ public class SWDubboEnhanceFilter implements Filter {
|
|||
viewPoint.append(":" + invoker.getUrl().getPort());
|
||||
viewPoint.append(invoker.getUrl().getAbsolutePath());
|
||||
viewPoint.append("." + invocation.getMethodName() + "(");
|
||||
for (Class classes : invocation.getParameterTypes()) {
|
||||
for (Class<?> classes : invocation.getParameterTypes()) {
|
||||
viewPoint.append(classes.getSimpleName() + ",");
|
||||
}
|
||||
|
||||
|
|
@ -96,7 +96,7 @@ public class SWDubboEnhanceFilter implements Filter {
|
|||
}
|
||||
|
||||
viewPoint.append(")");
|
||||
return Identification.newBuilder().viewPoint(viewPoint.toString()).spanType('D').build();
|
||||
return Identification.newBuilder().viewPoint(viewPoint.toString()).spanType("D").build();
|
||||
}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -15,7 +15,7 @@ public class HttpClientTracing {
|
|||
httpRequest.setHeader(traceHearName,
|
||||
"ContextData=" + sender.beforeSend(Identification.newBuilder()
|
||||
.viewPoint(url)
|
||||
.spanType('W')
|
||||
.spanType("W")
|
||||
.build())
|
||||
.toString());
|
||||
return executor.execute();
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ public class CallableStatementTracing {
|
|||
"callableStatement."
|
||||
+ method
|
||||
+ (sql == null || sql.length() == 0 ? ""
|
||||
: ":" + sql)).spanType('J').build());
|
||||
: ":" + sql)).spanType("J").build());
|
||||
return exec.exe(realStatement, sql);
|
||||
} catch (SQLException e) {
|
||||
sender.handleException(e);
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ public class ConnectionTracing {
|
|||
"connection."
|
||||
+ method
|
||||
+ (sql == null || sql.length() == 0 ? ""
|
||||
: ":" + sql)).spanType('J').build());
|
||||
: ":" + sql)).spanType("J").build());
|
||||
return exec.exe(realConnection, sql);
|
||||
} catch (SQLException e) {
|
||||
sender.handleException(e);
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ public class PreparedStatementTracing {
|
|||
"preaparedStatement."
|
||||
+ method
|
||||
+ (sql == null || sql.length() == 0 ? ""
|
||||
: ":" + sql)).spanType('J').build());
|
||||
: ":" + sql)).spanType("J").build());
|
||||
return exec.exe(realStatement, sql);
|
||||
} catch (SQLException e) {
|
||||
sender.handleException(e);
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ public class StatementTracing {
|
|||
"statement."
|
||||
+ method
|
||||
+ (sql == null || sql.length() == 0 ? ""
|
||||
: ":" + sql)).spanType('J').build());
|
||||
: ":" + sql)).spanType("J").build());
|
||||
return exec.exe(realStatement, sql);
|
||||
} catch (SQLException e) {
|
||||
sender.handleException(e);
|
||||
|
|
|
|||
|
|
@ -58,7 +58,7 @@ public class SkyWalkingFilter implements Filter {
|
|||
private Identification generateIdentification(HttpServletRequest request) {
|
||||
return Identification.newBuilder()
|
||||
.viewPoint(request.getRequestURL().toString())
|
||||
.spanType('W')
|
||||
.spanType("W")
|
||||
.build();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -27,17 +27,17 @@ public class PersistenceThread extends Thread {
|
|||
BufferedReader bufferedReader = null;
|
||||
int offset;
|
||||
while (true) {
|
||||
file1 = getDataFiles();
|
||||
if (file1 == null) {
|
||||
try {
|
||||
Thread.sleep(SWITCH_FILE_WAIT_TIME);
|
||||
} catch (InterruptedException e) {
|
||||
logger.error("Failure sleep", e);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
try {
|
||||
file1 = getDataFiles();
|
||||
if (file1 == null) {
|
||||
try {
|
||||
Thread.sleep(SWITCH_FILE_WAIT_TIME);
|
||||
} catch (InterruptedException e) {
|
||||
logger.error("Failure sleep", e);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
bufferedReader = new BufferedReader(new FileReader(file1));
|
||||
offset = moveOffSet(file1, bufferedReader);
|
||||
if (logger.isDebugEnabled()) {
|
||||
|
|
|
|||
|
|
@ -135,9 +135,13 @@ public class TraceLogEntry extends Span {
|
|||
if (StringUtil.isBlank(spanTypeStr) || Constants.SPAN_TYPE_MAP.containsKey(spanTypeStr)) {
|
||||
result.spanTypeStr = Constants.SPAN_TYPE_U;
|
||||
}
|
||||
String spanTypeName = Constants.SPAN_TYPE_MAP.get(spanTypeStr);
|
||||
result.spanTypeStr = spanTypeStr;
|
||||
result.spanTypeName = spanTypeName;
|
||||
if(Constants.SPAN_TYPE_MAP.containsKey(spanTypeStr)){
|
||||
result.spanTypeName = Constants.SPAN_TYPE_MAP.get(spanTypeStr);;
|
||||
}else{
|
||||
//非默认支持的类型,使用原文中的类型,不需要解析
|
||||
result.spanTypeName = result.spanTypeStr;
|
||||
}
|
||||
|
||||
// 处理状态key-value
|
||||
String statusCodeStr = String.valueOf(result.getStatusCode());
|
||||
|
|
|
|||
Loading…
Reference in New Issue