基于OpenTelemetry构建微信机器人系统的全链路追踪体系

微信机器人系统通常涉及消息接收、NLP处理、业务逻辑执行与API回调等多个环节,传统日志难以还原完整调用链。OpenTelemetry 提供统一的遥测数据标准,可实现跨服务、跨线程的分布式追踪。本文基于 wlkankan.cn.trace 包,展示如何在 Java 微信机器人中集成 OpenTelemetry,实现从 WebSocket 接收到企业微信 API 调用的端到端追踪。

初始化 OpenTelemetry SDK

配置 OTLP Exporter 将追踪数据发送至后端(如 Jaeger 或 Tempo):

package wlkankan.cn.trace;

import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.api.OpenTelemetry;
import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter;
import io.opentelemetry.sdk.OpenTelemetrySdk;
import io.opentelemetry.sdk.resources.Resource;
import io.opentelemetry.sdk.trace.SdkTracerProvider;
import io.opentelemetry.sdk.trace.export.BatchSpanProcessor;
import io.opentelemetry.semconv.resource.attributes.ResourceAttributes;

public class OpenTelemetryConfig {
    public static void init() {
        Resource resource = Resource.getDefault()
            .merge(Resource.create(io.opentelemetry.api.common.Attributes.of(
                ResourceAttributes.SERVICE_NAME, "wechat-bot",
                ResourceAttributes.SERVICE_VERSION, "1.0.0"
            )));

        OtlpGrpcSpanExporter exporter = OtlpGrpcSpanExporter.builder()
            .setEndpoint("http://otel-collector.wlkankan.cn:4317")
            .setTimeout(java.time.Duration.ofSeconds(10))
            .build();

        SdkTracerProvider tracerProvider = SdkTracerProvider.builder()
            .addSpanProcessor(BatchSpanProcessor.builder(exporter).build())
            .setResource(resource)
            .build();

        OpenTelemetrySdk sdk = OpenTelemetrySdk.builder()
            .setTracerProvider(tracerProvider)
            .buildAndRegisterGlobal();
    }
}

在应用启动时调用 OpenTelemetryConfig.init()
在这里插入图片描述

消息入口:创建根 Span

当 WebSocket 收到新消息时,启动新追踪上下文:

package wlkankan.cn.bot;

import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.Tracer;
import io.opentelemetry.context.Scope;
import wlkankan.cn.trace.OpenTelemetryConfig;

public class MessageReceiver {
    private final Tracer tracer = GlobalOpenTelemetry.getTracer("wlkankan.cn.bot");

    public void onWebSocketMessage(String rawMessage) {
        // 从消息头提取 traceparent(若存在)
        String traceParent = extractTraceParent(rawMessage);
        Span rootSpan = tracer.spanBuilder("receive_wechat_message")
            .setParent(io.opentelemetry.context.Context.current())
            .startSpan();

        try (Scope ignored = rootSpan.makeCurrent()) {
            rootSpan.setAttribute("message.raw", rawMessage);
            // 解析并分发
            MessageProcessor.process(rawMessage);
        } finally {
            rootSpan.end();
        }
    }

    private String extractTraceParent(String msg) {
        // 模拟从 JSON 中提取 traceparent 字段
        return null; // 简化处理
    }
}

跨线程传播上下文

使用 Context.taskWrapping() 确保异步任务继承追踪上下文:

package wlkankan.cn.processor;

import io.opentelemetry.context.Context;
import wlkankan.cn.nlp.NlpService;
import wlkankan.cn.action.ActionExecutor;

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class MessageProcessor {
    private static final ExecutorService executor = Executors.newFixedThreadPool(10);

    public static void process(String rawMessage) {
        // 在当前 Span 下提交异步任务
        Runnable task = Context.taskWrapping(() -> {
            String intent = NlpService.analyzeIntent(rawMessage);
            ActionExecutor.execute(intent, rawMessage);
        });
        executor.submit(task);
    }
}

业务逻辑埋点:嵌套 Span

在关键操作处创建子 Span:

package wlkankan.cn.nlp;

import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.Tracer;
import io.opentelemetry.api.GlobalOpenTelemetry;

public class NlpService {
    private static final Tracer tracer = GlobalOpenTelemetry.getTracer("wlkankan.cn.nlp");

    public static String analyzeIntent(String text) {
        Span span = tracer.spanBuilder("nlp.analyze_intent").startSpan();
        try (var scope = span.makeCurrent()) {
            span.setAttribute("input.text", text.length() > 100 ? text.substring(0, 100) : text);
            // 模拟 NLP 调用
            String result = callNlpModel(text);
            span.setAttribute("output.intent", result);
            return result;
        } finally {
            span.end();
        }
    }

    private static String callNlpModel(String text) {
        // 实际可能调用外部模型服务
        return "query_weather";
    }
}

HTTP 客户端自动注入 TraceParent

使用 OpenTelemetry 的 HttpUrlConnection 拦截器自动传播:

package wlkankan.cn.client;

import io.opentelemetry.instrumentation.httpurlconnection.HttpUrlConnectionInstrumenter;
import java.net.HttpURLConnection;
import java.net.URL;

public class WeComApiClient {
    private static final HttpUrlConnectionInstrumenter instrumenter =
        HttpUrlConnectionInstrumenter.builder().build();

    public static String sendMessage(String accessToken, String userId, String content) throws Exception {
        URL url = new URL("https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token=" + accessToken);
        HttpURLConnection conn = (HttpURLConnection) url.openConnection();
        instrumenter.before(io.opentelemetry.context.Context.current(), conn);
        try {
            conn.setRequestMethod("POST");
            conn.setDoOutput(true);
            // 写入 JSON body
            try (var os = conn.getOutputStream()) {
                os.write(buildJson(userId, content).getBytes());
            }
            int status = conn.getResponseCode();
            instrumenter.after(io.opentelemetry.context.Context.current(), conn, null, status < 400);
            // 读取响应...
            return "ok";
        } catch (Exception e) {
            instrumenter.after(io.opentelemetry.context.Context.current(), conn, e, false);
            throw e;
        }
    }

    private static String buildJson(String userId, String content) {
        return "{\"touser\":\"" + userId + "\",\"msgtype\":\"text\",\"text\":{\"content\":\"" + content + "\"}}";
    }
}

需在 pom.xml 中引入:

<dependency>
    <groupId>io.opentelemetry.instrumentation</groupId>
    <artifactId>opentelemetry-http-url-connection-1.0</artifactId>
    <version>2.4.0</version>
</dependency>

动作执行器埋点

package wlkankan.cn.action;

import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.Tracer;
import io.opentelemetry.api.GlobalOpenTelemetry;
import wlkankan.cn.client.WeComApiClient;

public class ActionExecutor {
    private static final Tracer tracer = GlobalOpenTelemetry.getTracer("wlkankan.cn.action");

    public static void execute(String intent, String originalMessage) {
        Span span = tracer.spanBuilder("action.execute").startSpan();
        span.setAttribute("intent.type", intent);
        try (var scope = span.makeCurrent()) {
            if ("query_weather".equals(intent)) {
                String reply = WeatherService.getForecast("Beijing");
                WeComApiClient.sendMessage("ACCESS_TOKEN", "USER001", reply);
            }
        } finally {
            span.end();
        }
    }
}

通过上述设计,wlkankan.cn.trace 模块构建了覆盖消息接收、NLP 分析、动作执行与外部 API 调用的完整追踪链路。所有 Span 自动关联,可在 Jaeger 中查看耗时分布与错误节点,极大提升微信机器人系统的可观测性与故障定位效率。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