5. Collector 전송 로직
전체 흐름
DefaultBaseTraceFactory.newDefaultTrace(...)
-> StorageFactory.createStorage(...)
-> BufferedStorage
-> DataSender<SpanType>
-> SpanGrpcDataSender.send(...)
-> MessageConverter<SpanType, GeneratedMessageV3>
-> PSpan or PSpanChunk
-> PSpanMessage
-> SpanGrpc.sendSpan()
-> collector SpanService.sendSpan()
-> SimpleHandler<PSpan/PSpanChunk>
5.1 Span sender 바인딩
파일:
agent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/context/module/GrpcModule.javaagent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/context/provider/StorageFactoryProvider.javaagent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/context/storage/BufferedStorageFactory.javaagent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/context/provider/grpc/SpanGrpcDataSenderProvider.java
GrpcModule.bindSpanDataSender()가 @SpanDataSender DataSender<SpanType>을 SpanGrpcDataSenderProvider에 바인딩한다.
GrpcModule.bindSpanDataSender()
-> bind(MessageConverter<SpanType, GeneratedMessageV3>).toProvider(GrpcSpanMessageConverterProvider)
-> bind(SpanProcessor<PSpan.Builder, PSpanChunk.Builder>).toProvider(GrpcSpanProcessorProvider)
-> bind(DataSender<SpanType>).annotatedWith(SpanDataSender).toProvider(SpanGrpcDataSenderProvider)
StorageFactoryProvider는 이 sender를 받아 BufferedStorageFactory를 만든다.
StorageFactoryProvider.get()
-> contextConfig.isIoBufferingEnable()
-> new BufferedStorageFactory(bufferSize, spanDataSender)
5.2 SpanGrpcDataSender
파일:
agent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/sender/grpc/SpanGrpcDataSender.java
핵심 메소드/필드:
SpanGrpcDataSender(...)dispatcher.onDispatch(...)startStream()attemptRenew()close()
생성자에서 span collector로 향하는 async stub과 client streaming task를 만든다.
new SpanGrpcDataSender(...)
-> SpanGrpc.newStub(managedChannel)
-> reconnectExecutor.newReconnector(...)
-> new ClientStreamingService(...)
-> reconnectJob.run()
-> startStream()
실제 data dispatch는 dispatcher.onDispatch(...)에서 일어난다.
onDispatch(stream, SpanType data)
-> messageConverter.toMessage(data)
-> if PSpanChunk: PSpanMessage.setSpanChunk(...)
-> if PSpan: PSpanMessage.setSpan(...)
-> stream.onNext(spanMessage)
-> attemptRenew()
startStream()은 queue를 읽는 DefaultStreamTask를 시작한다.
startStream()
-> new DefaultStreamTask(id, clientStreamService, streamExecutorFactory, queue, dispatcher, failState)
-> streamTask.start()
여기서 queue는 상위 GrpcDataSender 계층의 전송 queue다. trace thread는 dataSender.send(span)까지만 하고, 실제 network write는 stream task가 비동기로 처리한다.
5.3 Collector SpanService
파일:
collector/src/main/java/com/navercorp/pinpoint/collector/receiver/grpc/service/SpanService.java
핵심 메소드:
sendSpan(StreamObserver<Empty>)messageDispatch(...)
agent의 spanStub.sendSpan(response)에 대응하는 collector endpoint가 SpanService.sendSpan()이다.
SpanService.sendSpan(...)
-> new ServerCallStream<>(..., this::messageDispatch, ...)
stream으로 들어온 PSpanMessage는 span과 span chunk로 갈라져 handler에 전달된다.
messageDispatch(...)
-> if spanMessage.hasSpan()
serverRequestFactory.newServerRequest(MessageTypes.SPAN, span)
spanHandler.handleSimple(request)
-> else if spanMessage.hasSpanChunk()
serverRequestFactory.newServerRequest(MessageTypes.SPANCHUNK, spanChunk)
spanChunkHandler.handleSimple(request)
5.4 AgentInfo 전송
파일:
agent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentInfoSender.javaagent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/sender/grpc/AgentGrpcDataSender.javacollector/src/main/java/com/navercorp/pinpoint/collector/receiver/grpc/service/AgentService.java
DefaultApplicationContext.start()가 agentInfoSender.start()를 호출하면 AgentInfoSender scheduler가 agent info를 보낸다.
AgentInfoSender.start()
-> Scheduler.start()
-> AgentInfoSendTask.run()
-> agentInfoFactory.createAgentInfo()
-> dataSender.request(agentInfo)
gRPC 구현은 AgentGrpcDataSender.request(...)가 담당한다.
AgentGrpcDataSender.request(...)
-> messageConverter.toMessage(agentInfo)
-> agentInfoStub.requestAgentInfo(pAgentInfo, observer)
collector 쪽은 AgentService.requestAgentInfo(...)가 받는다.
AgentService.requestAgentInfo(...)
-> jobRunner.execute(MessageTypes.AGENT_INFO, agentInfo, responseObserver, handler::handleRequest)
-> pingEventHandler.update(transportId)'개발 메모' 카테고리의 다른 글
| Pinpoint 분석 #6 (0) | 2026.08.04 |
|---|---|
| Pinpoint 분석 #5 (0) | 2026.08.04 |
| Pinpoint 분석 #4 (0) | 2026.08.04 |
| Pinpoint 분석 #3 (0) | 2026.08.04 |
| Pinpoint 분석 #2 (0) | 2026.08.04 |