diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java index 2ddaac3d4..2064fafc7 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java @@ -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()); } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/Chain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/Chain.java new file mode 100644 index 000000000..fc022a63d --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/Chain.java @@ -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 chains = new ArrayList(); + + 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); + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/IStorageChain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/IStorageChain.java new file mode 100644 index 000000000..87b54edc6 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/IStorageChain.java @@ -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); +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java new file mode 100644 index 000000000..8d850d144 --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java @@ -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); + } + } + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/hbase/HBaseOperator.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/HBaseOperator.java similarity index 98% rename from skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/hbase/HBaseOperator.java rename to skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/HBaseOperator.java index 3137277ac..a3695b773 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/hbase/HBaseOperator.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/HBaseOperator.java @@ -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; diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java new file mode 100644 index 000000000..61599c78a --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java @@ -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); + } + + + } + +}