From 734b14bf688397ade9dfbd51155a932d79404e0c Mon Sep 17 00:00:00 2001 From: wusheng Date: Fri, 19 Feb 2016 14:50:45 +0800 Subject: [PATCH] =?UTF-8?q?1.=E6=9C=8D=E5=8A=A1=E7=AB=AF=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=E8=BF=90=E8=A1=8C=E6=97=B6=E6=8A=A5=E5=91=8A=EF=BC=8C=E5=AE=9A?= =?UTF-8?q?=E5=88=B6=E6=97=A5=E5=BF=97=E8=BE=93=E5=87=BA=E7=BA=BF=E7=A8=8B?= =?UTF-8?q?=E7=9A=84=E8=BF=90=E8=A1=8C=E6=83=85=E5=86=B5=EF=BC=8C=E6=8F=90?= =?UTF-8?q?=E9=AB=98=E8=BE=A8=E8=AF=86=E5=BA=A6=E3=80=82=E7=A4=BA=E4=BE=8B?= =?UTF-8?q?=E5=A6=82=E4=B8=8B=EF=BC=9A=20---------Server=20Health=20Collec?= =?UTF-8?q?tor=20Report---------=20id=20[INFO]read=203=20chars=20from?= =?UTF-8?q?=20local=20file:1455864122412-37cd0a49fbb84ae9b111a08557cbd827(?= =?UTF-8?q?t:1455864245081)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit id [INFO]read 3 chars from local file:1455864122406-e19d28ec360d44d2a06335511f61d7c6(t:1455864245080) id [INFO]read 3 chars from local file:1455864122413-c0efe8bf31194c1887301ff3098eed00(t:1455864245080) id [INFO]read 3 chars from local file:1455864122413-79b44d640a7c4f9899d599c8354b861b(t:1455864245080) id [INFO]read 3 chars from local file:1455864122413-db2d14a825db4b5a93fa1a0a0e25d0c6(t:1455864245082) id [INFO]flush memory register to file.(t:1455864250063) ------------------------------------------------ --- ...864245067-d118af4fb24f47ff8e63be040a22a952 | 0 ...864245073-a95ac07cb29942a38eaf2c77892aad8d | 0 ...864245074-25f40f2c16594f2e9c069dfd5fffc362 | 0 ...864245074-340dfbbd62e04b9fb7fa95c0a1b6cb0c | 0 ...864245074-b765b5d1f7eb4bfc8a9d94cd177e1224 | 0 .../skywalking/reciever/CollectionServer.java | 25 +++++--- .../skywalking/reciever/conf/Config.java | 9 ++- .../skywalking/reciever/conf/Constants.java | 2 - .../ServerHealthCollector.java | 64 +++++++++++++++---- .../selfexamination/ServerHeathReading.java | 16 ++--- 10 files changed, 82 insertions(+), 34 deletions(-) create mode 100644 skywalking-server/D:/test-data/data/buffer/1455864245067-d118af4fb24f47ff8e63be040a22a952 create mode 100644 skywalking-server/D:/test-data/data/buffer/1455864245073-a95ac07cb29942a38eaf2c77892aad8d create mode 100644 skywalking-server/D:/test-data/data/buffer/1455864245074-25f40f2c16594f2e9c069dfd5fffc362 create mode 100644 skywalking-server/D:/test-data/data/buffer/1455864245074-340dfbbd62e04b9fb7fa95c0a1b6cb0c create mode 100644 skywalking-server/D:/test-data/data/buffer/1455864245074-b765b5d1f7eb4bfc8a9d94cd177e1224 diff --git a/skywalking-server/D:/test-data/data/buffer/1455864245067-d118af4fb24f47ff8e63be040a22a952 b/skywalking-server/D:/test-data/data/buffer/1455864245067-d118af4fb24f47ff8e63be040a22a952 new file mode 100644 index 000000000..e69de29bb diff --git a/skywalking-server/D:/test-data/data/buffer/1455864245073-a95ac07cb29942a38eaf2c77892aad8d b/skywalking-server/D:/test-data/data/buffer/1455864245073-a95ac07cb29942a38eaf2c77892aad8d new file mode 100644 index 000000000..e69de29bb diff --git a/skywalking-server/D:/test-data/data/buffer/1455864245074-25f40f2c16594f2e9c069dfd5fffc362 b/skywalking-server/D:/test-data/data/buffer/1455864245074-25f40f2c16594f2e9c069dfd5fffc362 new file mode 100644 index 000000000..e69de29bb diff --git a/skywalking-server/D:/test-data/data/buffer/1455864245074-340dfbbd62e04b9fb7fa95c0a1b6cb0c b/skywalking-server/D:/test-data/data/buffer/1455864245074-340dfbbd62e04b9fb7fa95c0a1b6cb0c new file mode 100644 index 000000000..e69de29bb diff --git a/skywalking-server/D:/test-data/data/buffer/1455864245074-b765b5d1f7eb4bfc8a9d94cd177e1224 b/skywalking-server/D:/test-data/data/buffer/1455864245074-b765b5d1f7eb4bfc8a9d94cd177e1224 new file mode 100644 index 000000000..e69de29bb diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java index 535da096c..a0812f5e7 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/CollectionServer.java @@ -1,12 +1,11 @@ package com.ai.cloud.skywalking.reciever; -import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; -import com.ai.cloud.skywalking.reciever.conf.Config; -import com.ai.cloud.skywalking.reciever.conf.ConfigInitializer; -import com.ai.cloud.skywalking.reciever.handler.CollectionServerDataHandler; -import com.ai.cloud.skywalking.reciever.persistance.PersistenceThreadLauncher; import io.netty.bootstrap.ServerBootstrap; -import io.netty.channel.*; +import io.netty.channel.ChannelFuture; +import io.netty.channel.ChannelInitializer; +import io.netty.channel.ChannelOption; +import io.netty.channel.ChannelPipeline; +import io.netty.channel.EventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.socket.nio.NioServerSocketChannel; import io.netty.handler.codec.LengthFieldBasedFrameDecoder; @@ -15,12 +14,20 @@ import io.netty.handler.codec.bytes.ByteArrayDecoder; import io.netty.handler.codec.bytes.ByteArrayEncoder; import io.netty.handler.logging.LogLevel; import io.netty.handler.logging.LoggingHandler; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; import java.io.IOException; import java.util.Properties; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; +import com.ai.cloud.skywalking.reciever.conf.Config; +import com.ai.cloud.skywalking.reciever.conf.ConfigInitializer; +import com.ai.cloud.skywalking.reciever.handler.CollectionServerDataHandler; +import com.ai.cloud.skywalking.reciever.persistance.PersistenceThreadLauncher; +import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; + public class CollectionServer { static Logger logger = LogManager.getLogger(CollectionServer.class); @@ -62,6 +69,8 @@ public class CollectionServer { public static void main(String[] args) throws IOException, IllegalAccessException, InterruptedException { logger.info("To initialize the collect server configuration parameters...."); initializeParam(); + logger.info("To init server health collector..."); + ServerHealthCollector.init(); logger.info("To launch register persistence thread...."); PersistenceThreadLauncher.doLaunch(); logger.info("To init data buffer thread container..."); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java index 3ab8f3383..a3f357e9a 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java @@ -72,9 +72,9 @@ public class Config { } public static class HBaseConfig { - // + public static String TABLE_NAME = "sw-call-chain"; - // + public static String FAMILY_COLUMN_NAME = "call-chain"; public static String ZK_HOSTNAME; @@ -104,4 +104,9 @@ public class Config { public static boolean ALARM_OFF_FLAG = false; } + + public static class HealthCollector { + // 默认健康检查上报时间 + public static long REPORT_INTERVAL = 5 * 60 * 1000L; + } } \ No newline at end of file diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Constants.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Constants.java index 9107b73f7..13f210b9b 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Constants.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Constants.java @@ -1,7 +1,5 @@ package com.ai.cloud.skywalking.reciever.conf; public class Constants { - public static final String HEALTH_DATA_SPILT_PATTERN = "^~"; - public static final String DATA_SPILT = "#&"; } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHealthCollector.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHealthCollector.java index d100ba6d9..ea2e17a1f 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHealthCollector.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHealthCollector.java @@ -1,33 +1,71 @@ package com.ai.cloud.skywalking.reciever.selfexamination; -import java.util.HashMap; +import java.util.Arrays; import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.util.MachineUtil; -public class ServerHealthCollector extends Thread{ - private static Map heathReadings = new HashMap(); +public class ServerHealthCollector extends Thread { + private Logger logger = LogManager.getLogger(ServerHealthCollector.class); + + private static Map heathReadings = new ConcurrentHashMap(); + + private ServerHealthCollector(){ + super("ServerHealthCollector"); + } - public static ServerHeathReading getCurrentHeathReading(String extraId){ + public static void init(){ + new ServerHealthCollector().start(); + } + + public static ServerHeathReading getCurrentHeathReading(String extraId) { String id = getId(extraId); - if(!heathReadings.containsKey(id)){ + if (!heathReadings.containsKey(id)) { synchronized (heathReadings) { - if(!heathReadings.containsKey(id)){ + if (!heathReadings.containsKey(id)) { heathReadings.put(id, new ServerHeathReading(id)); } } } return heathReadings.get(id); } - - private static String getId(String extraId){ - return "SkyWalkingServer,M:" + MachineUtil.getHostDesc() +",P:" + MachineUtil.getProcessNo() + ",T:" + + private static String getId(String extraId) { + return "SkyWalkingServer,M:" + MachineUtil.getHostDesc() + ",P:" + + MachineUtil.getProcessNo() + ",T:" + Thread.currentThread().getName() + "(" - + Thread.currentThread().getId() + ")" + (extraId == null? "" : ",extra:" + extraId); + + Thread.currentThread().getId() + ")" + + (extraId == null ? "" : ",extra:" + extraId); } - + @Override - public void run(){ - //TODO: 服务端本地存储,用于将信息存储如数据库,并供前台展现,完成定时刷新 + public void run() { + while (true) { + try { + String[] keyList = heathReadings.keySet().toArray(new String[0]); + Arrays.sort(keyList); + StringBuilder log = new StringBuilder(); + log.append("\n---------Server Health Collector Report---------\n"); + for(String key : keyList){ + log.append(heathReadings.get(key)).append("\n"); + } + log.append("------------------------------------------------\n"); + + logger.info(log); + + try { + Thread.sleep(Config.HealthCollector.REPORT_INTERVAL); + } catch (InterruptedException e) { + logger.warn("sleep error.", e); + } + } catch (Throwable t) { + logger.error("ServerHealthCollector report error.", t); + } + } } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHeathReading.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHeathReading.java index 9dfcdfb9c..8731dfe65 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHeathReading.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/selfexamination/ServerHeathReading.java @@ -3,12 +3,10 @@ package com.ai.cloud.skywalking.reciever.selfexamination; import java.util.HashMap; import java.util.Map; -import com.ai.cloud.skywalking.reciever.conf.Constants; - public class ServerHeathReading { - public static final String ERROR = "ERROR"; - public static final String WARNING = "WARNING"; - public static final String INFO = "INFO"; + public static final String ERROR = "[ERROR]"; + public static final String WARNING = "[WARNING]"; + public static final String INFO = "[INFO]"; private String id; @@ -32,10 +30,10 @@ public class ServerHeathReading { @Override public String toString(){ - StringBuilder sb = new StringBuilder(this.id); - sb.append(Constants.HEALTH_DATA_SPILT_PATTERN); + StringBuilder sb = new StringBuilder(); + sb.append("id<").append(this.id).append(">\n"); for(Map.Entry data : datas.entrySet()){ - sb.append(data.getKey()).append(Constants.HEALTH_DATA_SPILT_PATTERN).append(data.getValue().toString()).append(Constants.HEALTH_DATA_SPILT_PATTERN); + sb.append(data.getKey()).append(data.getValue().toString()).append("\n"); } //reset data @@ -68,7 +66,7 @@ public class ServerHeathReading { @Override public String toString(){ - return "d:" + data + Constants.HEALTH_DATA_SPILT_PATTERN + "t:" + statusTime; + return data + "(t:" + statusTime + ")"; } } }