Using waiting grpc send status, instead of reponse. Avoid performance loss.

This commit is contained in:
wusheng 2016-11-27 17:24:08 +08:00
parent 8a7a6f04fc
commit 77f096d46f
4 changed files with 38 additions and 32 deletions

View File

@ -10,6 +10,7 @@ import io.grpc.stub.CallStreamObserver;
import io.grpc.stub.StreamObserver;
import java.util.List;
import java.util.concurrent.locks.LockSupport;
public class SpanStorageClient {
@ -23,54 +24,59 @@ public class SpanStorageClient {
}
public void sendRequestSpan(List<RequestSpan> requestSpan) {
StreamObserver<RequestSpan> requestSpanStreamObserver =
spanStorageStub.storageRequestSpan(new StreamObserver<SendResult>() {
@Override
public void onNext(SendResult sendResult) {
listener.onBatchFinished(sendResult);
}
StreamObserver<RequestSpan> requestSpanStreamObserver = spanStorageStub.storageRequestSpan(new StreamObserver<SendResult>() {
@Override
public void onNext(SendResult sendResult) {
}
@Override
public void onError(Throwable throwable) {
listener.onError(throwable);
}
@Override
public void onError(Throwable throwable) {
listener.onError(throwable);
}
@Override
public void onCompleted() {
}
});
@Override
public void onCompleted() {
listener.onBatchFinished();
}
});
for (RequestSpan span : requestSpan) {
requestSpanStreamObserver.onNext(span);
}
while (!((CallStreamObserver<RequestSpan>) requestSpanStreamObserver).isReady()) {
LockSupport.parkNanos(1);
}
requestSpanStreamObserver.onCompleted();
}
public void sendACKSpan(List<AckSpan> ackSpan) {
StreamObserver<AckSpan> ackSpanStreamObserver =
spanStorageStub.storageACKSpan(new StreamObserver<SendResult>() {
@Override
public void onNext(SendResult sendResult) {
listener.onBatchFinished(sendResult);
}
StreamObserver<AckSpan> ackSpanStreamObserver = spanStorageStub.storageACKSpan(new StreamObserver<SendResult>() {
@Override
public void onNext(SendResult sendResult) {
@Override
public void onError(Throwable throwable) {
listener.onError(throwable);
}
}
@Override
public void onCompleted() {
@Override
public void onError(Throwable throwable) {
listener.onError(throwable);
}
}
});
@Override
public void onCompleted() {
listener.onBatchFinished();
}
});
for (AckSpan span : ackSpan) {
ackSpanStreamObserver.onNext(span);
}
while (!((CallStreamObserver<AckSpan>) ackSpanStreamObserver).isReady()) {
LockSupport.parkNanos(1);
}
ackSpanStreamObserver.onCompleted();
}

View File

@ -8,5 +8,5 @@ import com.a.eye.skywalking.network.grpc.SendResult;
public interface StorageClientListener {
void onError(Throwable throwable);
void onBatchFinished(SendResult sendResult);
void onBatchFinished();
}

View File

@ -134,7 +134,7 @@ public class Agent2RoutingClient extends Thread {
}
@Override
public void onBatchFinished(SendResult sendResult) {
public void onBatchFinished() {
batchFinished = true;
HealthCollector.getCurrentHeathReading("Agent2RoutingClient").updateData(HeathReading.INFO, "batch send data to routing node.");
}

View File

@ -76,7 +76,7 @@ public class StorageThread extends Thread {
}
@Override
public void onBatchFinished(SendResult sendResult) {
public void onBatchFinished() {
isCompleted = true;
}