解决RegisterPersistenceThread没有合理的使用offset文件和offset.bak文件
This commit is contained in:
parent
9bfc3ce1a8
commit
ccbfbe578f
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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<FileRegisterEntry> 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<FileRegisterEntry> 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.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue