1.修改了全新的链状存储模式,用于存储和后续相关操作的代码隔离
This commit is contained in:
parent
2596de4306
commit
9994fb383b
|
|
@ -1,17 +1,22 @@
|
|||
package com.ai.cloud.skywalking.reciever.persistance;
|
||||
|
||||
import com.ai.cloud.skywalking.reciever.conf.Config;
|
||||
import com.ai.cloud.skywalking.reciever.hbase.HBaseOperator;
|
||||
import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry;
|
||||
import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.OFFSET_FILE_READ_BUFFER_SIZE;
|
||||
import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.OFFSET_FILE_SKIP_LENGTH;
|
||||
import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.SWITCH_FILE_WAIT_TIME;
|
||||
|
||||
import java.io.BufferedReader;
|
||||
import java.io.File;
|
||||
import java.io.FileNotFoundException;
|
||||
import java.io.FileReader;
|
||||
import java.io.IOException;
|
||||
|
||||
import org.apache.commons.io.FileUtils;
|
||||
import org.apache.commons.io.comparator.NameFileComparator;
|
||||
import org.apache.commons.lang.StringUtils;
|
||||
import org.apache.logging.log4j.LogManager;
|
||||
import org.apache.logging.log4j.Logger;
|
||||
|
||||
import java.io.*;
|
||||
|
||||
import static com.ai.cloud.skywalking.reciever.conf.Config.Persistence.*;
|
||||
import com.ai.cloud.skywalking.reciever.conf.Config;
|
||||
import com.ai.cloud.skywalking.reciever.storage.StorageChainController;
|
||||
|
||||
public class PersistenceThread extends Thread {
|
||||
|
||||
|
|
@ -24,7 +29,6 @@ public class PersistenceThread extends Thread {
|
|||
BufferedReader bufferedReader;
|
||||
int offset;
|
||||
StringBuffer data;
|
||||
String[] buriedPointData;
|
||||
while (true) {
|
||||
file1 = getDataFiles();
|
||||
if (file1 == null) {
|
||||
|
|
@ -72,16 +76,8 @@ public class PersistenceThread extends Thread {
|
|||
break;
|
||||
}
|
||||
|
||||
buriedPointData = data.toString().split(";");
|
||||
for (String buriedPoint : buriedPointData) {
|
||||
BuriedPointEntry entry = BuriedPointEntry.convert(buriedPoint);
|
||||
if (StringUtils.isEmpty(entry.getParentLevel().trim())) {
|
||||
HBaseOperator.insert(entry.getTraceId(), String.valueOf(entry.getLevelId()), buriedPoint);
|
||||
} else {
|
||||
HBaseOperator.insert(entry.getTraceId(), entry.getParentLevel() + "." + entry.getLevelId(), buriedPoint);
|
||||
}
|
||||
}
|
||||
|
||||
StorageChainController.doStorage(data.toString());
|
||||
|
||||
data.delete(0, data.length());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,22 @@
|
|||
package com.ai.cloud.skywalking.reciever.storage;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry;
|
||||
|
||||
public class Chain {
|
||||
private List<IStorageChain> chains = new ArrayList<IStorageChain>();
|
||||
|
||||
private int index = 0;
|
||||
|
||||
public void doChain(BuriedPointEntry entry, String entryOriginData ){
|
||||
if(index < chains.size()){
|
||||
chains.get(index++).doChain(entry, entryOriginData, this);;
|
||||
}
|
||||
}
|
||||
|
||||
synchronized void addChain(IStorageChain chain){
|
||||
chains.add(chain);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
package com.ai.cloud.skywalking.reciever.storage;
|
||||
|
||||
import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry;
|
||||
|
||||
public interface IStorageChain {
|
||||
public void doChain(BuriedPointEntry entry, String entryOriginData, Chain chain);
|
||||
}
|
||||
|
|
@ -0,0 +1,32 @@
|
|||
package com.ai.cloud.skywalking.reciever.storage;
|
||||
|
||||
import org.apache.logging.log4j.LogManager;
|
||||
import org.apache.logging.log4j.Logger;
|
||||
|
||||
import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry;
|
||||
|
||||
public class StorageChainController {
|
||||
private static Logger logger = LogManager
|
||||
.getLogger(StorageChainController.class);
|
||||
|
||||
private static Chain globalChain = new Chain();
|
||||
|
||||
public static void doStorage(String buriedPointDatas) {
|
||||
String[] buriedPointData = buriedPointDatas.split(";");
|
||||
if (buriedPointData == null || buriedPointData.length == 0) {
|
||||
return;
|
||||
}
|
||||
for (String buriedPoint : buriedPointData) {
|
||||
try {
|
||||
if(buriedPoint == null || buriedPoint.trim().length() == 0){
|
||||
continue;
|
||||
}
|
||||
BuriedPointEntry entry = BuriedPointEntry.convert(buriedPoint);
|
||||
globalChain.doChain(entry, buriedPoint);
|
||||
} catch (Throwable e) {
|
||||
logger.error("ready to save buriedPoint error, choose to ignore. data="
|
||||
+ buriedPoint, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package com.ai.cloud.skywalking.reciever.hbase;
|
||||
package com.ai.cloud.skywalking.reciever.storage.chain;
|
||||
|
||||
import com.ai.cloud.skywalking.reciever.conf.Config;
|
||||
import com.ai.cloud.skywalking.reciever.conf.ConfigInitializer;
|
||||
|
|
@ -0,0 +1,22 @@
|
|||
package com.ai.cloud.skywalking.reciever.storage.chain;
|
||||
|
||||
import org.apache.commons.lang.StringUtils;
|
||||
|
||||
import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry;
|
||||
import com.ai.cloud.skywalking.reciever.storage.Chain;
|
||||
import com.ai.cloud.skywalking.reciever.storage.IStorageChain;
|
||||
|
||||
public class SaveToHBaseChain implements IStorageChain{
|
||||
|
||||
@Override
|
||||
public void doChain(BuriedPointEntry entry, String entryOriginData, Chain chain) {
|
||||
if (StringUtils.isEmpty(entry.getParentLevel().trim())) {
|
||||
HBaseOperator.insert(entry.getTraceId(), String.valueOf(entry.getLevelId()), entryOriginData);
|
||||
} else {
|
||||
HBaseOperator.insert(entry.getTraceId(), entry.getParentLevel() + "." + entry.getLevelId(), entryOriginData);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
Loading…
Reference in New Issue