diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java index 6fee9986b..617f341c6 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java @@ -3,6 +3,7 @@ package com.ai.cloud.skywalking.reciever.buffer; import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.processor.AbstractSpanProcessor; +import com.ai.cloud.skywalking.reciever.processor.IProcessor; import com.ai.cloud.skywalking.reciever.processor.ProcessorFactory; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; @@ -50,7 +51,7 @@ public class DataBufferThread extends Thread { } for (Map.Entry> entry : serializeObjects.entrySet()) { - AbstractSpanProcessor processor = ProcessorFactory.chooseProcessor(entry.getKey()); + IProcessor processor = ProcessorFactory.chooseProcessor(entry.getKey()); if (processor != null) { processor.process(entry.getValue()); } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AbstractSpanProcessor.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AbstractSpanProcessor.java index 5b7c4f746..95f5ec147 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AbstractSpanProcessor.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AbstractSpanProcessor.java @@ -48,6 +48,4 @@ public abstract class AbstractSpanProcessor implements IProcessor { public abstract void doSaveHBase(Connection connection, List serializedObjects); - public abstract int getType(); - } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java index 901809c23..14c78c715 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/AckSpanProcessor.java @@ -45,7 +45,7 @@ public class AckSpanProcessor extends AbstractSpanProcessor { } @Override - public int getType() { + public int getProtocolType() { return 2; } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/IProcessor.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/IProcessor.java index 82f9ab911..a095981b2 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/IProcessor.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/IProcessor.java @@ -5,5 +5,7 @@ import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; import java.util.List; public interface IProcessor { + int getProtocolType(); + void process(List serializedObjects); } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java index ed875d49c..e1694a12c 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/ProcessorFactory.java @@ -9,18 +9,18 @@ import java.util.ServiceLoader; public class ProcessorFactory { private static Logger logger = LogManager.getLogger(ProcessorFactory.class); - private static Map type_processor_mapping = new HashMap(); + private static Map type_processor_mapping = new HashMap(); static { - ServiceLoader processors = ServiceLoader.load(AbstractSpanProcessor.class); + ServiceLoader processors = ServiceLoader.load(IProcessor.class); - for (AbstractSpanProcessor processor : processors) { - logger.info("Init protocol type and processor mapping : {} --> {}.", processor.getType(), processor.getClass().getName()); - type_processor_mapping.put(processor.getType(), processor); + for (IProcessor processor : processors) { + logger.info("Init protocol type and processor mapping : {} --> {}.", processor.getProtocolType(), processor.getClass().getName()); + type_processor_mapping.put(processor.getProtocolType(), processor); } } - public static AbstractSpanProcessor chooseProcessor(int dataType) { + public static IProcessor chooseProcessor(int dataType) { return type_processor_mapping.get(dataType); } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java index c081c92fb..715186a37 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/processor/RequestSpanProcessor.java @@ -44,7 +44,7 @@ public class RequestSpanProcessor extends AbstractSpanProcessor { } @Override - public int getType() { + public int getProtocolType() { return 1; } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/util/HBaseUtil.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/util/HBaseUtil.java index baff69116..8d357bc15 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/util/HBaseUtil.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/util/HBaseUtil.java @@ -1,15 +1,19 @@ package com.ai.cloud.skywalking.reciever.util; +import com.ai.cloud.skywalking.reciever.processor.ProcessorFactory; import com.ai.cloud.skywalking.reciever.processor.exception.SaveToHBaseFailedException; import org.apache.hadoop.hbase.TableName; import org.apache.hadoop.hbase.client.Connection; import org.apache.hadoop.hbase.client.Put; import org.apache.hadoop.hbase.client.Table; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; import java.io.IOException; import java.util.List; public class HBaseUtil { + private static Logger logger = LogManager.getLogger(HBaseUtil.class); public static void batchSavePuts(Connection connection, String tableName, List puts) { Object[] resultArrays = new Object[puts.size()]; @@ -18,9 +22,9 @@ public class HBaseUtil { table.batch(puts, resultArrays); // ignore failed data } catch (IOException e) { - throw new SaveToHBaseFailedException(e); + logger.error("batchSavePuts failure.", e); } catch (InterruptedException e) { - throw new SaveToHBaseFailedException(e); + logger.error("batchSavePuts failure.", e); } } }