From d6f49bc931cefa70f7ff0e3a4b3f71be5de16532 Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Thu, 24 Dec 2015 17:41:16 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8Doffset=E6=96=87=E4=BB=B6?= =?UTF-8?q?=E4=B8=AD=E7=9A=84=E6=97=A0=E7=94=A8=E6=96=87=E4=BB=B6=E6=B2=A1?= =?UTF-8?q?=E6=9C=89=E8=A2=AB=E5=88=A0=E9=99=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../reciever/persistance/MemoryRegister.java | 23 +++++++++++++++---- 1 file changed, 19 insertions(+), 4 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 7147f31ee..7e29cb537 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 @@ -1,10 +1,13 @@ package com.ai.cloud.skywalking.reciever.persistance; +import com.ai.cloud.skywalking.reciever.conf.Config; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import java.io.*; +import java.util.Arrays; import java.util.Collection; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -24,7 +27,7 @@ public class MemoryRegister { if (logger.isDebugEnabled()) { logger.debug("Register entry[{}] into the memory register", fileName); } - if (entries.containsKey(fileName)){ + if (entries.containsKey(fileName)) { entries.get(fileName).setOffset(offset); } } @@ -39,7 +42,7 @@ public class MemoryRegister { } - public void removeEntry(String fileName){ + public void removeEntry(String fileName) { entries.remove(fileName); } @@ -96,6 +99,12 @@ public class MemoryRegister { private MemoryRegister() { BufferedReader reader; // 在处理数据之前需要初始化处理文件的处理状态 + + //去掉entries中无法与缓存数据文件匹配的文件 + File parentDir = new File(Config.Buffer.DATA_BUFFER_FILE_PARENT_DIRECTORY); + //上次未处理的缓存数据文件,entries内的数据主要以缓存 + List bufferFileNameList = Arrays.asList(parentDir.list()); + try { // 读取offset文件 file = new File(REGISTER_FILE_PARENT_DIRECTORY, REGISTER_FILE_NAME); @@ -108,7 +117,9 @@ public class MemoryRegister { 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)); + if (bufferFileNameList.contains(ss[0])) { + entries.put(ss[0], new FileRegisterEntry(ss[0], Integer.valueOf(ss[1]), FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); + } } } // 创建offset文件 @@ -120,12 +131,16 @@ public class MemoryRegister { while ((offsetData = reader.readLine()) != null && !"EOF".equals(offsetData)) { try { String[] ss = offsetData.split("\t"); - entries.put(ss[0], new FileRegisterEntry(ss[0], Integer.valueOf(ss[1]), FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); + if (bufferFileNameList.contains(ss[0])) { + entries.put(ss[0], new FileRegisterEntry(ss[0], Integer.valueOf(ss[1]), FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); + } } catch (Exception e) { continue; } } } + + } catch (FileNotFoundException e) { logger.error("The offset file does not exist.", e); checkOffSetExists();