From ccbfbe578fc67acd97e7dee0e44314dc009249c1 Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Mon, 30 Nov 2015 22:59:36 +0800 Subject: [PATCH] =?UTF-8?q?=E8=A7=A3=E5=86=B3RegisterPersistenceThread?= =?UTF-8?q?=E6=B2=A1=E6=9C=89=E5=90=88=E7=90=86=E7=9A=84=E4=BD=BF=E7=94=A8?= =?UTF-8?q?offset=E6=96=87=E4=BB=B6=E5=92=8Coffset.bak=E6=96=87=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../reciever/persistance/MemoryRegister.java | 34 +++-- .../RegisterPersistenceThread.java | 131 ++++++++++-------- 2 files changed, 95 insertions(+), 70 deletions(-) diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/MemoryRegister.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/MemoryRegister.java index 1d9c29575..230629d0c 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/MemoryRegister.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/MemoryRegister.java @@ -8,8 +8,7 @@ import java.util.Collection; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; -import static com.ai.cloud.skywalking.reciever.conf.Config.RegisterPersistence.REGISTER_FILE_NAME; -import static com.ai.cloud.skywalking.reciever.conf.Config.RegisterPersistence.REGISTER_FILE_PARENT_DIRECTORY; +import static com.ai.cloud.skywalking.reciever.conf.Config.RegisterPersistence.*; public class MemoryRegister { private Logger logger = LogManager.getLogger(MemoryRegister.class); @@ -72,16 +71,33 @@ public class MemoryRegister { } private MemoryRegister() { - // 读取offset文件 - checkOffSetExists(); BufferedReader reader; // 在处理数据之前需要初始化处理文件的处理状态 try { - reader = new BufferedReader(new FileReader(file)); - String offsetData; - while ((offsetData = reader.readLine()) != null && !"EOF".equals(offsetData)) { - String[] ss = offsetData.split("\t"); - entries.put(ss[0], new FileRegisterEntry(ss[0], Integer.valueOf(ss[1]), FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); + // 读取offset文件 + file = new File(REGISTER_FILE_PARENT_DIRECTORY, REGISTER_FILE_NAME); + // offset File不存在 + if (!file.exists()) { + File offsetBackUpFile = new File(REGISTER_FILE_PARENT_DIRECTORY, REGISTER_BAK_FILE_NAME); + // offset备份文件存在 + if (offsetBackUpFile.exists()) { + reader = new BufferedReader(new FileReader(offsetBackUpFile)); + String offsetData; + while ((offsetData = reader.readLine()) != null && !"EOF".equals(offsetData)) { + String[] ss = offsetData.split("\t"); + entries.put(ss[0], new FileRegisterEntry(ss[0], Integer.valueOf(ss[1]), FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); + } + } + // 创建offset文件 + file.createNewFile(); + } else { + // 如果存在 + reader = new BufferedReader(new FileReader(file)); + String offsetData; + while ((offsetData = reader.readLine()) != null && !"EOF".equals(offsetData)) { + String[] ss = offsetData.split("\t"); + entries.put(ss[0], new FileRegisterEntry(ss[0], Integer.valueOf(ss[1]), FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); + } } } catch (FileNotFoundException e) { logger.error("The offset file does not exist.", e); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/RegisterPersistenceThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/RegisterPersistenceThread.java index 9d17daf26..4e2a76e50 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/RegisterPersistenceThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/RegisterPersistenceThread.java @@ -1,11 +1,9 @@ package com.ai.cloud.skywalking.reciever.persistance; -import org.apache.commons.io.FileUtils; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; import java.io.BufferedWriter; import java.io.File; @@ -17,68 +15,79 @@ import static com.ai.cloud.skywalking.reciever.conf.Config.RegisterPersistence.* public class RegisterPersistenceThread extends Thread { - private Logger logger = LogManager - .getLogger(RegisterPersistenceThread.class); + private Logger logger = LogManager + .getLogger(RegisterPersistenceThread.class); - private BufferedWriter writer; + private BufferedWriter writer; - public RegisterPersistenceThread() { - super("RegisterPersistenceThread"); - } + public RegisterPersistenceThread() { + super("RegisterPersistenceThread"); + } - @Override - public void run() { - while (true) { - try { - Thread.sleep(OFFSET_WRITTEN_FILE_WAIT_CYCLE); - } catch (InterruptedException e) { - logger.error("Sleep failure", e); - } + @Override + public void run() { + while (true) { + try { + Thread.sleep(OFFSET_WRITTEN_FILE_WAIT_CYCLE); + } catch (InterruptedException e) { + logger.error("Sleep failure", e); + } - File file = new File(REGISTER_FILE_PARENT_DIRECTORY, - REGISTER_FILE_NAME); - File bakFile = new File(REGISTER_FILE_PARENT_DIRECTORY, - REGISTER_BAK_FILE_NAME); - try { - FileUtils.copyFile(file, bakFile); - } catch (IOException e) { - logger.error("Sleep failure", e); - } + try { + File file = new File(REGISTER_FILE_PARENT_DIRECTORY, + REGISTER_FILE_NAME); + File bakFile = new File(REGISTER_FILE_PARENT_DIRECTORY, + REGISTER_BAK_FILE_NAME); + // 先删除备份文件 + if (bakFile.exists()) { + bakFile.delete(); + } - Collection fileRegisterEntries = MemoryRegister - .instance().getEntries(); - logger.debug("file Register Entries size [{}]", - fileRegisterEntries.size()); - try { - writer = new BufferedWriter(new FileWriter(file)); - } catch (IOException e) { - logger.error("Write The offset file anomalies."); - } + // 将文件改名字 + file.renameTo(bakFile); - for (FileRegisterEntry fileRegisterEntry : fileRegisterEntries) { - try { - writer.write(fileRegisterEntry.toString() + "\n"); - } catch (IOException e) { - logger.error( - "Write file register entry to offset file failure", - e); - } - } - try { - writer.write("EOF\n"); - writer.flush(); - } catch (IOException e) { - logger.error("Flush offset file failure", e); - } finally { - try { - writer.close(); - } catch (IOException e) { - logger.error("close offset file failure", e); - } - } + // + if (!file.exists()) { + file.createNewFile(); + } - ServerHealthCollector.getCurrentHeathReading(null).updateData( - ServerHeathReading.INFO, "flush memory register to file."); - } - } + Collection fileRegisterEntries = MemoryRegister + .instance().getEntries(); + logger.debug("file Register Entries size [{}]", + fileRegisterEntries.size()); + try { + writer = new BufferedWriter(new FileWriter(file)); + } catch (IOException e) { + logger.error("Write The offset file anomalies."); + } + + for (FileRegisterEntry fileRegisterEntry : fileRegisterEntries) { + try { + writer.write(fileRegisterEntry.toString() + "\n"); + } catch (IOException e) { + logger.error( + "Write file register entry to offset file failure", e); + } + } + try { + writer.write("EOF\n"); + writer.flush(); + } catch (IOException e) { + logger.error("Flush offset file failure", e); + } finally { + try { + writer.close(); + } catch (IOException e) { + logger.error("close offset file failure", e); + } + } + } catch (IOException e) { + logger.error("Failed to back up offset file.", e); + } + + + ServerHealthCollector.getCurrentHeathReading(null).updateData( + ServerHeathReading.INFO, "flush memory register to file."); + } + } }