From f22986e30e5c99e4677aefd97721457d6e6aa8bb Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Tue, 1 Dec 2015 16:55:17 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BF=AE=E6=94=B9Span=E7=9A=84=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E5=88=86=E9=9A=94=E7=AC=A6=202.=E4=BF=AE=E6=94=B9?= =?UTF-8?q?=E5=86=99=E5=85=A5=E6=96=87=E4=BB=B6=E5=92=8C=E8=AF=BB=E5=85=A5?= =?UTF-8?q?HBase=E7=BA=BF=E7=A8=8B=E6=95=B0=E9=87=8F=E6=94=B9=E4=B8=BA?= =?UTF-8?q?=E4=B8=80=E4=B8=AA=E5=8F=82=E6=95=B0=E6=8E=A7=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../buffer/DataBufferThreadContainer.java | 6 +- .../skywalking/reciever/conf/Config.java | 14 +- .../skywalking/reciever/conf/Constants.java | 2 + .../handler/CollectionServerDataHandler.java | 13 +- .../persistance/PersistenceThread.java | 330 +++++++++--------- .../PersistenceThreadLauncher.java | 2 +- .../storage/StorageChainController.java | 3 +- .../src/main/resources/config.properties | 9 +- 8 files changed, 200 insertions(+), 179 deletions(-) diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java index 9c4449c3d..34f988dda 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java @@ -11,7 +11,6 @@ import java.util.Arrays; import java.util.List; import java.util.concurrent.ThreadLocalRandom; -import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.MAX_THREAD_NUMBER; import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.MAX_APPEND_EOF_FLAGS_THREAD_NUMBER; public class DataBufferThreadContainer { @@ -50,8 +49,9 @@ public class DataBufferThreadContainer { logger.debug("start:" + start + "\tend:" + end); } } - logger.info("Data buffer thread size {} begin to init ", MAX_THREAD_NUMBER); - for (int i = 0; i < MAX_THREAD_NUMBER; i++) { + logger.info("Data buffer thread size {} begin to init ", Config.Server. + MAX_DEAL_DATA_THREAD_NUMBER); + for (int i = 0; i < Config.Server.MAX_DEAL_DATA_THREAD_NUMBER; i++) { DataBufferThread dataBufferThread = new DataBufferThread(); dataBufferThread.start(); buffers.add(dataBufferThread); 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 5c55e4f92..25094507a 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 @@ -6,12 +6,12 @@ public class Config { public static class Server { // 采集服务器的端口 public static int PORT = 34000; + // 最大数据处理线程数量 + public static int MAX_DEAL_DATA_THREAD_NUMBER = 3; } // 数据缓存配置类 public static class Buffer { - // 最大数据缓存线程数量 - public static int MAX_THREAD_NUMBER = 3; //每个线程最大缓存数量 public static int PER_THREAD_MAX_BUFFER_NUMBER = 1024; @@ -37,13 +37,10 @@ public class Config { } public static class Persistence { - // 最大持久化的线程数量 - public static int MAX_THREAD_NUMBER = 1; - // 定位文件时,每次读取偏移量跳过大小 public static int STEP_SIZE_FOR_LOCATING_FILE_OFFSET = 2048; - // 处理文件完成之后,等待时间 + // 切换文件,等待时间 public static long SWITCH_FILE_WAIT_TIME = 5000L; // 追加EOF标志位的线程数量 @@ -51,6 +48,9 @@ public class Config { // 每次存储的最大数量 public static final int MAX_STORAGE_SIZE_PER_TIME = 1024 * 1024; + + // 当读取到文件结束时等待时间 + public static long READ_ENDING_FILE_MAX_WAITE_TIME = 500L; } public static class RegisterPersistence { @@ -80,7 +80,7 @@ public class Config { public static class StorageChain { public static long RETRY_STORAGE_WAIT_TIME = 50L; - + public static String STORAGE_TYPE = "hbase"; } } \ 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 318b23178..a3a477846 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 @@ -4,4 +4,6 @@ public class Constants { public static final String spiltRegx = "\\^\\~"; 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/handler/CollectionServerDataHandler.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java index 30b122e23..d52d1a888 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java @@ -1,16 +1,17 @@ package com.ai.cloud.skywalking.reciever.handler; +import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; -import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; - public class CollectionServerDataHandler extends SimpleChannelInboundHandler { - + @Override protected void channelRead0(ChannelHandlerContext ctx, byte[] msg) throws Exception { - Thread.currentThread().setName("ServerReceiver"); - - DataBufferThreadContainer.getDataBufferThread().saveTemporarily(msg); + Thread.currentThread().setName("ServerReceiver"); + // 当接受到这条消息的是空,则忽略 + if (msg != null && msg.length >= 0) { + DataBufferThreadContainer.getDataBufferThread().saveTemporarily(msg); + } } } 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 61fd18529..1b0dafd11 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 @@ -4,7 +4,6 @@ import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; import com.ai.cloud.skywalking.reciever.storage.StorageChainController; - import org.apache.commons.io.FileUtils; import org.apache.commons.io.comparator.NameFileComparator; import org.apache.logging.log4j.LogManager; @@ -15,171 +14,190 @@ import java.io.*; import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.*; public class PersistenceThread extends Thread { - private Logger logger = LogManager.getLogger(PersistenceThread.class); - PersistenceThread() { - super("PersistenceThread"); - } + private Logger logger = LogManager.getLogger(PersistenceThread.class); - @Override - public void run() { - File file1; - 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; - } + PersistenceThread() { + super("PersistenceThread"); + } - try { - bufferedReader = new BufferedReader(new FileReader(file1)); - offset = moveOffSet(file1, bufferedReader); - if (logger.isDebugEnabled()) { - logger.debug("Get file[{}] offset [{}]", file1.getName(), - offset); - } - StringBuilder stringBuilder = new StringBuilder( - MAX_STORAGE_SIZE_PER_TIME); - String tmpData; - while (true) { - tmpData = bufferedReader.readLine(); + @Override + public void run() { + File file1; + 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; + } - if (tmpData == null || tmpData.length() <= 0) { - if (stringBuilder != null && stringBuilder.length() > 0) { - StorageChainController.doStorage(stringBuilder - .toString()); - stringBuilder.delete(0, stringBuilder.length()); - } - MemoryRegister - .instance() - .doRegisterStatus( - new FileRegisterEntry( - file1.getName(), - offset, - FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); - break; - } - ServerHealthCollector.getCurrentHeathReading(null) - .updateData( - ServerHeathReading.INFO, - "read " + tmpData.length() - + " chars from local file:" + file1.getName()); + try { + bufferedReader = new BufferedReader(new FileReader(file1)); + offset = moveOffSet(file1, bufferedReader); + if (logger.isDebugEnabled()) { + logger.debug("Get file[{}] offset [{}]", file1.getName(), + offset); + } + StringBuilder stringBuilder = new StringBuilder( + MAX_STORAGE_SIZE_PER_TIME); + String tmpData; + while (true) { + tmpData = bufferedReader.readLine(); + //文件结束 + if (tmpData == null) { + if (stringBuilder != null && stringBuilder.length() > 0) { + StorageChainController.doStorage(stringBuilder + .toString()); + stringBuilder.delete(0, stringBuilder.length()); + } - if ("EOF".equals(tmpData)) { - if (stringBuilder != null && stringBuilder.length() > 0) { - StorageChainController.doStorage(stringBuilder - .toString()); - } + try { + Thread.sleep(READ_ENDING_FILE_MAX_WAITE_TIME); + } catch (InterruptedException e) { + logger.error("Sleep failed", e); + } - bufferedReader.close(); - logger.info( - "Data in file[{}] has been successfully processed", - file1.getName()); - boolean deleteSuccess = false; - while (!deleteSuccess) { - deleteSuccess = FileUtils.deleteQuietly(new File( - file1.getParent(), file1.getName())); - } - logger.info("Delete file[{}] {}", file1.getName(), - (deleteSuccess ? "success" : "failed")); - MemoryRegister.instance().unRegister(file1.getName()); - break; - } + continue; + } - if (stringBuilder.length() + tmpData.length() >= MAX_STORAGE_SIZE_PER_TIME) { - StorageChainController.doStorage(stringBuilder - .toString()); - stringBuilder.delete(0, stringBuilder.length()); - MemoryRegister - .instance() - .doRegisterStatus( - new FileRegisterEntry( - file1.getName(), - offset, - FileRegisterEntry.FileRegisterEntryStatus.REGISTER)); - } + //文件读入/n字符串 + if (tmpData.length() <= 0) { + // 加上回车的字符串长度 + offset += 1; + continue; + } +// +// if (tmpData == null || tmpData.length() <= 0) { +// MemoryRegister +// .instance() +// .doRegisterStatus( +// new FileRegisterEntry( +// file1.getName(), +// offset, +// FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); +// break; +// } + ServerHealthCollector.getCurrentHeathReading(null) + .updateData( + ServerHeathReading.INFO, + "read " + tmpData.length() + + " chars from local file:" + file1.getName()); - stringBuilder.append(tmpData); - // 加上回车的字符串长度 - offset += tmpData.length() + 1; - } - } catch (FileNotFoundException e) { - logger.error("The data file could not be found", e); - } catch (IOException e) { - logger.error("The data file could not be found", e); - } finally { - try { - if (bufferedReader != null) - bufferedReader.close(); - } catch (IOException e) { - logger.error("can't close data file", e); - } - } + if ("EOF".equals(tmpData)) { + if (stringBuilder != null && stringBuilder.length() > 0) { + StorageChainController.doStorage(stringBuilder + .toString()); + } - try { - Thread.sleep(SWITCH_FILE_WAIT_TIME); - } catch (InterruptedException e) { - logger.error("Failure sleep.", e); - } - } - } + bufferedReader.close(); + logger.info( + "Data in file[{}] has been successfully processed", + file1.getName()); + boolean deleteSuccess = false; + while (!deleteSuccess) { + deleteSuccess = FileUtils.deleteQuietly(new File( + file1.getParent(), file1.getName())); + } + logger.info("Delete file[{}] {}", file1.getName(), + (deleteSuccess ? "success" : "failed")); + MemoryRegister.instance().unRegister(file1.getName()); + break; + } - private int moveOffSet(File file1, BufferedReader bufferedReader) - throws IOException { - int offset = MemoryRegister.instance().getOffSet(file1.getName()); - if (-1 == offset || offset == 0) { - // 以前该文件没有被任何人处理过,需要重新注册 - MemoryRegister - .instance() - .doRegisterStatus( - new FileRegisterEntry( - file1.getName(), - 0, - FileRegisterEntry.FileRegisterEntryStatus.REGISTER)); - offset = 0; - } else { - char[] cha = new char[STEP_SIZE_FOR_LOCATING_FILE_OFFSET]; - int length = 0; - while (length + STEP_SIZE_FOR_LOCATING_FILE_OFFSET < offset) { - length += STEP_SIZE_FOR_LOCATING_FILE_OFFSET; - bufferedReader.read(cha); - } - bufferedReader.read(cha, 0, Math.abs(offset - length)); - cha = null; - } - return offset; - } + if (stringBuilder.length() + tmpData.length() >= MAX_STORAGE_SIZE_PER_TIME) { + StorageChainController.doStorage(stringBuilder + .toString()); + stringBuilder.delete(0, stringBuilder.length()); + MemoryRegister + .instance() + .doRegisterStatus( + new FileRegisterEntry( + file1.getName(), + offset, + FileRegisterEntry.FileRegisterEntryStatus.REGISTER)); + } - private File getDataFiles() { - File file1 = null; - File parentDir = new File( - Config.Buffer.DATA_BUFFER_FILE_PARENT_DIRECTORY); - NameFileComparator sizeComparator = new NameFileComparator(); - File[] dataFileList = sizeComparator.sort(parentDir.listFiles()); - for (File file : dataFileList) { - if (file.getName().startsWith(".")) { - continue; - } - if (MemoryRegister.instance().isRegister(file.getName())) { - if (logger.isDebugEnabled()) - logger.debug( - "The file [{}] is being used by another thread ", - file); - continue; - } - if (logger.isDebugEnabled()) { - logger.debug("Begin to deal data file [{}]", file.getName()); - } - file1 = file; - break; - } + stringBuilder.append(tmpData); + // 加上回车的字符串长度 + offset += tmpData.length() + 1; + } + } catch (FileNotFoundException e) { + logger.error("The data file could not be found", e); + } catch (IOException e) { + logger.error("The data file could not be found", e); + } finally { + try { + if (bufferedReader != null) + bufferedReader.close(); + } catch (IOException e) { + logger.error("can't close data file", e); + } + } - return file1; - } + try { + Thread.sleep(SWITCH_FILE_WAIT_TIME); + } catch (InterruptedException e) { + logger.error("Failure sleep.", e); + } + } + } + + private int moveOffSet(File file1, BufferedReader bufferedReader) + throws IOException { + int offset = MemoryRegister.instance().getOffSet(file1.getName()); + if (-1 == offset || offset == 0) { + // 以前该文件没有被任何人处理过,需要重新注册 + MemoryRegister + .instance() + .doRegisterStatus( + new FileRegisterEntry( + file1.getName(), + 0, + FileRegisterEntry.FileRegisterEntryStatus.REGISTER)); + offset = 0; + } else { + char[] cha = new char[STEP_SIZE_FOR_LOCATING_FILE_OFFSET]; + int length = 0; + while (length + STEP_SIZE_FOR_LOCATING_FILE_OFFSET < offset) { + length += STEP_SIZE_FOR_LOCATING_FILE_OFFSET; + bufferedReader.read(cha); + } + bufferedReader.read(cha, 0, Math.abs(offset - length)); + cha = null; + } + return offset; + } + + private File getDataFiles() { + File file1 = null; + File parentDir = new File( + Config.Buffer.DATA_BUFFER_FILE_PARENT_DIRECTORY); + NameFileComparator sizeComparator = new NameFileComparator(); + File[] dataFileList = sizeComparator.sort(parentDir.listFiles()); + for (File file : dataFileList) { + if (file.getName().startsWith(".")) { + continue; + } + if (MemoryRegister.instance().isRegister(file.getName())) { + if (logger.isDebugEnabled()) + logger.debug( + "The file [{}] is being used by another thread ", + file); + continue; + } + if (logger.isDebugEnabled()) { + logger.debug("Begin to deal data file [{}]", file.getName()); + } + file1 = file; + break; + } + + return file1; + } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThreadLauncher.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThreadLauncher.java index 5426f31ab..b88982f5d 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThreadLauncher.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThreadLauncher.java @@ -6,7 +6,7 @@ public class PersistenceThreadLauncher { public static void doLaunch() { new RegisterPersistenceThread().start(); - for (int i = 0; i < Config.Persistence.MAX_THREAD_NUMBER; i++) { + for (int i = 0; i < Config.Server.MAX_DEAL_DATA_THREAD_NUMBER; i++) { new PersistenceThread().start(); } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java index c57bbb189..d76e7c8e8 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java @@ -2,6 +2,7 @@ package com.ai.cloud.skywalking.reciever.storage; import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.reciever.conf.Config; +import com.ai.cloud.skywalking.reciever.conf.Constants; import com.ai.cloud.skywalking.reciever.storage.chain.SaveToHBaseChain; import com.ai.cloud.skywalking.reciever.storage.chain.SaveToMySQLChain; import org.apache.logging.log4j.LogManager; @@ -29,7 +30,7 @@ public class StorageChainController { } public static void doStorage(String buriedPointDatas) { - String[] buriedPointData = buriedPointDatas.split(";"); + String[] buriedPointData = buriedPointDatas.split(Constants.DATA_SPILT); if (buriedPointData == null || buriedPointData.length == 0) { return; } diff --git a/skywalking-server/src/main/resources/config.properties b/skywalking-server/src/main/resources/config.properties index 9d5284d22..893ef4e5f 100644 --- a/skywalking-server/src/main/resources/config.properties +++ b/skywalking-server/src/main/resources/config.properties @@ -1,8 +1,7 @@ #采集服务器的端口 server.port=34000 +server.max_deal_data_thread_number=1 -#最大数据缓存线程数量 -buffer.max_thread_number=1 #每个线程最大缓存数量 buffer.per_thread_max_buffer_number=1024 #无数据处理时轮询等待时间(单位:毫秒) @@ -16,14 +15,14 @@ buffer.data_file_max_length=104857600 #每次Flush的缓存数据的个数 buffer.flush_number_of_cache=30 -#最大持久化的线程数量 -persistence.max_thread_number=3 #定位文件时,每次读取偏移量跳过大小 persistence.step_size_for_location_file_offset=20480 -#处理文件完成之后,等待时间(单位:毫秒) +#切换数据文件,等待时间(单位:毫秒) persistence.switch_file_wait_time=5000 #追加EOF标志位的线程数量 persistence.max_append_eof_flags_thread_number=2 +#当读取文件结束时最大等待时间 +persistence.read_ending_file_max_waite_time=50 #偏移量注册文件的目录 registerpersistence.register_file_parent_directory=d:/test-data/data/offset