Java 클라이언트 분산 추적

Java 클라이언트 분산 추적 (Java client tracing)

Pulsar Java 클라이언트는 OpenTelemetry 분산 추적을 내장 지원해요. 프로듀서가 메시지를 보낼 때 스팬(span)이 생기고 컨슈머가 처리할 때까지 이어지는 end-to-end 트레이스를 만들 수 있어요. 이 문서는 추적을 활성화하는 방법과 스팬의 속성, ack 동작에 따라 스팬이 어떻게 끝나는지를 정리해줘요.

출처: 문서

본문

이 문서는 Pulsar Java 클라이언트와 함께 OpenTelemetry 분산 추적을 사용하는 방법을 설명해요.

개요 (Overview)

Pulsar Java 클라이언트는 OpenTelemetry 분산 추적을 내장 지원해요. 이를 통해 다음과 같은 것들을 할 수 있어요.

  • 프로듀서에서 브로커로의 메시지 게시 추적
  • 브로커에서 컨슈머로의 메시지 소비 추적
  • 메시지 속성을 통한 서비스 간 트레이스 컨텍스트 전파
  • 외부 소스(예: HTTP 요청)에서 트레이스 컨텍스트 추출
  • 분산 시스템 전반에 걸친 end-to-end 트레이스 생성

기능 (Features)

프로듀서 추적 (Producer Tracing)

프로듀서 추적은 다음에 대한 스팬을 만들어요.

  • sendsend() 또는 sendAsync()가 호출될 때 스팬이 시작되고, 브로커가 수신을 ack할 때 완료돼요.

컨슈머 추적 (Consumer Tracing)

컨슈머 추적은 다음에 대한 스팬을 만들어요.

  • process — 메시지가 수신될 때 스팬이 시작되고, 메시지가 ack되거나 부정 ack되거나 ack 타임아웃이 발생할 때 완료돼요.

트레이스 컨텍스트 전파 (Trace Context Propagation)

트레이스 컨텍스트는 W3C TraceContext 형식으로 자동 전파돼요.

  • traceparent — 트레이스 ID, 스팬 ID, 트레이스 플래그를 포함해요.
  • tracestate — 벤더별 트레이스 정보를 포함해요.

컨텍스트는 메시지 속성에 주입되고 추출되어, 서비스 간의 매끄러운 트레이스 전파를 가능하게 해요.

빠른 시작 (Quick Start)

1. 의존성 추가 (Add Dependencies)

Pulsar 클라이언트는 이미 OpenTelemetry API 의존성을 포함해요. SDK와 exporter를 추가하면 돼요.

<dependency>
    <groupId>io.opentelemetry</groupId>
    <artifactId>opentelemetry-sdk</artifactId>
    <version>${opentelemetry.version}</version>
</dependency>

<dependency>
    <groupId>io.opentelemetry</groupId>
    <artifactId>opentelemetry-exporter-otlp</artifactId>
    <version>${opentelemetry.version}</version>
</dependency>

2. 추적 활성화 (Enable Tracing)

추적을 활성화하는 방법은 세 가지가 있어요.

옵션 1: OpenTelemetry Java Agent 사용 (가장 쉬움)

# Start your application with the Java Agent
java -javaagent:opentelemetry-javaagent.jar \
     -Dotel.service.name=my-service \
     -Dotel.exporter.otlp.endpoint=http://localhost:4317 \
     -jar your-application.jar
// Just enable tracing - uses GlobalOpenTelemetry from the agent
PulsarClient client = PulsarClient.builder()
    .serviceUrl("pulsar://localhost:6650")
    .enableTracing(true)  // That's it!
    .build();

옵션 2: 명시적 OpenTelemetry 인스턴스 사용

OpenTelemetry openTelemetry = // configure your OpenTelemetry instance
PulsarClient client = PulsarClient.builder()
    .serviceUrl("pulsar://localhost:6650")
    .openTelemetry(openTelemetry, true)  // Set OpenTelemetry AND enable tracing
    .build();

옵션 3: GlobalOpenTelemetry 사용

// Configure GlobalOpenTelemetry once in your application
GlobalOpenTelemetry.set(myOpenTelemetry);

// Enable tracing in the client - uses GlobalOpenTelemetry
PulsarClient client = PulsarClient.builder()
    .serviceUrl("pulsar://localhost:6650")
    .enableTracing(true)
    .build();

추적이 활성화되면 일어나는 일:

  • 프로듀서 send 연산에 대한 스팬 생성
  • 메시지 속성에 트레이스 컨텍스트 자동 주입
  • 컨슈머 receive/ack 연산에 대한 스팬 생성
  • 메시지 속성에서 트레이스 컨텍스트 자동 추출
  • 모든 스팬을 연결해 end-to-end 분산 트레이스 생성

