From 9561a50f6d06b63b501000858a3b9876cbb7b7fc Mon Sep 17 00:00:00 2001 From: wusheng Date: Thu, 24 Nov 2016 17:32:51 +0800 Subject: [PATCH 1/5] use disruptor to replace DataCarrier, to improve performance --- skywalking-storage-center/pom.xml | 6 +-- .../eye/skywalking/storage/config/Config.java | 4 ++ .../storage/data/SpanDataConsumer.java | 49 ------------------ .../storage/data/file/DataFileNameDesc.java | 3 +- .../storage/data/file/DataFileWriter.java | 5 +- .../storage/data/spandata/AckSpanData.java | 4 ++ .../data/spandata/RequestSpanData.java | 4 ++ .../storage/disruptor/ack/AckSpanFactory.java | 14 +++++ .../ack/StoreAckSpanEventHandler.java | 45 ++++++++++++++++ .../disruptor/request/RequestSpanFactory.java | 14 +++++ .../request/StoreRequestSpanEventHandler.java | 45 ++++++++++++++++ .../storage/listener/StorageListener.java | 47 +++++++++++------ .../storage/util/AtomicRangeInteger.java | 51 +++++++++++++++++++ 13 files changed, 222 insertions(+), 69 deletions(-) delete mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanFactory.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/StoreAckSpanEventHandler.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanFactory.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/StoreRequestSpanEventHandler.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/util/AtomicRangeInteger.java diff --git a/skywalking-storage-center/pom.xml b/skywalking-storage-center/pom.xml index 7f2cff9c3..1115373a7 100644 --- a/skywalking-storage-center/pom.xml +++ b/skywalking-storage-center/pom.xml @@ -25,9 +25,9 @@ ${project.version} - com.a.eye - data-carrier - 1.2 + com.lmax + disruptor + 3.3.6 com.a.eye 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 ebdab7afb..755cc4ede 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,6 +8,10 @@ public class Config { public static int PORT = 34000; } + public static class Disruptor{ + public static int BUFFER_SIZE = 1024 * 128; + } + public static class DataConsumer { diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java deleted file mode 100644 index d5efd6cb7..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java +++ /dev/null @@ -1,49 +0,0 @@ -package com.a.eye.skywalking.storage.data; - -import com.a.eye.datacarrier.consumer.IConsumer; -import com.a.eye.skywalking.health.report.HealthCollector; -import com.a.eye.skywalking.health.report.HeathReading; -import com.a.eye.skywalking.logging.api.ILog; -import com.a.eye.skywalking.logging.api.LogManager; -import com.a.eye.skywalking.storage.data.file.DataFileWriter; -import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; -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.SpanData; - -import java.util.List; - -public class SpanDataConsumer implements IConsumer { - - private static ILog logger = LogManager.getLogger(SpanDataConsumer.class); - private DataFileWriter fileWriter; - private IndexOperator operator; - - @Override - public void init() { - fileWriter = new DataFileWriter(); - operator = IndexOperatorFactory.createIndexOperator(); - } - - @Override - public void consume(List data) { - IndexMetaCollection collection = fileWriter.write(data); - - operator.batchUpdate(collection); - - HealthCollector.getCurrentHeathReading("SpanDataConsumer") - .updateData(HeathReading.INFO, "%s messages were successful consumed .", data.size()); - } - - @Override - public void onError(List span, Throwable throwable) { - logger.error("Failed to consumer span data.", throwable); - HealthCollector.getCurrentHeathReading("SpanDataConsumer").updateData(HeathReading.ERROR, - "Failed to consume span data. error message : " + throwable.getMessage()); - } - - @Override - public void onExit() { - fileWriter.close(); - } -} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java index 89fb28be4..8626b97d5 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileNameDesc.java @@ -1,6 +1,7 @@ package com.a.eye.skywalking.storage.data.file; -import com.a.eye.datacarrier.common.AtomicRangeInteger; + +import com.a.eye.skywalking.storage.util.AtomicRangeInteger; import java.text.ParseException; import java.text.SimpleDateFormat; diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java index d818cf04a..54fbda980 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java @@ -23,11 +23,14 @@ public class DataFileWriter { for (SpanData data : spanData) { collections.add(dataFile.write(data)); } - dataFile.flush(); return collections; } + public void flush(){ + dataFile.flush(); + } + public void close(){ dataFile.close(); } 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 16530939e..ad8d81533 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 @@ -16,6 +16,10 @@ public class AckSpanData extends AbstractSpanData { public AckSpanData() { } + public void setAckSpan(AckSpan ackSpan) { + this.ackSpan = ackSpan; + } + @Override public SpanType getSpanType() { return SpanType.ACKSpan; diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java index 52733af6e..2646a9246 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/spandata/RequestSpanData.java @@ -16,6 +16,10 @@ public class RequestSpanData extends AbstractSpanData { public RequestSpanData() { } + public void setRequestSpan(RequestSpan requestSpan) { + this.requestSpan = requestSpan; + } + @Override public SpanType getSpanType() { return SpanType.RequestSpan; 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 new file mode 100644 index 000000000..5bcebd928 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/AckSpanFactory.java @@ -0,0 +1,14 @@ +package com.a.eye.skywalking.storage.disruptor.ack; + +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; +import com.lmax.disruptor.EventFactory; + +/** + * Created by wusheng on 2016/11/24. + */ +public class AckSpanFactory implements EventFactory { + @Override + public AckSpanData newInstance() { + return new AckSpanData(); + } +} 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 new file mode 100644 index 000000000..3f150d730 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/ack/StoreAckSpanEventHandler.java @@ -0,0 +1,45 @@ +package com.a.eye.skywalking.storage.disruptor.ack; + +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.storage.data.file.DataFileWriter; +import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; +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.lmax.disruptor.EventHandler; + +import java.util.ArrayList; +import java.util.List; + +/** + * Created by wusheng on 2016/11/24. + */ +public class StoreAckSpanEventHandler implements EventHandler { + private static ILog logger = LogManager.getLogger(StoreAckSpanEventHandler.class); + private DataFileWriter fileWriter; + private IndexOperator operator; + private int bufferSize = 100; + private List buffer = new ArrayList<>(bufferSize); + + public StoreAckSpanEventHandler() { + fileWriter = new DataFileWriter(); + operator = IndexOperatorFactory.createIndexOperator(); + } + + @Override + public void onEvent(AckSpanData event, long sequence, boolean endOfBatch) throws Exception { + buffer.add(event); + + if (endOfBatch || buffer.size() == bufferSize) { + IndexMetaCollection collection = fileWriter.write(buffer); + + operator.batchUpdate(collection); + + HealthCollector.getCurrentHeathReading("StoreAckSpanEventHandler").updateData(HeathReading.INFO, "%s messages were successful consumed .", buffer.size()); + } + } +} 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 new file mode 100644 index 000000000..312ccf613 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/RequestSpanFactory.java @@ -0,0 +1,14 @@ +package com.a.eye.skywalking.storage.disruptor.request; + +import com.a.eye.skywalking.storage.data.spandata.RequestSpanData; +import com.lmax.disruptor.EventFactory; + +/** + * Created by wusheng on 2016/11/24. + */ +public class RequestSpanFactory implements EventFactory { + @Override + public RequestSpanData newInstance() { + return new RequestSpanData(); + } +} 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 new file mode 100644 index 000000000..741f64f7a --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/disruptor/request/StoreRequestSpanEventHandler.java @@ -0,0 +1,45 @@ +package com.a.eye.skywalking.storage.disruptor.request; + +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; +import com.a.eye.skywalking.logging.api.ILog; +import com.a.eye.skywalking.logging.api.LogManager; +import com.a.eye.skywalking.storage.data.file.DataFileWriter; +import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; +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.RequestSpanData; +import com.a.eye.skywalking.storage.data.spandata.SpanData; +import com.lmax.disruptor.EventHandler; + +import java.util.ArrayList; +import java.util.List; + +/** + * Created by wusheng on 2016/11/24. + */ +public class StoreRequestSpanEventHandler implements EventHandler { + private static ILog logger = LogManager.getLogger(StoreRequestSpanEventHandler.class); + private DataFileWriter fileWriter; + private IndexOperator operator; + private int bufferSize = 100; + private List buffer = new ArrayList<>(bufferSize); + + public StoreRequestSpanEventHandler() { + fileWriter = new DataFileWriter(); + operator = IndexOperatorFactory.createIndexOperator(); + } + + @Override + public void onEvent(RequestSpanData event, long sequence, boolean endOfBatch) throws Exception { + buffer.add(event); + + if (endOfBatch || buffer.size() == bufferSize) { + IndexMetaCollection collection = fileWriter.write(buffer); + + operator.batchUpdate(collection); + + HealthCollector.getCurrentHeathReading("StoreRequestSpanEventHandler").updateData(HeathReading.INFO, "%s messages were successful consumed .", buffer.size()); + } + } +} 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 673860e8c..1af869527 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 @@ -1,6 +1,5 @@ package com.a.eye.skywalking.storage.listener; -import com.a.eye.datacarrier.DataCarrier; import com.a.eye.skywalking.health.report.HealthCollector; import com.a.eye.skywalking.health.report.HeathReading; import com.a.eye.skywalking.logging.api.ILog; @@ -9,32 +8,48 @@ import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.network.grpc.RequestSpan; import com.a.eye.skywalking.network.listener.SpanStorageListener; import com.a.eye.skywalking.storage.config.Config; -import com.a.eye.skywalking.storage.data.SpanDataConsumer; -import com.a.eye.skywalking.storage.data.spandata.SpanData; -import com.a.eye.skywalking.storage.data.spandata.SpanDataBuilder; +import com.a.eye.skywalking.storage.data.spandata.AckSpanData; +import com.a.eye.skywalking.storage.data.spandata.RequestSpanData; +import com.a.eye.skywalking.storage.disruptor.ack.AckSpanFactory; +import com.a.eye.skywalking.storage.disruptor.ack.StoreAckSpanEventHandler; +import com.a.eye.skywalking.storage.disruptor.request.RequestSpanFactory; +import com.a.eye.skywalking.storage.disruptor.request.StoreRequestSpanEventHandler; +import com.lmax.disruptor.RingBuffer; +import com.lmax.disruptor.dsl.Disruptor; +import com.lmax.disruptor.util.DaemonThreadFactory; public class StorageListener implements SpanStorageListener { private ILog logger = LogManager.getLogger(StorageListener.class); - private DataCarrier spanDataDataCarrier; + private Disruptor requestSpanDisruptor; + private RingBuffer requestSpanRingBuffer; + + private Disruptor ackSpanDisruptor; + private RingBuffer ackSpanRingBuffer; public StorageListener() { - spanDataDataCarrier = new DataCarrier<>(Config.DataConsumer.CHANNEL_SIZE, Config.DataConsumer.BUFFER_SIZE); - spanDataDataCarrier.consume(SpanDataConsumer.class, Config.DataConsumer.CONSUMER_SIZE); + requestSpanDisruptor = new Disruptor(new RequestSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); + requestSpanDisruptor.handleEventsWith(new StoreRequestSpanEventHandler()); + requestSpanRingBuffer = requestSpanDisruptor.getRingBuffer(); + + ackSpanDisruptor = new Disruptor(new AckSpanFactory(), Config.Disruptor.BUFFER_SIZE, DaemonThreadFactory.INSTANCE); + ackSpanDisruptor.handleEventsWith(new StoreAckSpanEventHandler()); + ackSpanRingBuffer = ackSpanDisruptor.getRingBuffer(); } @Override public boolean storage(RequestSpan requestSpan) { try { - spanDataDataCarrier.produce(SpanDataBuilder.build(requestSpan)); - HealthCollector.getCurrentHeathReading("StorageListener") - .updateData(HeathReading.INFO, "RequestSpan stored."); + long sequence = requestSpanRingBuffer.next(); // Grab the next sequence + RequestSpanData data = requestSpanRingBuffer.get(sequence); + data.setRequestSpan(requestSpan); + + HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.INFO, "RequestSpan stored."); return true; } catch (Exception e) { logger.error("RequestSpan trace-id[{}] store failure..", requestSpan.getTraceId(), e); - HealthCollector.getCurrentHeathReading("StorageListener") - .updateData(HeathReading.ERROR, "RequestSpan store failure."); + HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.ERROR, "RequestSpan store failure."); return false; } } @@ -42,13 +57,15 @@ public class StorageListener implements SpanStorageListener { @Override public boolean storage(AckSpan ackSpan) { try { - spanDataDataCarrier.produce(SpanDataBuilder.build(ackSpan)); + long sequence = ackSpanRingBuffer.next(); // Grab the next sequence + AckSpanData data = ackSpanRingBuffer.get(sequence); + data.setAckSpan(ackSpan); + HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.INFO, "AckSpan stored."); return true; } catch (Exception e) { logger.error("AckSpan trace-id[{}] store failure..", ackSpan.getTraceId(), e); - HealthCollector.getCurrentHeathReading("StorageListener") - .updateData(HeathReading.ERROR, "AckSpan store failure."); + HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.ERROR, "AckSpan store failure."); return false; } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/util/AtomicRangeInteger.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/util/AtomicRangeInteger.java new file mode 100644 index 000000000..66b34f6d6 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/util/AtomicRangeInteger.java @@ -0,0 +1,51 @@ +package com.a.eye.skywalking.storage.util; + +import java.io.Serializable; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Created by wusheng on 2016/10/25. + */ +public class AtomicRangeInteger extends Number implements Serializable { + private static final long serialVersionUID = -4099792402691141643L; + private AtomicInteger value; + private int startValue; + private int endValue; + + public AtomicRangeInteger(int startValue, int maxValue) { + this.value = new AtomicInteger(startValue); + this.startValue = startValue; + this.endValue = maxValue - 1; + } + + public final int getAndIncrement() { + int current; + int next; + do { + current = this.value.get(); + next = current >= this.endValue?this.startValue:current + 1; + } while(!this.value.compareAndSet(current, next)); + + return current; + } + + public final int get() { + return this.value.get(); + } + + public int intValue() { + return this.value.intValue(); + } + + public long longValue() { + return this.value.longValue(); + } + + public float floatValue() { + return this.value.floatValue(); + } + + public double doubleValue() { + return this.value.doubleValue(); + } +} From 340dd221c30f558cdbb31750f12b0a5a8072d7f7 Mon Sep 17 00:00:00 2001 From: wusheng Date: Thu, 24 Nov 2016 17:50:27 +0800 Subject: [PATCH 2/5] fix flush missing --- .../a/eye/skywalking/storage/data/file/DataFileWriter.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java index 54fbda980..87e3ca525 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java @@ -24,11 +24,9 @@ public class DataFileWriter { collections.add(dataFile.write(data)); } - return collections; - } - - public void flush(){ dataFile.flush(); + + return collections; } public void close(){ From d1431cad59e3b566ef4574f03feb066409348eb7 Mon Sep 17 00:00:00 2001 From: wusheng Date: Thu, 24 Nov 2016 17:52:33 +0800 Subject: [PATCH 3/5] fix publish missing --- .../eye/skywalking/storage/listener/StorageListener.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) 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 1af869527..15b997693 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 @@ -40,8 +40,8 @@ public class StorageListener implements SpanStorageListener { @Override public boolean storage(RequestSpan requestSpan) { + long sequence = requestSpanRingBuffer.next(); // Grab the next sequence try { - long sequence = requestSpanRingBuffer.next(); // Grab the next sequence RequestSpanData data = requestSpanRingBuffer.get(sequence); data.setRequestSpan(requestSpan); @@ -51,13 +51,15 @@ public class StorageListener implements SpanStorageListener { logger.error("RequestSpan trace-id[{}] store failure..", requestSpan.getTraceId(), e); HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.ERROR, "RequestSpan store failure."); return false; + } finally{ + requestSpanRingBuffer.publish(sequence); } } @Override public boolean storage(AckSpan ackSpan) { + long sequence = ackSpanRingBuffer.next(); // Grab the next sequence try { - long sequence = ackSpanRingBuffer.next(); // Grab the next sequence AckSpanData data = ackSpanRingBuffer.get(sequence); data.setAckSpan(ackSpan); @@ -67,6 +69,8 @@ public class StorageListener implements SpanStorageListener { logger.error("AckSpan trace-id[{}] store failure..", ackSpan.getTraceId(), e); HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.ERROR, "AckSpan store failure."); return false; + } finally{ + requestSpanRingBuffer.publish(sequence); } } } From e750b540ccd6a0583b39fd91b5874aedae55f191 Mon Sep 17 00:00:00 2001 From: wusheng Date: Thu, 24 Nov 2016 18:29:07 +0800 Subject: [PATCH 4/5] fix: 1.disruptor not start. 2.init data file from an empty dir failure. --- .../eye/skywalking/storage/data/file/DataFileLoader.java | 3 ++- .../eye/skywalking/storage/listener/StorageListener.java | 8 +++++--- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileLoader.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileLoader.java index 1fd7971bd..feac53fac 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileLoader.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileLoader.java @@ -21,7 +21,8 @@ public class DataFileLoader { List allDataFile = new ArrayList(); for (File fileEntry : dataFileDir.listFiles()) { - allDataFile.add(new DataFile(fileEntry)); + if (fileEntry.getName().split("_").length == 8) + allDataFile.add(new DataFile(fileEntry)); } return allDataFile; } 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 15b997693..eceede1ea 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 @@ -31,10 +31,12 @@ public class StorageListener implements SpanStorageListener { public StorageListener() { 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.handleEventsWith(new StoreAckSpanEventHandler()); + ackSpanDisruptor.start(); ackSpanRingBuffer = ackSpanDisruptor.getRingBuffer(); } @@ -51,7 +53,7 @@ public class StorageListener implements SpanStorageListener { logger.error("RequestSpan trace-id[{}] store failure..", requestSpan.getTraceId(), e); HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.ERROR, "RequestSpan store failure."); return false; - } finally{ + } finally { requestSpanRingBuffer.publish(sequence); } } @@ -69,8 +71,8 @@ public class StorageListener implements SpanStorageListener { logger.error("AckSpan trace-id[{}] store failure..", ackSpan.getTraceId(), e); HealthCollector.getCurrentHeathReading("StorageListener").updateData(HeathReading.ERROR, "AckSpan store failure."); return false; - } finally{ - requestSpanRingBuffer.publish(sequence); + } finally { + ackSpanRingBuffer.publish(sequence); } } } From b5e00fe20291f3b8ad6c2770134f1d45352991ce Mon Sep 17 00:00:00 2001 From: wusheng Date: Fri, 25 Nov 2016 09:28:25 +0800 Subject: [PATCH 5/5] fix: event handler memory leak, buffer never be cleared --- .../disruptor/ack/StoreAckSpanEventHandler.java | 10 +++++++--- .../request/StoreRequestSpanEventHandler.java | 10 +++++++--- 2 files changed, 14 insertions(+), 6 deletions(-) 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 3f150d730..4f361b404 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 @@ -35,11 +35,15 @@ public class StoreAckSpanEventHandler implements EventHandler { buffer.add(event); if (endOfBatch || buffer.size() == bufferSize) { - IndexMetaCollection collection = fileWriter.write(buffer); + try { + IndexMetaCollection collection = fileWriter.write(buffer); - operator.batchUpdate(collection); + operator.batchUpdate(collection); - HealthCollector.getCurrentHeathReading("StoreAckSpanEventHandler").updateData(HeathReading.INFO, "%s messages were successful consumed .", buffer.size()); + HealthCollector.getCurrentHeathReading("StoreAckSpanEventHandler").updateData(HeathReading.INFO, "%s messages were successful consumed .", buffer.size()); + } finally { + buffer.clear(); + } } } } 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 741f64f7a..3b7371c04 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 @@ -35,11 +35,15 @@ public class StoreRequestSpanEventHandler implements EventHandler