基于OpenTelemetry构建微信机器人系统的全链路追踪体系
·
基于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 中查看耗时分布与错误节点,极大提升微信机器人系统的可观测性与故障定位效率。
更多推荐
所有评论(0)