parent
d545b24fcd
commit
f22986e30e
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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";
|
||||
}
|
||||
}
|
||||
|
|
@ -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 = ",";
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<byte[]> {
|
||||
|
||||
|
||||
@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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in New Issue