From c99679a8680e3e0d47ad1a3b9cc2de65dd2a4470 Mon Sep 17 00:00:00 2001 From: wusheng Date: Sun, 27 Nov 2016 11:29:36 +0800 Subject: [PATCH] 1. change grpc facade, to open more interface to client. --- .../com/a/eye/skywalking/network/Client.java | 20 +++++------- .../com/a/eye/skywalking/network/Server.java | 14 ++++----- .../ConsumeSpanDataFailedException.java | 9 ------ .../grpc/client/SpanStorageClient.java | 30 +++++++----------- .../grpc/client/TraceSearchClient.java | 6 ++-- .../grpc/server/AsyncTraceSearchServer.java | 6 ++-- .../grpc/server/SpanStorageServer.java | 6 ++-- .../grpc/server/TraceSearchServer.java | 2 +- .../client/StorageClientListener.java | 12 +++++++ .../AsyncTraceSearchServerListener.java} | 4 +-- .../SpanStorageServerListener.java} | 4 +-- .../{ => server}/TraceSearchListener.java | 2 +- .../storage/listener/SearchListener.java | 5 ++- .../storage/listener/StorageListener.java | 4 +-- .../src/test/java/StorageThread.java | 31 ++++++++++++++++++- 15 files changed, 87 insertions(+), 68 deletions(-) delete mode 100644 skywalking-network/src/main/java/com/a/eye/skywalking/network/exception/ConsumeSpanDataFailedException.java create mode 100644 skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/client/StorageClientListener.java rename skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/{AsyncTraceSearchListener.java => server/AsyncTraceSearchServerListener.java} (66%) rename skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/{SpanStorageListener.java => server/SpanStorageServerListener.java} (66%) rename skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/{ => server}/TraceSearchListener.java (73%) diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/Client.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/Client.java index 29e1d696d..6f9991f28 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/Client.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/Client.java @@ -1,30 +1,24 @@ package com.a.eye.skywalking.network; -import com.a.eye.skywalking.network.grpc.SpanStorageServiceGrpc; -import com.a.eye.skywalking.network.grpc.TraceSearchServiceGrpc; import com.a.eye.skywalking.network.grpc.client.SpanStorageClient; import com.a.eye.skywalking.network.grpc.client.TraceSearchClient; +import com.a.eye.skywalking.network.listener.client.StorageClientListener; import io.grpc.ManagedChannel; import io.grpc.ManagedChannelBuilder; public class Client { private ManagedChannel channel; - private SpanStorageServiceGrpc.SpanStorageServiceStub spanStorageStub; - private TraceSearchServiceGrpc.TraceSearchServiceStub traceSearchServiceStub; public Client(String ip, int address) { channel = ManagedChannelBuilder.forAddress(ip, address).usePlaintext(true).build(); - spanStorageStub = SpanStorageServiceGrpc.newStub(channel); - traceSearchServiceStub = TraceSearchServiceGrpc.newStub(channel); + } + + public SpanStorageClient newSpanStorageClient(StorageClientListener listener) { + return new SpanStorageClient(channel, listener); } - public SpanStorageClient newSpanStorageClient() { - return new SpanStorageClient(spanStorageStub); - } - - - public TraceSearchClient newTraceSearchClient(){ - return new TraceSearchClient(traceSearchServiceStub); + public TraceSearchClient newTraceSearchClient(StorageClientListener listener){ + return new TraceSearchClient(channel, listener); } } diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/Server.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/Server.java index 6770e18af..b558d77bc 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/Server.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/Server.java @@ -3,9 +3,9 @@ package com.a.eye.skywalking.network; import com.a.eye.skywalking.network.grpc.server.AsyncTraceSearchServer; import com.a.eye.skywalking.network.grpc.server.SpanStorageServer; import com.a.eye.skywalking.network.grpc.server.TraceSearchServer; -import com.a.eye.skywalking.network.listener.AsyncTraceSearchListener; -import com.a.eye.skywalking.network.listener.SpanStorageListener; -import com.a.eye.skywalking.network.listener.TraceSearchListener; +import com.a.eye.skywalking.network.listener.server.AsyncTraceSearchServerListener; +import com.a.eye.skywalking.network.listener.server.SpanStorageServerListener; +import com.a.eye.skywalking.network.listener.server.TraceSearchListener; import io.grpc.netty.NettyServerBuilder; import io.netty.channel.nio.NioEventLoopGroup; @@ -51,8 +51,8 @@ public class Server { .workerEventLoopGroup(new NioEventLoopGroup()).build()); } - public TransferServiceBuilder addSpanStorageService(SpanStorageListener spanStorageListener) { - serverBuilder.addService(new SpanStorageServer(spanStorageListener)); + public TransferServiceBuilder addSpanStorageService(SpanStorageServerListener spanStorageServerListener) { + serverBuilder.addService(new SpanStorageServer(spanStorageServerListener)); return this; } @@ -61,8 +61,8 @@ public class Server { return this; } - public TransferServiceBuilder addAsyncTraceSearchService(AsyncTraceSearchListener asyncTraceSearchListener){ - serverBuilder.addService(new AsyncTraceSearchServer(asyncTraceSearchListener)); + public TransferServiceBuilder addAsyncTraceSearchService(AsyncTraceSearchServerListener asyncTraceSearchServerListener){ + serverBuilder.addService(new AsyncTraceSearchServer(asyncTraceSearchServerListener)); return this; } } diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/exception/ConsumeSpanDataFailedException.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/exception/ConsumeSpanDataFailedException.java deleted file mode 100644 index cd50960bb..000000000 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/exception/ConsumeSpanDataFailedException.java +++ /dev/null @@ -1,9 +0,0 @@ -package com.a.eye.skywalking.network.exception; - -/** - * Created by xin on 2016/11/23. - */ -public class ConsumeSpanDataFailedException extends RuntimeException { - public ConsumeSpanDataFailedException(Exception e) { - } -} diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/SpanStorageClient.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/SpanStorageClient.java index 5a6aa8eed..14e261448 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/SpanStorageClient.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/SpanStorageClient.java @@ -1,10 +1,11 @@ package com.a.eye.skywalking.network.grpc.client; -import com.a.eye.skywalking.network.exception.ConsumeSpanDataFailedException; import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.network.grpc.RequestSpan; import com.a.eye.skywalking.network.grpc.SendResult; import com.a.eye.skywalking.network.grpc.SpanStorageServiceGrpc; +import com.a.eye.skywalking.network.listener.client.StorageClientListener; +import io.grpc.ManagedChannel; import io.grpc.stub.CallStreamObserver; import io.grpc.stub.StreamObserver; @@ -12,8 +13,11 @@ public class SpanStorageClient { private final SpanStorageServiceGrpc.SpanStorageServiceStub spanStorageStub; - public SpanStorageClient(SpanStorageServiceGrpc.SpanStorageServiceStub spanStorageStub) { - this.spanStorageStub = spanStorageStub; + private final StorageClientListener listener; + + public SpanStorageClient(ManagedChannel channel, StorageClientListener listener) { + this.spanStorageStub = SpanStorageServiceGrpc.newStub(channel); + this.listener = listener; } public void sendRequestSpan(RequestSpan... requestSpan) { @@ -21,11 +25,12 @@ public class SpanStorageClient { spanStorageStub.storageRequestSpan(new StreamObserver() { @Override public void onNext(SendResult sendResult) { + listener.onBatchFinished(sendResult); } @Override public void onError(Throwable throwable) { - throwable.printStackTrace(); + listener.onError(throwable); } @Override @@ -37,13 +42,6 @@ public class SpanStorageClient { for (RequestSpan span : requestSpan) { requestSpanStreamObserver.onNext(span); } - while (!((CallStreamObserver) requestSpanStreamObserver).isReady()) { - try { - Thread.sleep(1); - } catch (InterruptedException e) { - throw new ConsumeSpanDataFailedException(e); - } - } requestSpanStreamObserver.onCompleted(); } @@ -53,11 +51,12 @@ public class SpanStorageClient { spanStorageStub.storageACKSpan(new StreamObserver() { @Override public void onNext(SendResult sendResult) { + listener.onBatchFinished(sendResult); } @Override public void onError(Throwable throwable) { - throwable.printStackTrace(); + listener.onError(throwable); } @Override @@ -69,13 +68,6 @@ public class SpanStorageClient { for (AckSpan span : ackSpan) { ackSpanStreamObserver.onNext(span); } - while (!((CallStreamObserver) ackSpanStreamObserver).isReady()) { - try { - Thread.sleep(1); - } catch (InterruptedException e) { - throw new ConsumeSpanDataFailedException(e); - } - } ackSpanStreamObserver.onCompleted(); } diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/TraceSearchClient.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/TraceSearchClient.java index ec4df2f05..9f4883437 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/TraceSearchClient.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/client/TraceSearchClient.java @@ -1,6 +1,8 @@ package com.a.eye.skywalking.network.grpc.client; import com.a.eye.skywalking.network.grpc.TraceSearchServiceGrpc; +import com.a.eye.skywalking.network.listener.client.StorageClientListener; +import io.grpc.ManagedChannel; /** * Created by wusheng on 2016/11/26. @@ -9,8 +11,8 @@ public class TraceSearchClient { private final TraceSearchServiceGrpc.TraceSearchServiceStub traceSearchServiceStub; - public TraceSearchClient(TraceSearchServiceGrpc.TraceSearchServiceStub traceSearchServiceStub) { - this.traceSearchServiceStub = traceSearchServiceStub; + public TraceSearchClient(ManagedChannel channel, StorageClientListener listener) { + this.traceSearchServiceStub = TraceSearchServiceGrpc.newStub(channel); } } diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/AsyncTraceSearchServer.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/AsyncTraceSearchServer.java index 3dd75a326..bfeaad84c 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/AsyncTraceSearchServer.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/AsyncTraceSearchServer.java @@ -4,7 +4,7 @@ import com.a.eye.skywalking.network.grpc.AsyncTraceSearchServiceGrpc; import com.a.eye.skywalking.network.grpc.QueryTask; import com.a.eye.skywalking.network.grpc.SearchResult; import com.a.eye.skywalking.network.grpc.Span; -import com.a.eye.skywalking.network.listener.AsyncTraceSearchListener; +import com.a.eye.skywalking.network.listener.server.AsyncTraceSearchServerListener; import io.grpc.stub.StreamObserver; import java.util.List; @@ -14,9 +14,9 @@ import java.util.List; */ public class AsyncTraceSearchServer extends AsyncTraceSearchServiceGrpc.AsyncTraceSearchServiceImplBase { - private AsyncTraceSearchListener searchListener; + private AsyncTraceSearchServerListener searchListener; - public AsyncTraceSearchServer(AsyncTraceSearchListener searchListener) { + public AsyncTraceSearchServer(AsyncTraceSearchServerListener searchListener) { this.searchListener = searchListener; } diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/SpanStorageServer.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/SpanStorageServer.java index 1db52ac51..cca7710c3 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/SpanStorageServer.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/SpanStorageServer.java @@ -4,13 +4,13 @@ import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.network.grpc.RequestSpan; import com.a.eye.skywalking.network.grpc.SendResult; import com.a.eye.skywalking.network.grpc.SpanStorageServiceGrpc; -import com.a.eye.skywalking.network.listener.SpanStorageListener; +import com.a.eye.skywalking.network.listener.server.SpanStorageServerListener; import io.grpc.stub.StreamObserver; public class SpanStorageServer extends SpanStorageServiceGrpc.SpanStorageServiceImplBase { - private SpanStorageListener listener; - public SpanStorageServer(SpanStorageListener listener) { + private SpanStorageServerListener listener; + public SpanStorageServer(SpanStorageServerListener listener) { this.listener = listener; } diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/TraceSearchServer.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/TraceSearchServer.java index b81c6d1c3..45235c076 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/TraceSearchServer.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/grpc/server/TraceSearchServer.java @@ -3,7 +3,7 @@ package com.a.eye.skywalking.network.grpc.server; import com.a.eye.skywalking.network.grpc.SearchResult; import com.a.eye.skywalking.network.grpc.TraceId; import com.a.eye.skywalking.network.grpc.TraceSearchServiceGrpc; -import com.a.eye.skywalking.network.listener.TraceSearchListener; +import com.a.eye.skywalking.network.listener.server.TraceSearchListener; import io.grpc.stub.StreamObserver; /** diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/client/StorageClientListener.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/client/StorageClientListener.java new file mode 100644 index 000000000..75bca8d64 --- /dev/null +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/client/StorageClientListener.java @@ -0,0 +1,12 @@ +package com.a.eye.skywalking.network.listener.client; + +import com.a.eye.skywalking.network.grpc.SendResult; + +/** + * Created by wusheng on 2016/11/27. + */ +public interface StorageClientListener { + void onError(Throwable throwable); + + void onBatchFinished(SendResult sendResult); +} diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/AsyncTraceSearchListener.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/AsyncTraceSearchServerListener.java similarity index 66% rename from skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/AsyncTraceSearchListener.java rename to skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/AsyncTraceSearchServerListener.java index d6fb79a5d..7959c0f2f 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/AsyncTraceSearchListener.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/AsyncTraceSearchServerListener.java @@ -1,4 +1,4 @@ -package com.a.eye.skywalking.network.listener; +package com.a.eye.skywalking.network.listener.server; import com.a.eye.skywalking.network.grpc.Span; import com.a.eye.skywalking.network.grpc.TraceId; @@ -8,6 +8,6 @@ import java.util.List; /** * Created by xin on 2016/11/15. */ -public interface AsyncTraceSearchListener { +public interface AsyncTraceSearchServerListener { List search(TraceId traceId); } diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/SpanStorageListener.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/SpanStorageServerListener.java similarity index 66% rename from skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/SpanStorageListener.java rename to skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/SpanStorageServerListener.java index 71c26794b..80884bee9 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/SpanStorageListener.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/SpanStorageServerListener.java @@ -1,9 +1,9 @@ -package com.a.eye.skywalking.network.listener; +package com.a.eye.skywalking.network.listener.server; import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.network.grpc.RequestSpan; -public interface SpanStorageListener{ +public interface SpanStorageServerListener { boolean storage(RequestSpan requestSpan); boolean storage(AckSpan ackSpan); diff --git a/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/TraceSearchListener.java b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/TraceSearchListener.java similarity index 73% rename from skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/TraceSearchListener.java rename to skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/TraceSearchListener.java index f19a6734a..bd5f247de 100644 --- a/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/TraceSearchListener.java +++ b/skywalking-network/src/main/java/com/a/eye/skywalking/network/listener/server/TraceSearchListener.java @@ -1,4 +1,4 @@ -package com.a.eye.skywalking.network.listener; +package com.a.eye.skywalking.network.listener.server; import com.a.eye.skywalking.network.grpc.Span; diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java index 0bd8e8d27..695bb2ea9 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/SearchListener.java @@ -6,8 +6,7 @@ import com.a.eye.skywalking.logging.api.ILog; import com.a.eye.skywalking.logging.api.LogManager; import com.a.eye.skywalking.network.grpc.Span; import com.a.eye.skywalking.network.grpc.TraceId; -import com.a.eye.skywalking.network.listener.AsyncTraceSearchListener; -import com.a.eye.skywalking.network.listener.TraceSearchListener; +import com.a.eye.skywalking.network.listener.server.AsyncTraceSearchServerListener; import com.a.eye.skywalking.storage.data.SpanDataFinder; import com.a.eye.skywalking.storage.data.spandata.SpanData; import com.a.eye.skywalking.storage.data.spandata.SpanDataHelper; @@ -15,7 +14,7 @@ import com.a.eye.skywalking.storage.data.spandata.SpanDataHelper; import java.util.ArrayList; import java.util.List; -public class SearchListener implements AsyncTraceSearchListener { +public class SearchListener implements AsyncTraceSearchServerListener { private static ILog logger = LogManager.getLogger(SearchListener.class); 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 eceede1ea..6fcd81817 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 @@ -6,7 +6,7 @@ import com.a.eye.skywalking.logging.api.ILog; import com.a.eye.skywalking.logging.api.LogManager; 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.network.listener.server.SpanStorageServerListener; import com.a.eye.skywalking.storage.config.Config; import com.a.eye.skywalking.storage.data.spandata.AckSpanData; import com.a.eye.skywalking.storage.data.spandata.RequestSpanData; @@ -18,7 +18,7 @@ import com.lmax.disruptor.RingBuffer; import com.lmax.disruptor.dsl.Disruptor; import com.lmax.disruptor.util.DaemonThreadFactory; -public class StorageListener implements SpanStorageListener { +public class StorageListener implements SpanStorageServerListener { private ILog logger = LogManager.getLogger(StorageListener.class); diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/StorageThread.java b/skywalking-storage-center/skywalking-storage/src/test/java/StorageThread.java index a7ab1489d..45121b940 100644 --- a/skywalking-storage-center/skywalking-storage/src/test/java/StorageThread.java +++ b/skywalking-storage-center/skywalking-storage/src/test/java/StorageThread.java @@ -1,22 +1,28 @@ import com.a.eye.skywalking.network.Client; import com.a.eye.skywalking.network.grpc.AckSpan; import com.a.eye.skywalking.network.grpc.RequestSpan; +import com.a.eye.skywalking.network.grpc.SendResult; import com.a.eye.skywalking.network.grpc.TraceId; import com.a.eye.skywalking.network.grpc.client.SpanStorageClient; +import com.a.eye.skywalking.network.listener.client.StorageClientListener; import com.a.eye.skywalking.storage.util.NetUtils; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.locks.LockSupport; public class StorageThread extends Thread { private SpanStorageClient client; private long count; private CountDownLatch countDownLatch; + private MyStorageClientListener listener; + StorageThread(long count, CountDownLatch countDownLatch) { - client = new Client("10.128.7.241", 34000).newSpanStorageClient(); + listener = new MyStorageClientListener(); + client = new Client("10.128.7.241", 34000).newSpanStorageClient(listener); this.count = count; this.countDownLatch = countDownLatch; } @@ -42,6 +48,11 @@ public class StorageThread extends Thread { client.sendACKSpan(ackSpanList); client.sendRequestSpan(requestSpanList); cycle = 0; + + while(!listener.isCompleted){ + LockSupport.parkNanos(1); + } + listener.begin(); } requestSpanList[cycle] = requestSpan; @@ -55,4 +66,22 @@ public class StorageThread extends Thread { countDownLatch.countDown(); } + + public class MyStorageClientListener implements StorageClientListener{ + volatile boolean isCompleted = false; + + @Override + public void onError(Throwable throwable) { + + } + + @Override + public void onBatchFinished(SendResult sendResult) { + isCompleted = true; + } + + public void begin(){ + isCompleted = false; + } + } }