3. 수동 인터셉터 구성 (고급) (Manual Interceptor Configuration)

수동으로 제어하려면 인터셉터를 명시적으로 추가할 수 있어요.

import org.apache.pulsar.client.impl.tracing.OpenTelemetryProducerInterceptor;
import org.apache.pulsar.client.impl.tracing.OpenTelemetryConsumerInterceptor;

// Create client (tracing not enabled globally)
PulsarClient client = PulsarClient.builder()
    .serviceUrl("pulsar://localhost:6650")
    .openTelemetry(openTelemetry)
    .build();

// Add interceptor manually to specific producer
Producer<String> producer = client.newProducer(Schema.STRING)
    .topic("my-topic")
    .intercept(new OpenTelemetryProducerInterceptor())
    .create();

// Add interceptor manually to specific consumer
Consumer<String> consumer = client.newConsumer(Schema.STRING)
    .topic("my-topic")
    .subscriptionName("my-subscription")
    .intercept(new OpenTelemetryConsumerInterceptor<>())
    .subscribe();

고급 사용법 (Advanced Usage)

End-to-End 추적 예시 (End-to-End Tracing Example)

이 예시는 HTTP 요청에서 Pulsar를 거쳐 컨슈머까지 이어지는 완전한 트레이스를 만드는 방법을 보여줘요.

// Service 1: HTTP API that publishes to Pulsar
@POST
@Path("/order")
public Response createOrder(@Context HttpHeaders headers, Order order) {
    // Extract trace context from incoming HTTP request
    Context context = TracingProducerBuilder.extractFromHeaders(
        convertHeaders(headers));

    // Publish to Pulsar with trace context
    TracingProducerBuilder tracingBuilder = new TracingProducerBuilder();
    producer.newMessage()
        .value(order)
        .let(builder -> tracingBuilder.injectContext(builder, context))
        .send();

    return Response.accepted().build();
}

// Service 2: Pulsar consumer that processes orders
Consumer<Order> consumer = client.newConsumer(Schema.JSON(Order.class))
    .topic("orders")
    .subscriptionName("order-processor")
    .intercept(new OpenTelemetryConsumerInterceptor<>())
    .subscribe();

while (true) {
    Message<Order> msg = consumer.receive();
    // Trace context is automatically extracted
    // Any spans created here will be part of the same trace
    processOrder(msg.getValue());
    consumer.acknowledge(msg);
}

사용자 정의 스팬 생성 (Custom Span Creation)

메시지 처리 중 사용자 정의 스팬을 만들 수 있어요.

import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.Tracer;
import io.opentelemetry.context.Scope;

Tracer tracer = GlobalOpenTelemetry.get().getTracer("my-app");

Message<String> msg = consumer.receive();

// Create a custom span for processing
Span span = tracer.spanBuilder("process-message")
    .setSpanKind(SpanKind.INTERNAL)
    .startSpan();

try (Scope scope = span.makeCurrent()) {
    // Your processing logic
    processMessage(msg.getValue());
    span.setStatus(StatusCode.OK);
} catch (Exception e) {
    span.recordException(e);
    span.setStatus(StatusCode.ERROR);
    throw e;
} finally {
    span.end();
    consumer.acknowledge(msg);
}

구성 (Configuration)

OpenTelemetry Java Agent와의 호환성 (Compatibility with OpenTelemetry Java Agent)

이 구현은 Pulsar용 OpenTelemetry Java Instrumentation과 완전히 호환돼요.

  • 둘 다 W3C TraceContext 형식(traceparent, tracestate 헤더)을 사용해요.
  • 둘 다 메시지 속성을 통해 컨텍스트를 전파해요.
  • 충돌 없음(No conflicts) — 우리 구현은 트레이스 컨텍스트가 이미 존재하는지(Java Agent에서) 확인하고 중복 주입을 피해요.
  • 두 접근 중 하나 또는 둘 다 함께 사용할 수 있어요.

OpenTelemetry Java Agent 사용 (Using OpenTelemetry Java Agent)

추적을 활성화하는 가장 쉬운 방법은 OpenTelemetry Java Agent(자동 계측)를 사용하는 거예요.

java -javaagent:path/to/opentelemetry-javaagent.jar \
     -Dotel.service.name=my-service \
     -Dotel.exporter.otlp.endpoint=http://localhost:4317 \
     -jar your-application.jar

