From cd60b5adf78ee6ac8d2b2d026a75d1ce0a096435 Mon Sep 17 00:00:00 2001 From: wusheng Date: Mon, 21 Dec 2015 23:20:56 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BF=AE=E6=94=B9SpanType=E4=B8=BA=E5=AD=97?= =?UTF-8?q?=E7=AC=A6=E4=B8=B2=EF=BC=8C=E4=BF=AE=E6=94=B9webui=EF=BC=8C?= =?UTF-8?q?=E6=8F=90=E9=AB=98=E7=B1=BB=E5=9E=8B=E7=9A=84=E6=89=A9=E5=B1=95?= =?UTF-8?q?=E6=80=A7=202.=E4=B8=BA=E5=90=8E=E7=BB=AD=E5=A2=9E=E5=8A=A0Data?= =?UTF-8?q?Sender=E6=94=AF=E6=8C=81=E5=89=AF=E6=9C=AC=E5=8F=91=E9=80=81?= =?UTF-8?q?=EF=BC=8C=E5=A2=9E=E5=8A=A0=E9=83=A8=E5=88=86=E4=BB=A3=E7=A0=81?= =?UTF-8?q?=E3=80=82=E5=8A=9F=E8=83=BD=E6=9A=82=E6=9C=AA=E5=AE=8C=E6=88=90?= =?UTF-8?q?=E3=80=82=203.=E4=B8=BADataSenderFactoryWithBalance.getSender?= =?UTF-8?q?=E6=96=B9=E6=B3=95=E5=A2=9E=E5=8A=A0=E5=BC=82=E5=B8=B8=E6=9C=BA?= =?UTF-8?q?=E5=88=B6=EF=BC=8C=E6=96=B9=E6=B3=95=E5=BC=82=E5=B8=B8=E9=80=80?= =?UTF-8?q?=E5=87=BA=E3=80=82=204.=E4=BF=AE=E5=A4=8D=E5=AE=A2=E6=88=B7?= =?UTF-8?q?=E7=AB=AF=E5=A4=A7=E9=87=8F=E7=BA=BF=E7=A8=8B=E7=BC=BA=E5=B0=91?= =?UTF-8?q?=E5=BC=82=E5=B8=B8=E4=BF=9D=E9=9A=9C=E8=83=BD=E5=8A=9B=E7=9A=84?= =?UTF-8?q?=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/main/resources/sky-walking.auth | 11 +- .../cloud/skywalking/buffer/BufferGroup.java | 173 +++++---- .../com/ai/cloud/skywalking/conf/Config.java | 3 + .../cloud/skywalking/model/ContextData.java | 9 +- .../skywalking/model/Identification.java | 6 +- .../cloud/skywalking/sender/DataSender.java | 7 +- .../sender/DataSenderFactoryWithBalance.java | 356 +++++++++--------- .../sender/DataSenderWithCopies.java | 40 ++ .../cloud/skywalking/sender/IDataSender.java | 5 + .../ai/cloud/skywalking/protocol/Span.java | 2 +- .../cloud/skywalking/protocol/SpanData.java | 6 +- .../plugin/dubbo/SWDubboEnhanceFilter.java | 4 +- .../httpclient/trace/HttpClientTracing.java | 2 +- .../tracing/CallableStatementTracing.java | 2 +- .../jdbc/tracing/ConnectionTracing.java | 2 +- .../tracing/PreparedStatementTracing.java | 2 +- .../plugin/jdbc/tracing/StatementTracing.java | 2 +- .../plugin/web/SkyWalkingFilter.java | 2 +- .../persistance/PersistenceThread.java | 20 +- .../com/ai/cloud/vo/mvo/TraceLogEntry.java | 8 +- 20 files changed, 376 insertions(+), 286 deletions(-) create mode 100644 skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java create mode 100644 skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/IDataSender.java diff --git a/samples/skywalking-auth/src/main/resources/sky-walking.auth b/samples/skywalking-auth/src/main/resources/sky-walking.auth index b6c09a441..d53f3c163 100644 --- a/samples/skywalking-auth/src/main/resources/sky-walking.auth +++ b/samples/skywalking-auth/src/main/resources/sky-walking.auth @@ -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 diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java index 4ed5cd364..6f2a8e546 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java @@ -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; + } } 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 efe946620..5bae26cad 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 @@ -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; diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java index 5a1ca009f..fbae131ce 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java @@ -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; } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/Identification.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/Identification.java index 4e33985b5..fd0507e05 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/Identification.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/Identification.java @@ -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; } 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 bb3a9dc69..0e1ebbae5 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 @@ -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 { 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 9e7f0105a..9cef15a32 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 @@ -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 unusedServerAddresses = new ArrayList(); + private static Logger logger = Logger + .getLogger(DataSenderFactoryWithBalance.class.getName()); + // unUsedServerAddress存放没有使用的服务器地址, + private static List unusedServerAddresses = new ArrayList(); - private static List usingDataSender = new ArrayList(); - private static int maxKeepConnectingSenderSize; + private static List usingDataSender = new ArrayList(); + 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 tmpInetSocketAddress = new HashSet(); - 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 tmpInetSocketAddress = new HashSet(); + 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 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 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); + } + } } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java new file mode 100644 index 000000000..aabb4147d --- /dev/null +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java @@ -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 senders = new ArrayList(); + + 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; + } + +} diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/IDataSender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/IDataSender.java new file mode 100644 index 000000000..ebe2b7edf --- /dev/null +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/IDataSender.java @@ -0,0 +1,5 @@ +package com.ai.cloud.skywalking.sender; + +public interface IDataSender { + public boolean send(String data); +} diff --git a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java index 24cdea067..32a752b60 100644 --- a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java +++ b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java @@ -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, diff --git a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java index 556d3f4fb..ab1d2d8b4 100644 --- a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java +++ b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java @@ -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; } diff --git a/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/SWDubboEnhanceFilter.java b/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/SWDubboEnhanceFilter.java index 3b32a3fdb..801e3a4fc 100644 --- a/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/SWDubboEnhanceFilter.java +++ b/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/SWDubboEnhanceFilter.java @@ -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(); } diff --git a/skywalking-sdk-plugin/httpclient-plugin/src/main/java/com/ai/cloud/skywalking/plugin/httpclient/trace/HttpClientTracing.java b/skywalking-sdk-plugin/httpclient-plugin/src/main/java/com/ai/cloud/skywalking/plugin/httpclient/trace/HttpClientTracing.java index 3336374a6..8a6af2144 100644 --- a/skywalking-sdk-plugin/httpclient-plugin/src/main/java/com/ai/cloud/skywalking/plugin/httpclient/trace/HttpClientTracing.java +++ b/skywalking-sdk-plugin/httpclient-plugin/src/main/java/com/ai/cloud/skywalking/plugin/httpclient/trace/HttpClientTracing.java @@ -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(); diff --git a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/CallableStatementTracing.java b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/CallableStatementTracing.java index 9995d0fd7..9838e6bae 100644 --- a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/CallableStatementTracing.java +++ b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/CallableStatementTracing.java @@ -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); diff --git a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/ConnectionTracing.java b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/ConnectionTracing.java index b92f72cdc..c5cb12143 100644 --- a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/ConnectionTracing.java +++ b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/ConnectionTracing.java @@ -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); diff --git a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/PreparedStatementTracing.java b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/PreparedStatementTracing.java index 974da70ab..5226bd844 100644 --- a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/PreparedStatementTracing.java +++ b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/PreparedStatementTracing.java @@ -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); diff --git a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/StatementTracing.java b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/StatementTracing.java index a599117cc..ed485eeaf 100644 --- a/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/StatementTracing.java +++ b/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/tracing/StatementTracing.java @@ -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); diff --git a/skywalking-sdk-plugin/web-plugin/src/main/java/com/ai/cloud/skywalking/plugin/web/SkyWalkingFilter.java b/skywalking-sdk-plugin/web-plugin/src/main/java/com/ai/cloud/skywalking/plugin/web/SkyWalkingFilter.java index b486841d9..f700eefd2 100644 --- a/skywalking-sdk-plugin/web-plugin/src/main/java/com/ai/cloud/skywalking/plugin/web/SkyWalkingFilter.java +++ b/skywalking-sdk-plugin/web-plugin/src/main/java/com/ai/cloud/skywalking/plugin/web/SkyWalkingFilter.java @@ -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(); } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java index 87c29e986..1c1f4f913 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java @@ -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()) { diff --git a/skywalking-webui/src/main/java/com/ai/cloud/vo/mvo/TraceLogEntry.java b/skywalking-webui/src/main/java/com/ai/cloud/vo/mvo/TraceLogEntry.java index b5a785f8f..b55f08146 100644 --- a/skywalking-webui/src/main/java/com/ai/cloud/vo/mvo/TraceLogEntry.java +++ b/skywalking-webui/src/main/java/com/ai/cloud/vo/mvo/TraceLogEntry.java @@ -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());