From efe90aa97017d7aaaa4aa664bc928ef8eab00231 Mon Sep 17 00:00:00 2001 From: ascrutae Date: Tue, 7 Feb 2017 22:03:46 +0800 Subject: [PATCH] fix some question --- .../ack/RouteAckSpanBufferEventHandler.java | 48 +++++++++---------- .../RouteSendRequestSpanEventHandler.java | 48 +++++++++---------- .../storage/data/file/DataFileWriter.java | 13 ++++- .../ack/StoreAckSpanEventHandler.java | 26 +++++----- .../request/StoreRequestSpanEventHandler.java | 27 +++++------ 5 files changed, 79 insertions(+), 83 deletions(-) diff --git a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/RouteAckSpanBufferEventHandler.java b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/RouteAckSpanBufferEventHandler.java index d8c7fc796..f2c532915 100644 --- a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/RouteAckSpanBufferEventHandler.java +++ b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/ack/RouteAckSpanBufferEventHandler.java @@ -29,38 +29,34 @@ public class RouteAckSpanBufferEventHandler extends AbstractRouteSpanEventHandle @Override public void onEvent(AckSpanHolder event, long sequence, boolean endOfBatch) throws Exception { - try { - buffer.add(event.getAckSpan()); + buffer.add(event.getAckSpan()); - if (stop) { - try { - for (AckSpan ackSpan : buffer) { - SpanDisruptor spanDisruptor = RoutingService.getRouter().lookup(ackSpan); - spanDisruptor.saveSpan(ackSpan); - } - } finally { - buffer.clear(); + if (stop) { + try { + for (AckSpan ackSpan : buffer) { + SpanDisruptor spanDisruptor = RoutingService.getRouter().lookup(ackSpan); + spanDisruptor.saveSpan(ackSpan); } - - return; + } finally { + buffer.clear(); } - wait2Finish(); + return; + } - if (endOfBatch || buffer.size() == bufferSize) { - try { - SpanStorageClient spanStorageClient = getStorageClient(); - spanStorageClient.sendACKSpan(buffer); - HealthCollector.getCurrentHeathReading("RouteAckSpanBufferEventHandler").updateData(HeathReading.INFO, "Batch consume %s messages successfully.", buffer.size()); - } catch (Throwable e) { - logger.error("Ack messages consume failure.", e); - HealthCollector.getCurrentHeathReading("RouteAckSpanBufferEventHandler").updateData(HeathReading.ERROR, "Batch consume %s messages failure.", buffer.size()); - } finally { - buffer.clear(); - } + wait2Finish(); + + if (endOfBatch || buffer.size() == bufferSize) { + try { + SpanStorageClient spanStorageClient = getStorageClient(); + spanStorageClient.sendACKSpan(buffer); + HealthCollector.getCurrentHeathReading("RouteAckSpanBufferEventHandler").updateData(HeathReading.INFO, "Batch consume %s messages successfully.", buffer.size()); + } catch (Throwable e) { + logger.error("Ack messages consume failure.", e); + HealthCollector.getCurrentHeathReading("RouteAckSpanBufferEventHandler").updateData(HeathReading.ERROR, "Batch consume %s messages failure.", buffer.size()); + } finally { + buffer.clear(); } - } finally { - event.setAckSpan(null); } } } diff --git a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/request/RouteSendRequestSpanEventHandler.java b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/request/RouteSendRequestSpanEventHandler.java index 7ed6cc388..6c15fb6c6 100644 --- a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/request/RouteSendRequestSpanEventHandler.java +++ b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/request/RouteSendRequestSpanEventHandler.java @@ -32,38 +32,34 @@ public class RouteSendRequestSpanEventHandler extends AbstractRouteSpanEventHand @Override public void onEvent(RequestSpanHolder event, long sequence, boolean endOfBatch) throws Exception { - try { - buffer.add(event.getRequestSpan()); + buffer.add(event.getRequestSpan()); - if (stop) { - try { - for (RequestSpan requestSpan : buffer) { - SpanDisruptor spanDisruptor = RoutingService.getRouter().lookup(requestSpan); - spanDisruptor.saveSpan(requestSpan); - } - } finally { - buffer.clear(); + if (stop) { + try { + for (RequestSpan requestSpan : buffer) { + SpanDisruptor spanDisruptor = RoutingService.getRouter().lookup(requestSpan); + spanDisruptor.saveSpan(requestSpan); } - - return; + } finally { + buffer.clear(); } - wait2Finish(); + return; + } - if (endOfBatch || buffer.size() == bufferSize) { - try { - SpanStorageClient spanStorageClient = getStorageClient(); - spanStorageClient.sendRequestSpan(buffer); - HealthCollector.getCurrentHeathReading("RouteSendRequestSpanEventHandler").updateData(HeathReading.INFO, "Batch consume %s messages successfully.", buffer.size()); - } catch (Throwable e) { - logger.error("RequestSpan messages consume failure.", e); - HealthCollector.getCurrentHeathReading("RouteSendRequestSpanEventHandler").updateData(HeathReading.ERROR, "Batch consume %s messages failure.", buffer.size()); - } finally { - buffer.clear(); - } + wait2Finish(); + + if (endOfBatch || buffer.size() == bufferSize) { + try { + SpanStorageClient spanStorageClient = getStorageClient(); + spanStorageClient.sendRequestSpan(buffer); + HealthCollector.getCurrentHeathReading("RouteSendRequestSpanEventHandler").updateData(HeathReading.INFO, "Batch consume %s messages successfully.", buffer.size()); + } catch (Throwable e) { + logger.error("RequestSpan messages consume failure.", e); + HealthCollector.getCurrentHeathReading("RouteSendRequestSpanEventHandler").updateData(HeathReading.ERROR, "Batch consume %s messages failure.", buffer.size()); + } finally { + buffer.clear(); } - } finally { - event.setRequestSpan(null); } } } 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 1d4a72ad6..259b54ef4 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 @@ -1,5 +1,7 @@ package com.a.eye.skywalking.storage.data.file; +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; import com.a.eye.skywalking.storage.data.spandata.SpanData; import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; @@ -20,14 +22,23 @@ public class DataFileWriter { } IndexMetaCollection collections = new IndexMetaCollection(); + int failedCount = 0; try { for (SpanData data : spanData) { - collections.add(dataFile.write(data)); + try { + collections.add(dataFile.write(data)); + }catch (Throwable e){ + failedCount++; + } } }finally { dataFile.flush(); } + if (failedCount > 0) { + HealthCollector.getCurrentHeathReading("DataFileWriter").updateData(HeathReading.ERROR ,"Failed to write %s span to data file.", Integer.valueOf(failedCount)); + } + return collections; } 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 2edc3cb76..fa1916539 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,24 +35,20 @@ public class StoreAckSpanEventHandler implements EventHandler { @Override public void onEvent(AckSpanData event, long sequence, boolean endOfBatch) throws Exception { - try { - buffer.add(event); + 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 39c1e37fe..d472ca701 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,25 +35,22 @@ public class StoreRequestSpanEventHandler implements EventHandler