Note: Java Agent를 사용할 때는 .openTelemetry(otel, true)를 호출할 필요가 없어요 — 에이전트가 Pulsar를 자동으로 계측하거든요. 하지만 호출해도 충돌은 발생하지 않아요.

프로그래밍 방식 구성 (Programmatic Configuration)

OpenTelemetry를 프로그래밍 방식으로도 구성할 수 있어요.

import io.opentelemetry.sdk.OpenTelemetrySdk;
import io.opentelemetry.sdk.trace.SdkTracerProvider;
import io.opentelemetry.sdk.trace.export.BatchSpanProcessor;
import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter;

OtlpGrpcSpanExporter spanExporter = OtlpGrpcSpanExporter.builder()
    .setEndpoint("http://localhost:4317")
    .build();

SdkTracerProvider tracerProvider = SdkTracerProvider.builder()
    .addSpanProcessor(BatchSpanProcessor.builder(spanExporter).build())
    .build();

OpenTelemetrySdk openTelemetry = OpenTelemetrySdk.builder()
    .setTracerProvider(tracerProvider)
    .buildAndRegisterGlobal();

환경 변수 (Environment Variables)

환경 변수로 구성해요.

export OTEL_SERVICE_NAME=my-service
export OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317
export OTEL_TRACES_EXPORTER=otlp
export OTEL_METRICS_EXPORTER=otlp

스팬 속성 (Span Attributes)

추적 구현은 OpenTelemetry 메시징 시맨틱 컨벤션에 따라 스팬에 다음 속성을 추가해요.

프로듀서 스팬 (Producer Spans)

  • messaging.system : "pulsar"
  • messaging.destination.name : 토픽 이름
  • messaging.operation.name : "send"
  • messaging.message.id : 메시지 ID (브로커가 확인할 때 추가)

스팬 이름 : send {topic} (예: "send my-topic")

컨슈머 스팬 (Consumer Spans)

  • messaging.system : "pulsar"
  • messaging.destination.name : 토픽 이름
  • messaging.destination.subscription.name : 서브스크립션 이름
  • messaging.operation.name : "process"
  • messaging.message.id : 메시지 ID
  • messaging.pulsar.acknowledgment.type : 메시지가 어떻게 ack되었는지
    • "acknowledge" : 일반 개별 ack
    • "cumulative_acknowledge" : 누적 ack
    • "negative_acknowledge" : 메시지가 부정 ack됨 (재시도 예정)
    • "ack_timeout" : ack 타임아웃 발생 (재시도 예정)

스팬 이름 : process {topic} (예: "process my-topic")

스팬 수명 주기와 ack 동작 (Span Lifecycle and Acknowledgment Behavior)

다양한 ack 시나리오에서 스팬이 어떻게 처리되는지 이해해봐요. 모든 컨슈머 스팬에는 어떻게 완료됐는지를 나타내는 messaging.pulsar.acknowledgment.type 속성이 포함돼요.

성공적인 ack (Successful Acknowledgment)

  • 스팬이 OK 상태로 끝나요.
  • 속성: messaging.pulsar.acknowledgment.type = "acknowledge"

누적 ack (Cumulative Acknowledgment)

  • 스팬이 OK 상태로 끝나요.
  • 속성: messaging.pulsar.acknowledgment.type = "cumulative_acknowledge"
  • ack된 위치까지의 모든 스팬이 이 속성으로 끝나요.

부정 ack (Negative Acknowledgment)

  • 스팬이 OK 상태로 끝나요(오류가 아님).
  • 속성: messaging.pulsar.acknowledgment.type = "negative_acknowledge"
  • 이는 정상 흐름이지 실패가 아니에요 — 메시지가 재전달되고 새 스팬이 만들어질 거예요.

ack 타임아웃 (Acknowledgment Timeout)

  • 스팬이 OK 상태로 끝나요(오류가 아님).
  • 속성: messaging.pulsar.acknowledgment.type = "ack_timeout"
  • 이것은 ackTimeout이 구성됐을 때 기대되는 동작이에요 — 메시지가 재전달되고 새 스팬이 만들어질 거예요.

처리 중 애플리케이션 예외 (Application Exception During Processing)

  • 애플리케이션 코드가 예외를 던지면 자식 스팬을 만들고 ERROR 상태로 표시해요.
  • 컨슈머 스팬 자체는 negativeAcknowledge()를 호출할 때 정상적으로 끝나요.
  • 이렇게 하면 메시징 연산(OK)과 애플리케이션 로직(ERROR)을 명확히 분리할 수 있어요.

