diff --git a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanClearEventHandler.java b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanClearEventHandler.java new file mode 100644 index 000000000..4729e0ed5 --- /dev/null +++ b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanClearEventHandler.java @@ -0,0 +1,14 @@ +package com.a.eye.skywalking.routing.disruptor.ack; + +import com.lmax.disruptor.EventHandler; + +/** + * Created by xin on 2017/2/8. + */ +public class AckSpanClearEventHandler implements EventHandler { + + @Override + public void onEvent(AckSpanHolder event, long sequence, boolean endOfBatch) throws Exception { + event.setAckSpan(null); + } +} diff --git a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanDisruptor.java b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanDisruptor.java index ac70cb0e8..3d2e0d65e 100644 --- a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanDisruptor.java +++ b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/AckSpanDisruptor.java @@ -26,7 +26,7 @@ public class AckSpanDisruptor extends AbstractSpanDisruptor { public AckSpanDisruptor(String connectionURL) { ackSpanDisruptor = new Disruptor(new AckSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); ackSpanEventHandler = new RouteAckSpanBufferEventHandler(connectionURL); - ackSpanDisruptor.handleEventsWith(ackSpanEventHandler, new SpanAlarmHandler()); + ackSpanDisruptor.handleEventsWith(ackSpanEventHandler).then(new SpanAlarmHandler()).then(new AckSpanClearEventHandler()); ackSpanDisruptor.start(); ackSpanRingBuffer = ackSpanDisruptor.getRingBuffer(); } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanDataHolder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanDataHolder.java new file mode 100644 index 000000000..fd8adcc6a --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanDataHolder.java @@ -0,0 +1,23 @@ +package com.a.eye.skywalking.storage.disruptor.ack; + +import com.a.eye.skywalking.network.grpc.AckSpan; +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; + +/** + * @author zhangxin + */ +public class AckSpanDataHolder { + private AckSpanData ackSpanData; + + public AckSpanData getAckSpanData() { + return ackSpanData; + } + + public void clearData() { + this.ackSpanData = null; + } + + public void fillData(AckSpan ackSpan) { + this.ackSpanData = new AckSpanData(ackSpan); + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanFactory.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanFactory.java index 5bcebd928..4ac5d6d0d 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanFactory.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanFactory.java @@ -6,9 +6,9 @@ import com.lmax.disruptor.EventFactory; /** * Created by wusheng on 2016/11/24. */ -public class AckSpanFactory implements EventFactory { +public class AckSpanFactory implements EventFactory { @Override - public AckSpanData newInstance() { - return new AckSpanData(); + public AckSpanDataHolder newInstance() { + return new AckSpanDataHolder(); } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/StoreAckSpanEventHandler.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/StoreAckSpanEventHandler.java index fa1916539..e8bddefcd 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/StoreAckSpanEventHandler.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/StoreAckSpanEventHandler.java @@ -11,6 +11,7 @@ import com.a.eye.skywalking.storage.data.index.IndexOperator; import com.a.eye.skywalking.storage.data.index.IndexOperatorFactory; import com.a.eye.skywalking.storage.data.spandata.AckSpanData; import com.a.eye.skywalking.storage.data.spandata.SpanData; +import com.a.eye.skywalking.storage.disruptor.request.RequestSpanDataHolder; import com.lmax.disruptor.EventHandler; import java.util.ArrayList; @@ -19,7 +20,7 @@ import java.util.List; /** * Created by wusheng on 2016/11/24. */ -public class StoreAckSpanEventHandler implements EventHandler { +public class StoreAckSpanEventHandler implements EventHandler { private static ILog logger = LogManager.getLogger(StoreAckSpanEventHandler.class); private DataFileWriter fileWriter; private IndexOperator operator; @@ -34,21 +35,25 @@ public class StoreAckSpanEventHandler implements EventHandler { } @Override - public void onEvent(AckSpanData event, long sequence, boolean endOfBatch) throws Exception { - buffer.add(event); + public void onEvent(AckSpanDataHolder event, long sequence, boolean endOfBatch) throws Exception { + try { + buffer.add(event.getAckSpanData()); - if (endOfBatch || buffer.size() == bufferSize) { - try { - IndexMetaCollection collection = fileWriter.write(buffer); + if (endOfBatch || buffer.size() == bufferSize) { + try { + IndexMetaCollection collection = fileWriter.write(buffer); - operator.batchUpdate(collection); - HealthCollector.getCurrentHeathReading("StoreAckSpanEventHandler").updateData(HeathReading.INFO, "Batch consume %s messages successfully.", buffer.size()); - } catch (Throwable e) { - logger.error("Ack messages consume failure.", e); - HealthCollector.getCurrentHeathReading("StoreAckSpanEventHandler").updateData(HeathReading.ERROR, "Batch consume %s messages failure.", buffer.size()); - } finally { - buffer.clear(); + operator.batchUpdate(collection); + HealthCollector.getCurrentHeathReading("StoreAckSpanEventHandler").updateData(HeathReading.INFO, "Batch consume %s messages successfully.", buffer.size()); + } catch (Throwable e) { + logger.error("Ack messages consume failure.", e); + HealthCollector.getCurrentHeathReading("StoreAckSpanEventHandler").updateData(HeathReading.ERROR, "Batch consume %s messages failure.", buffer.size()); + } finally { + buffer.clear(); + } } + } finally { + event.clearData(); } } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanDataHolder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanDataHolder.java new file mode 100644 index 000000000..a4ea38da8 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanDataHolder.java @@ -0,0 +1,25 @@ +package com.a.eye.skywalking.storage.disruptor.request; + +import com.a.eye.skywalking.network.grpc.RequestSpan; +import com.a.eye.skywalking.storage.data.spandata.RequestSpanData; + +/** + * @author zhangxin + */ +public class RequestSpanDataHolder { + + private RequestSpanData requestSpanData; + + public void clearData() { + this.requestSpanData = null; + } + + public RequestSpanData getRequestSpanData() { + return requestSpanData; + } + + + public void fillData(RequestSpan requestSpan) { + requestSpanData = new RequestSpanData(requestSpan); + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanFactory.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanFactory.java index 312ccf613..170007da7 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanFactory.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanFactory.java @@ -6,9 +6,9 @@ import com.lmax.disruptor.EventFactory; /** * Created by wusheng on 2016/11/24. */ -public class RequestSpanFactory implements EventFactory { +public class RequestSpanFactory implements EventFactory { @Override - public RequestSpanData newInstance() { - return new RequestSpanData(); + public RequestSpanDataHolder newInstance() { + return new RequestSpanDataHolder(); } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/StoreRequestSpanEventHandler.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/StoreRequestSpanEventHandler.java index d472ca701..3e57d0841 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/StoreRequestSpanEventHandler.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/StoreRequestSpanEventHandler.java @@ -19,7 +19,7 @@ import java.util.List; /** * Created by wusheng on 2016/11/24. */ -public class StoreRequestSpanEventHandler implements EventHandler { +public class StoreRequestSpanEventHandler implements EventHandler { private static ILog logger = LogManager.getLogger(StoreRequestSpanEventHandler.class); private DataFileWriter fileWriter; private IndexOperator operator; @@ -34,23 +34,26 @@ public class StoreRequestSpanEventHandler implements EventHandler requestSpanDisruptor; - private RingBuffer requestSpanRingBuffer; + private Disruptor requestSpanDisruptor; + private RingBuffer requestSpanRingBuffer; - private Disruptor ackSpanDisruptor; - private RingBuffer ackSpanRingBuffer; + private Disruptor ackSpanDisruptor; + private RingBuffer ackSpanRingBuffer; public StorageListener() { - requestSpanDisruptor = new Disruptor(new RequestSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); + requestSpanDisruptor = new Disruptor(new RequestSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); requestSpanDisruptor.handleEventsWith(new StoreRequestSpanEventHandler()); requestSpanDisruptor.start(); requestSpanRingBuffer = requestSpanDisruptor.getRingBuffer(); - ackSpanDisruptor = new Disruptor(new AckSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); + ackSpanDisruptor = new Disruptor(new AckSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); ackSpanDisruptor.handleEventsWith(new StoreAckSpanEventHandler()); ackSpanDisruptor.start(); ackSpanRingBuffer = ackSpanDisruptor.getRingBuffer(); @@ -44,8 +44,8 @@ public class StorageListener implements SpanStorageServerListener { public boolean storage(RequestSpan requestSpan) { long sequence = requestSpanRingBuffer.next(); // Grab the next sequence try { - RequestSpanData data = requestSpanRingBuffer.get(sequence); - data.setRequestSpan(requestSpan); + RequestSpanDataHolder data = requestSpanRingBuffer.get(sequence); + data.fillData(requestSpan); HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.INFO, "RequestSpan stored."); return true; @@ -62,8 +62,8 @@ public class StorageListener implements SpanStorageServerListener { public boolean storage(AckSpan ackSpan) { long sequence = ackSpanRingBuffer.next(); // Grab the next sequence try { - AckSpanData data = ackSpanRingBuffer.get(sequence); - data.setAckSpan(ackSpan); + AckSpanDataHolder data = ackSpanRingBuffer.get(sequence); + data.fillData(ackSpan); HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.INFO, "AckSpan stored."); return true;