修复Span在StorageChain中发生异常时,程序进入死循环

This commit is contained in:
ascrutae 2016-02-25 15:54:13 +08:00
parent 6421ba0a30
commit 5b23fd09c3
4 changed files with 150 additions and 140 deletions

View File

@ -86,6 +86,8 @@ public class Config {
public static long RETRY_STORAGE_WAIT_TIME = 50L;
public static String STORAGE_TYPE = "hbase";
public static int RETRY_STORAGE_TIMES = 3;
}
public static class Alarm {

View File

@ -49,12 +49,13 @@ public class StorageChainController {
}
}
while (true) {
int retryTimes = 0;
while(retryTimes++ < Config.StorageChain.RETRY_STORAGE_TIMES) {
try {
Chain chain = new Chain(chainArray);
chain.doChain(spans);
break;
} catch (Throwable e) {
// 主要的异常可能跟环境有关系比如Redis,HBase将会重试N次
try {
Thread.sleep(Config.StorageChain.RETRY_STORAGE_WAIT_TIME);
} catch (InterruptedException e1) {

View File

@ -13,138 +13,139 @@ import redis.clients.jedis.exceptions.JedisConnectionException;
import java.util.List;
import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXPIRE_SECONDS;
import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXCEPTION_STACK_LENGTH;
import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXPIRE_SECONDS;
public class AlarmChain implements IStorageChain {
private static Logger logger = LogManager.getLogger(AlarmChain.class);
private static JedisPool jedisPool;
private static String[] config;
private static Object lock = new Object();
private static RedisInspector connector = new RedisInspector();;
private static Logger logger = LogManager.getLogger(AlarmChain.class);
private static JedisPool jedisPool;
private static String[] config;
private static Object lock = new Object();
private static RedisInspector connector = new RedisInspector();
;
static {
GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig();
String redisServerConfig = Config.Alarm.REDIS_SERVER;
if (redisServerConfig == null || redisServerConfig.length() <= 0) {
logger.error("Redis server is not setting. Switch off alarm module. ");
Config.Alarm.ALARM_OFF_FLAG = true;
} else {
config = redisServerConfig.split(":");
if (config.length != 2) {
logger.error("Redis server address is illegal setting, need to be 'ip:port'. Switch off alarm module. ");
Config.Alarm.ALARM_OFF_FLAG = true;
} else {
jedisPool = new JedisPool(genericObjectPoolConfig, config[0],
Integer.valueOf(config[1]));
// Test connect redis.
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
} catch (Exception e) {
handleFailedToConnectRedisServerException(e);
logger.error("can't connect to redis["
+ Config.Alarm.REDIS_SERVER + "]", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
}
}
static {
GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig();
String redisServerConfig = Config.Alarm.REDIS_SERVER;
if (redisServerConfig == null || redisServerConfig.length() <= 0) {
logger.error("Redis server is not setting. Switch off alarm module. ");
Config.Alarm.ALARM_OFF_FLAG = true;
} else {
config = redisServerConfig.split(":");
if (config.length != 2) {
logger.error("Redis server address is illegal setting, need to be 'ip:port'. Switch off alarm module. ");
Config.Alarm.ALARM_OFF_FLAG = true;
} else {
jedisPool = new JedisPool(genericObjectPoolConfig, config[0],
Integer.valueOf(config[1]));
// Test connect redis.
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
} catch (Exception e) {
handleFailedToConnectRedisServerException(e);
logger.error("can't connect to redis["
+ Config.Alarm.REDIS_SERVER + "]", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
}
}
@Override
public void doChain(List<Span> spans, Chain chain) {
for (Span span : spans) {
if (span.getStatusCode() != 1)
continue;
String exceptionStack = span.getExceptionStack();
if(exceptionStack == null){
exceptionStack = "";
}else if(exceptionStack.length() > ALARM_EXCEPTION_STACK_LENGTH){
exceptionStack = exceptionStack.substring(0, ALARM_EXCEPTION_STACK_LENGTH);
}
saveAlarmMessage(generateAlarmKey(span), span.getTraceId(), exceptionStack);
}
chain.doChain(spans);
}
@Override
public void doChain(List<Span> spans, Chain chain) {
for (Span span : spans) {
if (span.getStatusCode() != 1)
continue;
String exceptionStack = span.getExceptionStack();
if (exceptionStack == null) {
exceptionStack = "";
} else if (exceptionStack.length() > ALARM_EXCEPTION_STACK_LENGTH) {
exceptionStack = exceptionStack.substring(0, ALARM_EXCEPTION_STACK_LENGTH);
}
saveAlarmMessage(generateAlarmKey(span), span.getTraceId(), exceptionStack);
}
chain.doChain(spans);
}
private String generateAlarmKey(Span span) {
return span.getUserId() + "-" + span.getApplicationId() + "-"
+ (System.currentTimeMillis() / (10000 * 6));
}
private String generateAlarmKey(Span span) {
return span.getUserId() + "-" + span.getApplicationId() + "-"
+ (System.currentTimeMillis() / (10000 * 6));
}
private void saveAlarmMessage(String key, String traceId, String exceptionMsgOutline) {
if (Config.Alarm.ALARM_OFF_FLAG) {
return;
}
private void saveAlarmMessage(String key, String traceId, String exceptionMsgOutline) {
if (Config.Alarm.ALARM_OFF_FLAG) {
return;
}
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
jedis.hsetnx(key, traceId, exceptionMsgOutline);
jedis.expire(key, ALARM_EXPIRE_SECONDS);
} catch (Exception e) {
handleFailedToConnectRedisServerException(e);
logger.error("Failed to set data.", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
jedis.hsetnx(key, traceId, exceptionMsgOutline);
jedis.expire(key, ALARM_EXPIRE_SECONDS);
} catch (Exception e) {
handleFailedToConnectRedisServerException(e);
logger.error("Failed to set data.", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
private static class RedisInspector extends Thread {
@Override
public void run() {
logger.debug("Connecting to redis....");
Jedis jedis;
while (true) {
try {
jedisPool = new JedisPool(buildGenericObjectPoolConfig(),
config[0], Integer.valueOf(config[1]));
jedis = jedisPool.getResource();
jedis.get("ok");
break;
} catch (Exception e) {
if (e instanceof JedisConnectionException) {
try {
Thread.sleep(5000L);
} catch (InterruptedException e1) {
logger.error("Sleep failed", e);
}
continue;
}
}
}
logger.debug("Connected to redis success. Open alarm function.");
Config.Alarm.ALARM_OFF_FLAG = false;
// 清理当前线程
connector = null;
}
}
private static class RedisInspector extends Thread {
@Override
public void run() {
logger.debug("Connecting to redis....");
Jedis jedis;
while (true) {
try {
jedisPool = new JedisPool(buildGenericObjectPoolConfig(),
config[0], Integer.valueOf(config[1]));
jedis = jedisPool.getResource();
jedis.get("ok");
break;
} catch (Exception e) {
if (e instanceof JedisConnectionException) {
try {
Thread.sleep(5000L);
} catch (InterruptedException e1) {
logger.error("Sleep failed", e);
}
continue;
}
}
}
logger.debug("Connected to redis success. Open alarm function.");
Config.Alarm.ALARM_OFF_FLAG = false;
// 清理当前线程
connector = null;
}
}
private static void handleFailedToConnectRedisServerException(Exception e) {
if (e instanceof JedisConnectionException) {
// 发生连接不上Redis
if (connector == null || !connector.isAlive()) {
synchronized (lock) {
if (!connector.isAlive()) {
// 启动巡检线程
connector.start();
}
}
}
}
}
private static void handleFailedToConnectRedisServerException(Exception e) {
if (e instanceof JedisConnectionException) {
// 发生连接不上Redis
if (connector == null || !connector.isAlive()) {
synchronized (lock) {
if (!connector.isAlive()) {
// 启动巡检线程
connector.start();
}
}
}
}
}
private static GenericObjectPoolConfig buildGenericObjectPoolConfig() {
GenericObjectPoolConfig genericObjectPoolConfig = new GenericObjectPoolConfig();
genericObjectPoolConfig.setTestOnBorrow(true);
genericObjectPoolConfig.setMaxIdle(Config.Alarm.REDIS_MAX_IDLE);
genericObjectPoolConfig.setMinIdle(Config.Alarm.REDIS_MIN_IDLE);
genericObjectPoolConfig.setMaxTotal(Config.Alarm.REDIS_MAX_TOTAL);
return genericObjectPoolConfig;
}
private static GenericObjectPoolConfig buildGenericObjectPoolConfig() {
GenericObjectPoolConfig genericObjectPoolConfig = new GenericObjectPoolConfig();
genericObjectPoolConfig.setTestOnBorrow(true);
genericObjectPoolConfig.setMaxIdle(Config.Alarm.REDIS_MAX_IDLE);
genericObjectPoolConfig.setMinIdle(Config.Alarm.REDIS_MIN_IDLE);
genericObjectPoolConfig.setMaxTotal(Config.Alarm.REDIS_MAX_TOTAL);
return genericObjectPoolConfig;
}
}

View File

@ -77,24 +77,30 @@ public class SaveToHBaseChain implements IStorageChain {
if (spans == null || spans.size() <= 0)
return;
List<Put> puts = new ArrayList<Put>();
Put put;
Put put = null;
String columnName;
for (Span span : spans) {
put = new Put(Bytes.toBytes(span.getTraceId()), getTSBySpanTraceId(span));
if (StringUtils.isEmpty(span.getParentLevel().trim())) {
columnName = span.getLevelId() + "";
if (span.isReceiver()) {
columnName = span.getLevelId() + "-S";
try {
put = new Put(Bytes.toBytes(span.getTraceId()), getTSBySpanTraceId(span));
if (StringUtils.isEmpty(span.getParentLevel().trim())) {
columnName = span.getLevelId() + "";
if (span.isReceiver()) {
columnName = span.getLevelId() + "-S";
}
put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName),
Bytes.toBytes(span.getOriginData()));
} else {
columnName = span.getParentLevel() + "." + span.getLevelId();
if (span.isReceiver()) {
columnName = span.getParentLevel() + "." + span.getLevelId() + "-S";
}
put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName),
Bytes.toBytes(span.getOriginData()));
}
put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName),
Bytes.toBytes(span.getOriginData()));
} else {
columnName = span.getParentLevel() + "." + span.getLevelId();
if (span.isReceiver()) {
columnName = span.getParentLevel() + "." + span.getLevelId() + "-S";
}
put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName),
Bytes.toBytes(span.getOriginData()));
} catch (Exception e) {
// 不合规范的数据
logger.error("Failed to convert Span[" + span.getTraceId() + "] to put Object", e);
continue;
}
puts.add(put);
}