diff --git a/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/client/Agent2RoutingClient.java b/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/client/Agent2RoutingClient.java index 093119a7c..53704dad2 100644 --- a/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/client/Agent2RoutingClient.java +++ b/skywalking-sniffer/skywalking-api/src/main/java/com/a/eye/skywalking/client/Agent2RoutingClient.java @@ -26,7 +26,7 @@ public class Agent2RoutingClient extends Thread { private static ILog logger = LogManager.getLogger(Agent2RoutingClient.class); private List addrList; - private Client client; + private Client client; private SpanStorageClient spanStorageClient; private NetworkListener listener; private SendRequestSpanEventHandler requestSpanDataSupplier = null; @@ -43,7 +43,7 @@ public class Agent2RoutingClient extends Thread { if (addrSegments.length != 2) { throw new IllegalArgumentException("server addr should like ip:port, illegal addr:" + server); } - addrList.add(new ServerAddr(addrSegments[0], addrSegments[2])); + addrList.add(new ServerAddr(addrSegments[0], addrSegments[1])); } listener = new NetworkListener(); } @@ -63,7 +63,7 @@ public class Agent2RoutingClient extends Thread { private void connect() { try { - if(client != null && !client.isShutdown()){ + if (client != null && !client.isShutdown()) { client.shutdown(); } int addrIdx = new Random().nextInt(addrList.size()); @@ -84,26 +84,33 @@ public class Agent2RoutingClient extends Thread { List requestData = this.requestSpanDataSupplier.getBufferData(); List ackData = this.ackSpanDataSupplier.getBufferData(); - if (requestData.size() > 0 || ackData.size() > 0) { + boolean hasData = false; + if (requestData.size() > 0) { + hasData = true; listener.begin(); spanStorageClient.sendRequestSpan(requestData); + + + listener.wait2Confirm(); + } + if (ackData.size() > 0) { + hasData = true; + listener.begin(); + spanStorageClient.sendACKSpan(ackData); - while (!listener.isBatchFinished()) { - try { - Thread.sleep(10L); - } catch (InterruptedException e) { + listener.wait2Confirm(); + } - } - } - } else { + if(!hasData) { try { Thread.sleep(10 * 1000L); } catch (InterruptedException e) { } } + } try { @@ -143,6 +150,21 @@ public class Agent2RoutingClient extends Thread { HealthCollector.getCurrentHeathReading("Agent2RoutingClient").updateData(HeathReading.INFO, "batch send data to routing node."); } + void wait2Confirm(){ + // wait 20s, most. + int countDown = 100 * 20; + while (!listener.isBatchFinished()) { + try { + Thread.sleep(10L); + if(countDown-- < 0){ + batchFinished = true; + } + } catch (InterruptedException e) { + + } + } + } + } diff --git a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/AbstractRouteSpanEventHandler.java b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/AbstractRouteSpanEventHandler.java index 7eb776c12..9bbbfb569 100644 --- a/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/AbstractRouteSpanEventHandler.java +++ b/skywalking-storage-center/skywalking-routing/src/main/java/com/a/eye/skywalking/routing/disruptor/AbstractRouteSpanEventHandler.java @@ -56,9 +56,14 @@ public abstract class AbstractRouteSpanEventHandler implements EventHandler