提交部分代码修改,移除不正确的异常处理。
This commit is contained in:
parent
beb6276adf
commit
c83d684ab7
|
|
@ -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<Integer, List<AbstractDataSerializable>> entry : serializeObjects.entrySet()) {
|
||||
AbstractSpanProcessor processor = ProcessorFactory.chooseProcessor(entry.getKey());
|
||||
IProcessor processor = ProcessorFactory.chooseProcessor(entry.getKey());
|
||||
if (processor != null) {
|
||||
processor.process(entry.getValue());
|
||||
}
|
||||
|
|
|
|||
|
|
@ -48,6 +48,4 @@ public abstract class AbstractSpanProcessor implements IProcessor {
|
|||
|
||||
public abstract void doSaveHBase(Connection connection, List<AbstractDataSerializable> serializedObjects);
|
||||
|
||||
public abstract int getType();
|
||||
|
||||
}
|
||||
|
|
|
|||
|
|
@ -45,7 +45,7 @@ public class AckSpanProcessor extends AbstractSpanProcessor {
|
|||
}
|
||||
|
||||
@Override
|
||||
public int getType() {
|
||||
public int getProtocolType() {
|
||||
return 2;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,5 +5,7 @@ import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable;
|
|||
import java.util.List;
|
||||
|
||||
public interface IProcessor {
|
||||
int getProtocolType();
|
||||
|
||||
void process(List<AbstractDataSerializable> serializedObjects);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -9,18 +9,18 @@ import java.util.ServiceLoader;
|
|||
|
||||
public class ProcessorFactory {
|
||||
private static Logger logger = LogManager.getLogger(ProcessorFactory.class);
|
||||
private static Map<Integer, AbstractSpanProcessor> type_processor_mapping = new HashMap<Integer, AbstractSpanProcessor>();
|
||||
private static Map<Integer, IProcessor> type_processor_mapping = new HashMap<Integer, IProcessor>();
|
||||
|
||||
static {
|
||||
ServiceLoader<AbstractSpanProcessor> processors = ServiceLoader.load(AbstractSpanProcessor.class);
|
||||
ServiceLoader<IProcessor> 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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -44,7 +44,7 @@ public class RequestSpanProcessor extends AbstractSpanProcessor {
|
|||
}
|
||||
|
||||
@Override
|
||||
public int getType() {
|
||||
public int getProtocolType() {
|
||||
return 1;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<Put> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue