Adjust log of SegmentServiceClient, close #396
This commit is contained in:
parent
a46c01e34d
commit
815a47769c
|
|
@ -21,7 +21,7 @@ public class GRPCStreamServiceStatus {
|
||||||
/**
|
/**
|
||||||
* @param maxTimeout max wait time, milliseconds.
|
* @param maxTimeout max wait time, milliseconds.
|
||||||
*/
|
*/
|
||||||
public void wait4Finish(long maxTimeout) {
|
public boolean wait4Finish(long maxTimeout) {
|
||||||
long time = 0;
|
long time = 0;
|
||||||
while (!status) {
|
while (!status) {
|
||||||
if (time > maxTimeout) {
|
if (time > maxTimeout) {
|
||||||
|
|
@ -30,6 +30,7 @@ public class GRPCStreamServiceStatus {
|
||||||
try2Sleep(5);
|
try2Sleep(5);
|
||||||
time += 5;
|
time += 5;
|
||||||
}
|
}
|
||||||
|
return status;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
|
||||||
|
|
@ -28,6 +28,9 @@ public class TraceSegmentServiceClient implements BootService, IConsumer<TraceSe
|
||||||
private static final ILog logger = LogManager.getLogger(TraceSegmentServiceClient.class);
|
private static final ILog logger = LogManager.getLogger(TraceSegmentServiceClient.class);
|
||||||
private static final int TIMEOUT = 30 * 1000;
|
private static final int TIMEOUT = 30 * 1000;
|
||||||
|
|
||||||
|
private long lastLogTime;
|
||||||
|
private long segmentUplinkedCounter;
|
||||||
|
private long segmentAbandonedCounter;
|
||||||
private volatile DataCarrier<TraceSegment> carrier;
|
private volatile DataCarrier<TraceSegment> carrier;
|
||||||
private volatile TraceSegmentServiceGrpc.TraceSegmentServiceStub serviceStub;
|
private volatile TraceSegmentServiceGrpc.TraceSegmentServiceStub serviceStub;
|
||||||
private volatile GRPCChannelStatus status = GRPCChannelStatus.DISCONNECT;
|
private volatile GRPCChannelStatus status = GRPCChannelStatus.DISCONNECT;
|
||||||
|
|
@ -39,6 +42,9 @@ public class TraceSegmentServiceClient implements BootService, IConsumer<TraceSe
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void boot() throws Throwable {
|
public void boot() throws Throwable {
|
||||||
|
lastLogTime = System.currentTimeMillis();
|
||||||
|
segmentUplinkedCounter = 0;
|
||||||
|
segmentAbandonedCounter = 0;
|
||||||
carrier = new DataCarrier<TraceSegment>(CHANNEL_SIZE, BUFFER_SIZE);
|
carrier = new DataCarrier<TraceSegment>(CHANNEL_SIZE, BUFFER_SIZE);
|
||||||
carrier.setBufferStrategy(BufferStrategy.IF_POSSIBLE);
|
carrier.setBufferStrategy(BufferStrategy.IF_POSSIBLE);
|
||||||
carrier.consume(this, 1);
|
carrier.consume(this, 1);
|
||||||
|
|
@ -94,14 +100,27 @@ public class TraceSegmentServiceClient implements BootService, IConsumer<TraceSe
|
||||||
}
|
}
|
||||||
upstreamSegmentStreamObserver.onCompleted();
|
upstreamSegmentStreamObserver.onCompleted();
|
||||||
|
|
||||||
status.wait4Finish(TIMEOUT);
|
if (status.wait4Finish(TIMEOUT)) {
|
||||||
|
segmentUplinkedCounter += data.size();
|
||||||
if (logger.isDebugEnable()) {
|
|
||||||
logger.debug("{} trace segments have been sent to collector.", data.size());
|
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
if (logger.isDebugEnable()) {
|
segmentAbandonedCounter += data.size();
|
||||||
logger.debug("{} trace segments have been abandoned, cause by no available channel.", data.size());
|
}
|
||||||
|
|
||||||
|
printUplinkStatus();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void printUplinkStatus() {
|
||||||
|
long currentTimeMillis = System.currentTimeMillis();
|
||||||
|
if (lastLogTime - currentTimeMillis > 30 * 1000) {
|
||||||
|
lastLogTime = currentTimeMillis;
|
||||||
|
if (segmentUplinkedCounter > 0) {
|
||||||
|
logger.debug("{} trace segments have been sent to collector.", segmentUplinkedCounter);
|
||||||
|
segmentUplinkedCounter = 0;
|
||||||
|
}
|
||||||
|
if (segmentAbandonedCounter > 0) {
|
||||||
|
logger.debug("{} trace segments have been abandoned, cause by no available channel.", segmentAbandonedCounter);
|
||||||
|
segmentAbandonedCounter = 0;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue