diff --git a/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/ack/SendAckSpanEventHandler.java b/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/ack/SendAckSpanEventHandler.java index fe9d4d023..317c5d108 100644 --- a/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/ack/SendAckSpanEventHandler.java +++ b/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/ack/SendAckSpanEventHandler.java @@ -26,19 +26,23 @@ public class SendAckSpanEventHandler implements EventHandler { @Override public void onEvent(AckSpanHolder event, long sequence, boolean endOfBatch) throws Exception { - if (buffer[bufferIdx] != null) { - return; - } + try { + if (buffer[bufferIdx] != null) { + return; + } - buffer[bufferIdx] = event.getData(); - bufferIdx++; + buffer[bufferIdx] = event.getData(); + bufferIdx++; - if (bufferIdx == buffer.length) { - bufferIdx = 0; - } + if (bufferIdx == buffer.length) { + bufferIdx = 0; + } - if (endOfBatch) { - HealthCollector.getCurrentHeathReading("SendAckSpanEventHandler").updateData(HeathReading.INFO, "AckSpan messages were successful consumed ."); + if (endOfBatch) { + HealthCollector.getCurrentHeathReading("SendAckSpanEventHandler").updateData(HeathReading.INFO, "AckSpan messages were successful consumed ."); + } + }finally { + event.setData(null); } } diff --git a/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/request/SendRequestSpanEventHandler.java b/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/request/SendRequestSpanEventHandler.java index f35f7fc5a..3899accf6 100644 --- a/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/request/SendRequestSpanEventHandler.java +++ b/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/disruptor/request/SendRequestSpanEventHandler.java @@ -26,19 +26,23 @@ public class SendRequestSpanEventHandler implements EventHandler { private static ILog logger = LogManager.getLogger(StoreAckSpanEventHandler.class); private DataFileWriter fileWriter; - private IndexOperator operator; - private int bufferSize; + private IndexOperator operator; + private int bufferSize; private List buffer; public StoreAckSpanEventHandler() { @@ -35,20 +35,24 @@ public class StoreAckSpanEventHandler implements EventHandler { @Override public void onEvent(AckSpanData event, long sequence, boolean endOfBatch) throws Exception { - buffer.add(event); + try { + buffer.add(event); - 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.setAckSpan(null); } } } 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 be4d69c2b..39c1e37fe 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 @@ -22,8 +22,8 @@ import java.util.List; public class StoreRequestSpanEventHandler implements EventHandler { private static ILog logger = LogManager.getLogger(StoreRequestSpanEventHandler.class); private DataFileWriter fileWriter; - private IndexOperator operator; - private int bufferSize; + private IndexOperator operator; + private int bufferSize; private List buffer; public StoreRequestSpanEventHandler() { @@ -35,21 +35,25 @@ public class StoreRequestSpanEventHandler implements EventHandler