You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现跨模块全链路追踪并关联Log4j日志与Zipkin/Jaeger链路ID?

Alright, let's walk through implementing full distributed tracing with either Zipkin or Jaeger for your system, plus linking Log4j logs to trace IDs to speed up troubleshooting. Here's a tailored solution based on your setup:

1. Base Environment Setup

First, pick your tracing backend (both are OpenTracing-compatible; we'll cover both options briefly):

  • Zipkin: Deploy via Docker for quick testing:
    docker run -d -p 9411:9411 openzipkin/zipkin
    
  • Jaeger: Deploy the all-in-one image for easy setup:
    docker run -d -p 6831:6831/udp -p 16686:16686 jaegertracing/all-in-one:latest
    
2. Module 1: Tomcat-hosted Multi-endpoint Web App (Pushes to topic1)

Assuming this is a Spring Boot Web app (adjust for traditional Java Web if needed):

Step 1: Add Dependencies

Include these in your pom.xml (Maven) or build.gradle:

<!-- For Zipkin tracing -->
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-zipkin</artifactId>
</dependency>
<!-- For Kafka integration -->
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
<!-- For Jaeger, replace Zipkin dependency with this -->
<dependency>
    <groupId>io.opentracing.contrib</groupId>
    <artifactId>opentracing-spring-jaeger-web-starter</artifactId>
</dependency>

Step 2: Configure Application Properties

Add these to application.properties:

# Zipkin config
spring.zipkin.base-url=http://localhost:9411
spring.sleuth.sampler.probability=1.0 # Full sampling (adjust for production)

# Jaeger config (replace above if using Jaeger)
opentracing.jaeger.service-name=web-app-module
opentracing.jaeger.udp-sender.host=localhost
opentracing.jaeger.udp-sender.port=6831

# Kafka config
spring.kafka.bootstrap-servers=your-kafka-broker-address:9092

Step 3: Push Messages with Tracing Context

Spring Sleuth automatically injects trace/span IDs into Kafka message headers when using KafkaTemplate. If you need manual control (e.g., non-Spring Kafka clients):

@Autowired
private Tracer tracer;
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;

public void pushToTopic1(String payload) {
    // Start a new span for the message send operation
    Span sendSpan = tracer.nextSpan().name("web-app:push-to-topic1").start();
    try (Scope scope = tracer.scopeManager().activate(sendSpan)) {
        // Inject trace ID into Kafka headers
        Headers headers = new RecordHeaders();
        headers.add("uber-trace-id", tracer.currentSpan().context().traceIdString().getBytes());
        ProducerRecord<String, String> record = new ProducerRecord<>("topic1", null, headers, null, payload);
        kafkaTemplate.send(record);
    } finally {
        sendSpan.finish();
    }
}

For traditional Java Web apps (non-Spring), add an OpenTracing Servlet Filter to initialize traces, then manually inject trace IDs into Kafka headers as above.

3. Module 2: Executable Jar with Kafka Consumers (K1 & K2)

This module runs two consumers: K1 processes topic1 and pushes to topic2; K2 processes topic2. We need to extract tracing context from Kafka headers and propagate it through the pipeline.

Step 1: Add Dependencies

Same as Module 1 (use Zipkin or Jaeger dependency based on your choice).

Step 2: Consumer Logic with Tracing

K1 Consumer (topic1 → topic2)

@Autowired
private Tracer tracer;
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;

@KafkaListener(topics = "topic1", groupId = "consumer-group-k1")
public void consumeTopic1(ConsumerRecord<String, String> record) {
    // Extract trace ID from Kafka headers
    Header traceIdHeader = record.headers().lastHeader("uber-trace-id");
    if (traceIdHeader == null) {
        // Handle missing trace ID (start a new trace if needed)
        Span span = tracer.nextSpan().name("k1:consume-topic1").start();
        processAndPushToTopic2(record.value(), span);
        return;
    }

    // Reconstruct parent span context
    String traceId = new String(traceIdHeader.value());
    SpanContext parentContext = tracer.extract(Format.Builtin.TEXT_MAP, 
        new TextMapAdapter(Map.of("uber-trace-id", traceId)));
    
    // Start child span for processing
    Span processSpan = tracer.nextSpan(parentContext).name("k1:process-topic1").start();
    try (Scope scope = tracer.scopeManager().activate(processSpan)) {
        String processedPayload = processTopic1Message(record.value());
        pushToTopic2(processedPayload);
    } finally {
        processSpan.finish();
    }
}

private void pushToTopic2(String payload) {
    Span sendSpan = tracer.nextSpan().name("k1:push-to-topic2").start();
    try (Scope scope = tracer.scopeManager().activate(sendSpan)) {
        Headers headers = new RecordHeaders();
        headers.add("uber-trace-id", tracer.currentSpan().context().traceIdString().getBytes());
        ProducerRecord<String, String> record = new ProducerRecord<>("topic2", null, headers, null, payload);
        kafkaTemplate.send(record);
    } finally {
        sendSpan.finish();
    }
}

K2 Consumer (topic2 Final Processing)

@Autowired
private Tracer tracer;

@KafkaListener(topics = "topic2", groupId = "consumer-group-k2")
public void consumeTopic2(ConsumerRecord<String, String> record) {
    Header traceIdHeader = record.headers().lastHeader("uber-trace-id");
    if (traceIdHeader == null) {
        Span span = tracer.nextSpan().name("k2:consume-topic2").start();
        finalProcessTopic2Message(record.value(), span);
        return;
    }

    String traceId = new String(traceIdHeader.value());
    SpanContext parentContext = tracer.extract(Format.Builtin.TEXT_MAP, 
        new TextMapAdapter(Map.of("uber-trace-id", traceId)));
    
    Span processSpan = tracer.nextSpan(parentContext).name("k2:final-process-topic2").start();
    try (Scope scope = tracer.scopeManager().activate(processSpan)) {
        finalProcessTopic2Message(record.value());
    } finally {
        processSpan.finish();
    }
}

To tie logs directly to trace IDs (critical for quick troubleshooting):

Step 1: Update Log4j2 Configuration

Modify your log4j2.xml to include trace/span IDs in the log pattern:

<Appenders>
    <Console name="Console" target="SYSTEM_OUT">
        <PatternLayout pattern="%d{HH:mm:ss.SSS} [%t] %-5level %logger{36} - traceId: %X{traceId}, spanId: %X{spanId} - %msg%n"/>
    </Console>
    <!-- Add file appender with the same pattern if needed -->
</Appenders>

Step 2: Auto-populate MDC (Spring Boot)

Spring Sleuth automatically adds traceId and spanId to the MDC (Mapped Diagnostic Context), so no extra code is needed.

Step 3: Manual MDC Population (Non-Spring)

If you're not using Spring, manually inject trace/span IDs into MDC when activating spans:

// When starting a span
Span span = tracer.nextSpan().start();
try (Scope scope = tracer.scopeManager().activate(span)) {
    MDC.put("traceId", span.context().traceIdString());
    MDC.put("spanId", span.context().spanIdString());
    // Your business logic here
} finally {
    MDC.remove("traceId");
    MDC.remove("spanId");
    span.finish();
}
5. Validate & Troubleshoot
  1. Start all components: Tracing backend, Kafka, Module 1, Module 2
  2. Trigger a message from Module 1's web endpoint
  3. Access the tracing UI:
    • Zipkin: http://localhost:9411
    • Jaeger: http://localhost:16686
  4. Search for the trace ID (from logs or UI) to see the full flow:
    Web Endpoint → push-to-topic1 → k1:process-topic1 → push-to-topic2 → k2:final-process-topic2
  5. When errors occur, grab the trace ID from logs and use the tracing UI to pinpoint which step failed, check latency, or view error details.

内容的提问来源于stack exchange,提问作者user1808266

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 04:22:48