From 6c4213a9702815857bdc22f3ab6591cf42772e4a Mon Sep 17 00:00:00 2001 From: ascrutae Date: Thu, 8 Dec 2016 14:39:13 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E9=83=A8=E5=88=86=E9=97=AE?= =?UTF-8?q?=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../storage/alarm/SpanAlarmHandler.java | 17 +++---- .../alarm/checker/ExceptionChecker.java | 4 +- .../alarm/checker/ExecuteTimeChecker.java | 7 +-- .../storage/alarm/checker/ISpanChecker.java | 4 +- .../eye/skywalking/storage/config/Config.java | 5 ++- .../storage/data/spandata/AckSpanData.java | 12 +++++ .../storage/listener/StorageListener.java | 3 +- .../a/eye/skywalking/storage/TestMain.java | 44 +++++++++++++++++++ .../storage/alarm/SpanAlarmHandlerTest.java | 19 ++++---- 9 files changed, 89 insertions(+), 26 deletions(-) create mode 100644 skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/TestMain.java diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandler.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandler.java index 5b4c688ef..e392c8fbc 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandler.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandler.java @@ -3,6 +3,7 @@ package com.a.eye.skywalking.storage.alarm; import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.storage.alarm.checker.*; import com.a.eye.skywalking.storage.alarm.sender.AlarmMessageSenderFactory; +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; import com.lmax.disruptor.EventHandler; import java.util.ArrayList; @@ -11,7 +12,7 @@ import java.util.List; /** * Created by xin on 2016/12/8. */ -public class SpanAlarmHandler implements EventHandler { +public class SpanAlarmHandler implements EventHandler { private List spanCheckers = new ArrayList(); public SpanAlarmHandler() { @@ -20,17 +21,17 @@ public class SpanAlarmHandler implements EventHandler { spanCheckers.add(new ExecuteTimePossibleErrorChecker()); } + private String generateAlarmMessageKey(AckSpanData span, FatalReason reason) { + return span.getUserName() + "-" + span.getApplicationCode() + "-" + (System.currentTimeMillis() / (10000 * 6)) + "-" + reason; + } + @Override - public void onEvent(AckSpan span, long sequence, boolean endOfBatch) throws Exception { + public void onEvent(AckSpanData spanData, long sequence, boolean endOfBatch) throws Exception { for (ISpanChecker spanChecker : spanCheckers) { - CheckResult result = spanChecker.check(span); + CheckResult result = spanChecker.check(spanData); if (!result.isPassed()) { - AlarmMessageSenderFactory.getSender().send(generateAlarmMessageKey(span, result.getFatalReason()), result.getMessage()); + AlarmMessageSenderFactory.getSender().send(generateAlarmMessageKey(spanData, result.getFatalReason()), result.getMessage()); } } } - - private String generateAlarmMessageKey(AckSpan span, FatalReason reason) { - return span.getUsername() + "-" + span.getApplicationCode() + "-" + (System.currentTimeMillis() / (10000 * 6)) + "-" + reason; - } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExceptionChecker.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExceptionChecker.java index c7e0dbc92..566473026 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExceptionChecker.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExceptionChecker.java @@ -1,12 +1,12 @@ package com.a.eye.skywalking.storage.alarm.checker; -import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.storage.config.Config; +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; public class ExceptionChecker implements ISpanChecker { @Override - public CheckResult check(AckSpan span) { + public CheckResult check(AckSpanData span) { if (span.getStatusCode() != 1) return new CheckResult(); String exceptionStack = span.getExceptionStack(); diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExecuteTimeChecker.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExecuteTimeChecker.java index b6b43ceb6..1f11d5755 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExecuteTimeChecker.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ExecuteTimeChecker.java @@ -1,6 +1,7 @@ package com.a.eye.skywalking.storage.alarm.checker; import com.a.eye.skywalking.network.grpc.AckSpan; +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; /** * Created by xin on 2016/12/8. @@ -8,7 +9,7 @@ import com.a.eye.skywalking.network.grpc.AckSpan; public abstract class ExecuteTimeChecker implements ISpanChecker { @Override - public CheckResult check(AckSpan span) { + public CheckResult check(AckSpanData span) { long cost = span.getCost(); if (isOverThreshold(cost)) { return new CheckResult(getFatalLevel(), generateAlarmMessage(span)); @@ -21,8 +22,8 @@ public abstract class ExecuteTimeChecker implements ISpanChecker { protected abstract FatalReason getFatalLevel(); - protected String generateAlarmMessage(AckSpan span) { - return span.getViewpointId() + " cost " + span.getCost() + " ms."; + protected String generateAlarmMessage(AckSpanData span) { + return span.getViewPointId() + " cost " + span.getCost() + " ms."; } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ISpanChecker.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ISpanChecker.java index ef364efa0..4eca2b55b 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ISpanChecker.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/alarm/checker/ISpanChecker.java @@ -1,10 +1,10 @@ package com.a.eye.skywalking.storage.alarm.checker; -import com.a.eye.skywalking.network.grpc.AckSpan; +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; /** * Created by xin on 2016/12/8. */ public interface ISpanChecker { - CheckResult check(AckSpan span); + CheckResult check(AckSpanData span); } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java index f911e39e1..7dc5e82be 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java @@ -8,7 +8,7 @@ public class Config { public static int PORT = 34000; } - public static class Disruptor{ + public static class Disruptor { public static int BUFFER_SIZE = 1024 * 128; public static int FLUSH_SIZE = 100; @@ -43,4 +43,7 @@ public class Config { public static String PATH_PREFIX = "/skywalking/storage_list/"; } + public static class Alarm { + public static int ALARM_EXCEPTION_STACK_LENGTH = 300; + } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java index ad8d81533..9441a815b 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/AckSpanData.java @@ -56,4 +56,16 @@ public class AckSpanData extends AbstractSpanData { public int getStatusCode() { return ackSpan.getStatusCode(); } + + public String getViewPointId(){ + return ackSpan.getViewpointId(); + } + + public String getUserName(){ + return ackSpan.getUsername(); + } + + public String getApplicationCode(){ + return ackSpan.getApplicationCode(); + } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java index 6fcd81817..a5f8c1551 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java @@ -7,6 +7,7 @@ import com.a.eye.skywalking.logging.api.LogManager; import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.network.grpc.RequestSpan; import com.a.eye.skywalking.network.listener.server.SpanStorageServerListener; +import com.a.eye.skywalking.storage.alarm.SpanAlarmHandler; import com.a.eye.skywalking.storage.config.Config; import com.a.eye.skywalking.storage.data.spandata.AckSpanData; import com.a.eye.skywalking.storage.data.spandata.RequestSpanData; @@ -35,7 +36,7 @@ public class StorageListener implements SpanStorageServerListener { requestSpanRingBuffer = requestSpanDisruptor.getRingBuffer(); ackSpanDisruptor = new Disruptor(new AckSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); - ackSpanDisruptor.handleEventsWith(new StoreAckSpanEventHandler()); + ackSpanDisruptor.handleEventsWith(new StoreAckSpanEventHandler(), new SpanAlarmHandler()); ackSpanDisruptor.start(); ackSpanRingBuffer = ackSpanDisruptor.getRingBuffer(); } diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/TestMain.java b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/TestMain.java new file mode 100644 index 000000000..ecdc8502d --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/TestMain.java @@ -0,0 +1,44 @@ +package com.a.eye.skywalking.storage; + +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; +import com.a.eye.skywalking.storage.config.Config; +import com.a.eye.skywalking.storage.data.spandata.RequestSpanData; +import com.lmax.disruptor.EventHandler; +import com.lmax.disruptor.RingBuffer; +import com.lmax.disruptor.dsl.Disruptor; +import com.lmax.disruptor.util.DaemonThreadFactory; + +/** + * Created by xin on 2016/12/7. + */ +public class TestMain { + public static void main(String[] args) throws InterruptedException { + Disruptor requestSpanDataDisruptor = null; + requestSpanDataDisruptor = new Disruptor(new StringBuilderFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); + requestSpanDataDisruptor.handleEventsWith(new EventHandler() { + @Override + public void onEvent(StringBuilder event, long sequence, boolean endOfBatch) throws Exception { + System.out.println("AA: " + event); + } + }, new EventHandler() { + @Override + public void onEvent(StringBuilder event, long sequence, boolean endOfBatch) throws Exception { + System.out.println("BB: " + event); + } + }); + requestSpanDataDisruptor.start(); + RingBuffer stringBuilderRingBuffer = requestSpanDataDisruptor.getRingBuffer(); + + long sequence = stringBuilderRingBuffer.next(); // Grab the next sequence + try { + StringBuilder data = stringBuilderRingBuffer.get(sequence); + data.append("A"); + } catch (Exception e) { + } finally { + stringBuilderRingBuffer.publish(sequence); + } + + Thread.sleep(1000); + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandlerTest.java b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandlerTest.java index 25a658734..c2696e41f 100644 --- a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandlerTest.java +++ b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/alarm/SpanAlarmHandlerTest.java @@ -4,6 +4,7 @@ import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.network.grpc.TraceId; import com.a.eye.skywalking.storage.alarm.sender.AlarmMessageSender; import com.a.eye.skywalking.storage.alarm.sender.AlarmMessageSenderFactory; +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -24,10 +25,10 @@ public class SpanAlarmHandlerTest { private AlarmMessageSender messageHandler; @InjectMocks private SpanAlarmHandler handler; - private AckSpan normalAckSpan; - private AckSpan costMuchSpan; - private AckSpan costTooMuchSpan; - private AckSpan exceptionSpan; + private AckSpanData normalAckSpan; + private AckSpanData costMuchSpan; + private AckSpanData costTooMuchSpan; + private AckSpanData exceptionSpan; @Before public void setUp() { @@ -39,10 +40,10 @@ public class SpanAlarmHandlerTest { .addSegments(2016).addSegments(startTime).addSegments(2).addSegments(100).addSegments(30) .addSegments(1).build()); - normalAckSpan = builder.build(); - costMuchSpan = builder.setCost(600).build(); - costTooMuchSpan = builder.setCost(4000).build(); - exceptionSpan = builder.setCost(20).setStatusCode(1).setExceptionStack("occur exception").build(); + normalAckSpan = new AckSpanData(builder.build()); + costMuchSpan = new AckSpanData(builder.setCost(600).build()); + costTooMuchSpan = new AckSpanData(builder.setCost(4000).build()); + exceptionSpan = new AckSpanData(builder.setCost(20).setStatusCode(1).setExceptionStack("occur exception").build()); } @Test @@ -65,7 +66,7 @@ public class SpanAlarmHandlerTest { @Test public void testCostTooMuchSpan() throws Exception { - handler.onEvent(costTooMuchSpan,1, false); + handler.onEvent(costTooMuchSpan, 1, false); verify(messageHandler, times(1)).send(any(), anyString()); } }