Unary alternative to `TraceSegmentReportService.collect` (#5389)

This commit is contained in:
吴晟 Wu Sheng 2020-08-26 19:10:50 +08:00 committed by GitHub
parent 720c1dd92f
commit 086d730aed
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
2 changed files with 24 additions and 2 deletions

@ -1 +1 @@
Subproject commit 89381e14f9c29c4ee240c5e55c66ab44381924ca
Subproject commit cdd58617e720949f51c0ddf5adf10b2b188e94fe

View File

@ -21,6 +21,7 @@ package org.apache.skywalking.oap.server.receiver.trace.provider.handler.v8.grpc
import io.grpc.stub.StreamObserver;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.apm.network.common.v3.Commands;
import org.apache.skywalking.apm.network.language.agent.v3.SegmentCollection;
import org.apache.skywalking.apm.network.language.agent.v3.SegmentObject;
import org.apache.skywalking.apm.network.language.agent.v3.TraceSegmentReportServiceGrpc;
import org.apache.skywalking.oap.server.analyzer.module.AnalyzerModule;
@ -65,7 +66,7 @@ public class TraceSegmentReportServiceHandler extends TraceSegmentReportServiceG
@Override
public void onNext(SegmentObject segment) {
if (log.isDebugEnabled()) {
log.debug("receive segment");
log.debug("received segment in streaming");
}
HistogramMetrics.Timer timer = histogram.createTimer();
@ -91,4 +92,25 @@ public class TraceSegmentReportServiceHandler extends TraceSegmentReportServiceG
}
};
}
@Override
public void collectInSync(final SegmentCollection request, final StreamObserver<Commands> responseObserver) {
if (log.isDebugEnabled()) {
log.debug("received {} segments", request.getSegmentsCount());
}
request.getSegmentsList().forEach(segment -> {
HistogramMetrics.Timer timer = histogram.createTimer();
try {
segmentParserService.send(segment);
} catch (Exception e) {
errorCounter.inc();
} finally {
timer.finish();
}
});
responseObserver.onNext(Commands.newBuilder().build());
responseObserver.onCompleted();
}
}