Pinpoint 분석 #7

2026. 8. 4. 00:21·개발 메모

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.java
  • agent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/context/provider/StorageFactoryProvider.java
  • agent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/context/storage/BufferedStorageFactory.java
  • agent-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.java
  • agent-module/profiler/src/main/java/com/navercorp/pinpoint/profiler/sender/grpc/AgentGrpcDataSender.java
  • collector/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
'개발 메모' 카테고리의 다른 글
  • Pinpoint 분석 #6
  • Pinpoint 분석 #5
  • Pinpoint 분석 #4
  • Pinpoint 분석 #3
csb0710
csb0710
  • csb0710
    데모장
    csb0710
  • 전체
    오늘
    어제
    • 분류 전체보기 (60)
      • 스프링부트 메모 (7)
      • 개발 메모 (10)
      • 클라우드 메모 (10)
      • 설치&설정 메모 (2)
      • 알고리즘 메모 (18)
      • 인턴 메모 (7)
      • 데이터베이스 메모 (3)
      • 책 메모 (1)
  • 블로그 메뉴

    • 홈
    • 태그
    • 방명록
  • 링크

  • 공지사항

  • 인기 글

  • 태그

    코딩테스트
    스프링부트
    그리디
    디비설정
    서버 연결
    디비설치
    .gitmodules
    알고리즘
    GitHub
    HBase
    오블완
    코드트리
    이지퍼블리싱
    코드트리조별과제
    자동 답변 봇
    서버배포
    서버생성
    백준
    티스토리챌린지
    submodule
  • 최근 댓글

  • 최근 글

  • hELLO· Designed By정상우.v4.10.2
csb0710
Pinpoint 분석 #7
상단으로

티스토리툴바