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 + ")"; } } }