From 042e56cbde94677ee011b2168d77ffd84599e48a Mon Sep 17 00:00:00 2001 From: wusheng Date: Thu, 25 Feb 2016 13:30:07 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BF=AE=E6=94=B9health=20report=E7=9A=84?= =?UTF-8?q?=E9=83=A8=E5=88=86=E5=86=85=E5=AE=B9=EF=BC=8C=E6=98=8E=E7=A1=AE?= =?UTF-8?q?=E9=83=A8=E5=88=86=E5=B1=95=E7=8E=B0=E4=BF=A1=E6=81=AF=E3=80=82?= =?UTF-8?q?=202.=E4=B8=B0=E5=AF=8CPersistanceThread=E7=9A=84=E6=97=A5?= =?UTF-8?q?=E5=BF=97=EF=BC=8C=E6=96=B9=E4=BE=BF=E6=8E=92=E6=9F=A5=E9=94=99?= =?UTF-8?q?=E8=AF=AF=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../persistance/PersistenceThread.java | 315 +++++++++--------- .../ServerHealthCollector.java | 6 +- 2 files changed, 169 insertions(+), 152 deletions(-) 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 b42f1a011..4f374bd65 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 @@ -15,170 +15,185 @@ import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.*; public class PersistenceThread extends Thread { - private Logger logger = LogManager.getLogger(PersistenceThread.class); + private Logger logger = LogManager.getLogger(PersistenceThread.class); - PersistenceThread(int threadIdx) { - super("PersistenceThread_" + threadIdx); - } + PersistenceThread(int threadIdx) { + super("PersistenceThread_" + threadIdx); + } - @Override - public void run() { - File file1 = null; - BufferedReader bufferedReader = null; - int offset; - while (true) { - try { - file1 = getDataFiles(); - if (file1 == null) { - try { - Thread.sleep(SWITCH_FILE_WAIT_TIME); - } catch (InterruptedException e) { - logger.error("Failure sleep", e); - } - continue; - } + @Override + public void run() { + File file1 = null; + BufferedReader bufferedReader = null; + int offset; + while (true) { + 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()) { - 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) { - MemoryRegister.instance().updateOffSet(file1.getName(), offset); - StorageChainController.doStorage(stringBuilder - .toString()); - stringBuilder.delete(0, stringBuilder.length()); - } + 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) { + MemoryRegister.instance().updateOffSet( + file1.getName(), offset); + StorageChainController.doStorage(stringBuilder + .toString()); + stringBuilder.delete(0, stringBuilder.length()); + } - try { - Thread.sleep(READ_ENDING_FILE_MAX_WAITE_TIME); - } catch (InterruptedException e) { - logger.error("Sleep failed", e); - } + try { + Thread.sleep(READ_ENDING_FILE_MAX_WAITE_TIME); + } catch (InterruptedException e) { + logger.error("Sleep failed", e); + } - continue; - } + continue; + } - //文件读入/n字符串 - if (tmpData.length() <= 0) { - // 加上回车的字符串长度 - offset += 1; - continue; - } + // 文件读入/n字符串 + if (tmpData.length() <= 0) { + // 加上回车的字符串长度 + offset += 1; + continue; + } - ServerHealthCollector.getCurrentHeathReading(null) - .updateData( - ServerHeathReading.INFO, - "read " + tmpData.length() - + " chars from local file:" + file1.getName()); + ServerHealthCollector.getCurrentHeathReading(null) + .updateData( + ServerHeathReading.INFO, + "read " + tmpData.length() + + " chars from local file:" + + file1.getName()); - if ("EOF".equals(tmpData)) { - if (stringBuilder != null && stringBuilder.length() > 0) { - StorageChainController.doStorage(stringBuilder - .toString()); - } + if ("EOF".equals(tmpData)) { + if (stringBuilder != null && stringBuilder.length() > 0) { + StorageChainController.doStorage(stringBuilder + .toString()); + } - 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")); + 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().removeEntry(file1.getName()); - break; - } + MemoryRegister.instance().removeEntry(file1.getName()); + break; + } - if (stringBuilder.length() + tmpData.length() >= MAX_STORAGE_SIZE_PER_TIME) { - StorageChainController.doStorage(stringBuilder - .toString()); - stringBuilder.delete(0, stringBuilder.length()); - MemoryRegister.instance().updateOffSet(file1.getName(), offset); - } + if (stringBuilder.length() + tmpData.length() >= MAX_STORAGE_SIZE_PER_TIME) { + StorageChainController.doStorage(stringBuilder + .toString()); + stringBuilder.delete(0, stringBuilder.length()); + MemoryRegister.instance().updateOffSet(file1.getName(), + offset); + } - 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 { - if (file1 != null) { - MemoryRegister.instance().unRegister(file1.getName()); - } - try { - if (bufferedReader != null) - bufferedReader.close(); - } catch (IOException e) { - logger.error("can't close data file", e); - } - } + stringBuilder.append(tmpData); + // 加上回车的字符串长度 + offset += tmpData.length() + 1; + } + } catch (FileNotFoundException e) { + logger.error("The data file could not be found.", e); + ServerHealthCollector.getCurrentHeathReading(null).updateData( + ServerHeathReading.ERROR, e.getMessage()); + } catch (IOException e) { + logger.error("The data file I/O exception.", e); + ServerHealthCollector.getCurrentHeathReading(null).updateData( + ServerHeathReading.ERROR, e.getMessage()); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + ServerHealthCollector.getCurrentHeathReading(null).updateData( + ServerHeathReading.ERROR, t.getMessage()); + } finally { + try{ + if (file1 != null) { + MemoryRegister.instance().unRegister(file1.getName()); + } + }catch (Throwable t) { + logger.error("unRegister file[{}] failure", file1.getName(), t); + } + try { + if (bufferedReader != null) + bufferedReader.close(); + } catch (IOException e) { + logger.error("can't close data file", e); + } + } - try { - Thread.sleep(SWITCH_FILE_WAIT_TIME); - } catch (InterruptedException e) { - logger.error("Failure sleep.", e); - } - } - } + 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) { - 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 int moveOffSet(File file1, BufferedReader bufferedReader) + throws IOException { + int offset = MemoryRegister.instance().getOffSet(file1.getName()); + if (-1 == offset || offset == 0) { + 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().doRegister(file.getName()) == null) { - 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; - } + 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().doRegister(file.getName()) == null) { + 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; - } + return file1; + } } 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 ea2e17a1f..a9d3656ec 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 @@ -47,12 +47,14 @@ public class ServerHealthCollector extends Thread { public void run() { while (true) { try { - String[] keyList = heathReadings.keySet().toArray(new String[0]); + Map heathReadingsSnapshot = heathReadings; + heathReadings = new ConcurrentHashMap(); + String[] keyList = heathReadingsSnapshot.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(heathReadingsSnapshot.get(key)).append("\n"); } log.append("------------------------------------------------\n");