메시징 오류와 애플리케이션 오류를 분리하는 예시:

Message<String> msg = consumer.receive();
Span processingSpan = tracer.spanBuilder("business-logic").startSpan();

try (Scope scope = processingSpan.makeCurrent()) {
    processMessage(msg.getValue());
    processingSpan.setStatus(StatusCode.OK);
    consumer.acknowledge(msg);  // Consumer span ends with acknowledgment.type="acknowledge"
} catch (Exception e) {
    processingSpan.recordException(e);
    processingSpan.setStatus(StatusCode.ERROR);  // Business logic failed
    consumer.negativeAcknowledge(msg);  // Consumer span ends with acknowledgment.type="negative_acknowledge"
    throw e;
} finally {
    processingSpan.end();
}

ack 유형별 조회 (Querying by Acknowledgment Type)

messaging.pulsar.acknowledgment.type 속성으로 스팬을 필터링하고 분석할 수 있어요.

추적 백엔드의 예시 쿼리:

  • 모든 재시도된 메시지 찾기: messaging.pulsar.acknowledgment.type = "negative_acknowledge" OR "ack_timeout"
  • 재시도율 계산: count(negative_acknowledge) / count(acknowledge)
  • 타임아웃 문제 식별: messaging.pulsar.acknowledgment.type = "ack_timeout"
  • 누적 vs 개별 ack 분석: messaging.pulsar.acknowledgment.type로 그룹화

모범 사례 (Best Practices)

  • 항상 인터셉터 사용 — 완전한 가시성을 위해 프로듀서와 컨슈머 양쪽에 추적 인터셉터를 추가해요.
  • HTTP에서 컨텍스트 전파 — HTTP 엔드포인트에서 게시할 때 항상 트레이스 컨텍스트를 추출하고 전파해요.
  • 오류를 올바르게 처리 — 예외가 발생해도 스팬이 항상 끝나도록 해요.
  • 메시징 오류와 애플리케이션 오류 구분 — 메시징 연산(nack, 타임아웃)은 OK 상태와 이벤트로 끝나요. 애플리케이션 실패는 ERROR 상태의 별도 자식 스팬에서 추적해야 해요.
  • 의미 있는 스팬 이름 사용 — 기본 스팬 이름은 쉽게 식별할 수 있도록 토픽 이름을 포함해요.
  • 성능 고려 — 추적은 최소한의 오버헤드를 추가하지만, 높은 처리량 시나리오에서는 샘플링을 고려해요.
  • 리소스 정리 — 종료 시 인터셉터와 OpenTelemetry SDK를 제대로 닫아야 해요.

문제 해결 (Troubleshooting)

트레이스가 나타나지 않을 때 (Traces not appearing)

  • OpenTelemetry SDK가 구성되고 exporter가 설정됐는지 확인해요.
  • 프로듀서/컨슈머에 인터셉터가 추가됐는지 확인해요.
  • 트레이스 exporter 엔드포인트에 연결할 수 있는지 확인해요.
  • 디버그 로깅 활성화: -Dio.opentelemetry.javaagent.debug=true

부모-자식 관계가 없을 때 (Missing parent-child relationships)

  • TracingProducerBuilder.injectContext()로 트레이스 컨텍스트가 주입되고 있는지 확인해요.
  • 메시지 속성에 traceparent 헤더가 있는지 확인해요.
  • 프로듀서와 컨슈머 모두에 추적 인터셉터가 있는지 확인해요.

오버헤드가 클 때 (High overhead)

  • 샘플링 사용 고려: -Dotel.traces.sampler=parentbased_traceidratio -Dotel.traces.sampler.arg=0.1
  • 배치 스팬 프로세서(기본값) 사용
  • 필요하면 배치 프로세서 설정 조정

예시 (Examples)

완전한 예시는 다음 파일을 참고해요.

  • TracingExampleTest.java — 종합적인 사용 예시
  • OpenTelemetryTracingTest.java — API 사용을 보여주는 단위 테스트

API 참조 (API Reference)

주요 클래스 (Main Classes)

  • OpenTelemetryProducerInterceptor — 추적용 프로듀서 인터셉터
  • OpenTelemetryConsumerInterceptor — 추적용 컨슈머 인터셉터
  • TracingContext — 스팬 생성과 컨텍스트 전파를 위한 유틸리티 메서드
  • TracingProducerBuilder — 메시지에 트레이스 컨텍스트를 주입하는 헬퍼

추가 자료 (Additional Resources)

  • OpenTelemetry Java 문서
  • W3C Trace Context 규격
  • Pulsar 문서

더 알아보기 (Learn more)