diff --git a/apm-protocol/apm-network/src/main/proto b/apm-protocol/apm-network/src/main/proto index 89381e14f..cdd58617e 160000 --- a/apm-protocol/apm-network/src/main/proto +++ b/apm-protocol/apm-network/src/main/proto @@ -1 +1 @@ -Subproject commit 89381e14f9c29c4ee240c5e55c66ab44381924ca +Subproject commit cdd58617e720949f51c0ddf5adf10b2b188e94fe diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/TraceSegmentReportServiceHandler.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/TraceSegmentReportServiceHandler.java index b9bbc4ef6..405347533 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/TraceSegmentReportServiceHandler.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/TraceSegmentReportServiceHandler.java @@ -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 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(); + } }