学习目标
学完本章你应该能够:
- 跑通
dogapm基础框架:用函数式 option 模式初始化 MySQL / Redis,理解globalStart/globalClose全局注册表与endPoint的优雅启停机制。 - 用
dogapm封装的HttpServer/GrpcServer起服务,讲清拦截器(Interceptor)在 client / server 两侧的位置,以及它为什么是做监控埋点的天然切入点。 - 说清 OpenTelemetry 落地的"三部曲":初始化 SDK(Propagator / TracerProvider / MeterProvider)→ 用
otelhttp自动测 HTTP 边缘 → 手写自定义 Span 与 Metric 测业务内部。 - 解释为什么要把链路追踪数据发到 Jaeger 而不是只看控制台 stdout,以及 OTLP 上报用的 HTTP
4318端口是什么。 - 面试能讲清「监控 Metrics / 日志 Logs / 链路 Traces」三件套各自回答什么问题,以及
trace_id如何跨进程传播。
前置知识:
- Go 基础、
database/sql与redis客户端基本用法。 - gRPC 基础(拦截器概念)、Docker 基础(本地起组件)。
- 一点点可观测性常识:指标、日志、链路分别长什么样。
本章你会动手做的事:
- 用
docker-compose起 MySQL / Redis,跑TestFrame验证连接初始化成功。 - 把 OpenTelemetry 的
dice示例跑起来,浏览器访问/roll,数控制台里是不是出现了两个 Span。 - 起 Jaeger
all-in-one,把 Exporter 从 stdout 切到 OTLP HTTP4318,在 Jaeger UI 里看到上报的 trace。
基础代码框架
类比:
dogapm这个基础框架就像一间工厂的"总配电箱"。MySQL、Redis、HTTP 服务、gRPC 服务都是插在配电箱上的设备——每个设备启动时自己往globalStart名单里报个到,关停时往globalClose名单里挂个号;工厂拉闸(endPoint.Shutdown)时,配电箱就按顺序挨个通知它们断电。你新加一种设备,只要往两个名单里登记,启停流程全自动接管。
安装mysql和redis
本地使用docker-compose安装mysql和redis
version:"3"
services:
mysql:
restart: always
image: mysql:latest
container_name: mysql
environment:
- "MYSQL_ROOT_PASSWORD"
ports:
- "3306:3306"
redis:
restart: always
image: redis
container_name: redis
ports:
- "6379:6379"
创建数据库
ordersvc :订单表
create table t_order
(
id bigint auto_increment primary key,
order_id varchar(255) not null,
ctime timestamp default CURRENT_TIMESTAMP not null,
utime timestamp default CURRENT_TIMESTAMP not null,
sku_id bigint not null,
num int not null,
price int not null,
uid bigint not null,
constraint order_pk2
unique (order_id)
);
skusvc: 商品表
create table t_sku
(
id bigint auto_increment primary key,
name varchar(10) not null,
price int null,
ctime timestamp default CURRENT_TIMESTAMP not null,
utime timestamp default CURRENT_TIMESTAMP not null,
num int null
);
usersvc: 用户表
create table t_user
(
id bigint auto_increment primary key,
name varchar(20) not null,
ctime timestamp default CURRENT_TIMESTAMP not null,
utime timestamp default CURRENT_TIMESTAMP not null
);
mysql、redis初始化
package dogapm
import (
"context"
"database/sql"
_ "github.com/go-sql-driver/mysql"
"github.com/redis/go-redis/v9"
)
// 程序的框架初始化
type frame struct {
DB *sql.DB
RedisDB *redis.Client
}
var Frame = &frame{}
type option func(f *frame)
func FrameMysqlDBOption(addr string) option {
return func(f *frame) {
var err error
f.DB, err = sql.Open("mysql", addr)
if err != nil {
panic(err)
}
err = f.DB.Ping()
if err != nil {
panic(err)
}
}
}
func FrameRedisDBOption(addr string, pwd string) option {
return func(f *frame) {
client := redis.NewClient(&redis.Options{
Addr: addr,
DB: 0,
Password: pwd,
})
result, err := client.Ping(context.Background()).Result()
if err != nil {
panic(err)
}
if result != "PONG" {
panic("redis client init fail")
}
f.RedisDB = client
}
}
func (f *frame) Init(options ...option) {
for _, opt := range options {
opt(f)
}
}
启动所有服务
package dogapm
import (
"os"
"os/signal"
"syscall"
)
type start interface {
Start()
}
type close interface {
Close()
}
var (
globalStart = make([]start, 0)
globalClose = make([]close, 0)
)
var EndPotin = endPoint{stop: make(chan int, 1)}
type endPoint struct {
stop chan int
}
func (e *endPoint) Start() {
for _, starter := range globalStart {
starter.Start()
}
go func() {
quit := make(chan os.Signal)
signal.Notify(quit, syscall.SIGQUIT, syscall.SIGTERM, syscall.SIGINT)
<-quit
e.Shutdown()
}()
<-e.stop
}
func (e *endPoint) Shutdown() {
for _, closer := range globalClose {
closer.Close()
}
e.stop <- 1
}
下面这张图把 endPoint 的"启动 → 监听信号 → 优雅关停"生命周期画清楚,这也是 globalStart / globalClose 两个注册表真正被消费的地方:
flowchart TD
A[endPoint.Start 启动] --> B[遍历 globalStart 逐个调用 Start]
B --> C[起 goroutine 监听 SIGQUIT SIGTERM SIGINT]
C --> D{收到退出信号?}
D -- 是 --> E[调用 Shutdown]
E --> F[遍历 globalClose 逐个调用 Close]
F --> G[向 stop channel 发 1 主流程退出]
D -- 否 持续服务 --> C初始化HTTP
package dogapm
import (
"context"
"net/http"
)
/*
封装Http请求
mux 注册路由
Server启动服务
*/
type HttpServer struct {
mux *http.ServeMux
*http.Server
}
func NewHttpServer(addr string) *HttpServer {
mux := http.NewServeMux()
server := &http.Server{Addr: addr, Handler: mux}
httpServer := &HttpServer{mux: mux, Server: server}
globalStart = append(globalStart, httpServer) // 注册启动服务
globalClose = append(globalClose, httpServer) // 注册关闭服务
return httpServer
}
/*
将一个函数作为handler
*/
func (h *HttpServer) HandleFunc(pattern string, handler func(w http.ResponseWriter, r *http.Request)) {
h.mux.HandleFunc(pattern, handler)
}
/*
将一个struct必须实现ServeHTTP(w http.Response,r *http.Request)方法
作为handler
*/
func (h *HttpServer) Handle(pattern string, handler http.Handler) {
h.mux.Handle(pattern, handler)
}
// 启动服务
func (h *HttpServer) Start() {
go func() {
err := h.Server.ListenAndServe()
if err != nil {
panic(err)
}
}()
}
// 关闭服务
func (h *HttpServer) Close() {
h.Server.Shutdown(context.TODO())
}
初始化GRPC
1. Client
package dogapm
import (
"context"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
type GrpcClient struct {
*grpc.ClientConn
}
func NewGrpcClient(addr string) *GrpcClient {
dial, err := grpc.Dial(addr,
grpc.WithUnaryInterceptor(unaryInterceptor()),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
if err != nil {
panic(err)
}
return &GrpcClient{ClientConn: dial}
}
func unaryInterceptor() grpc.UnaryClientInterceptor {
return func(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
err := invoker(ctx, method, req, reply, cc, opts...)
return err
}
}
2. Server
package dogapm
import (
"context"
"google.golang.org/grpc"
"net"
)
type grpcServer struct {
*grpc.Server
addr string
}
func NewGrpcServer(addr string) *grpcServer {
svc := grpc.NewServer(grpc.UnaryInterceptor(unaryServerInterceptor()))
server := &grpcServer{
Server: svc,
addr: addr,
}
globalStart = append(globalStart, server) // 注册启动服务
globalClose = append(globalClose, server) // 注册关闭服务
return server
}
// 拦截器
func unaryServerInterceptor() grpc.UnaryServerInterceptor {
return func(ctx context.Context, req any, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (resp any, err error) {
return nil, nil
}
}
// 关闭服务
func (g *grpcServer) Close() {
g.Server.GracefulStop()
}
// 开启服务
func (g *grpcServer) Start() {
listen, err := net.Listen("tcp", g.addr)
if err != nil {
panic(err)
}
go func() {
err := g.Serve(listen)
if err != nil {
panic(err)
}
}()
}
HTTP状态码
package dogapm
import (
"encoding/json"
"fmt"
"net/http"
)
const jsonType string = "application/json"
type Status struct {
Code int `json:"code"`
Message string `json:"message"`
Body any `json:"body"`
}
type httpStatus struct {
}
var HttpStatus = &httpStatus{}
func (h *httpStatus) Success(w http.ResponseWriter) {
status := &Status{
Code: http.StatusOK,
Message: "Success",
}
marshal, err := json.Marshal(status)
if err != nil {
fmt.Println(err)
return
}
w.Header().Set("Content-Type", jsonType)
w.WriteHeader(http.StatusOK)
_, err = w.Write(marshal)
if err != nil {
fmt.Println(err)
return
}
}
func (h *httpStatus) SuccessBody(w http.ResponseWriter, msg string, body any) {
status := &Status{
Code: http.StatusOK,
Message: msg,
Body: body,
}
marshal, err := json.Marshal(status)
if err != nil {
fmt.Println(err)
return
}
w.Header().Set("Content-Type", jsonType)
w.WriteHeader(http.StatusOK)
_, err = w.Write(marshal)
if err != nil {
fmt.Println(err)
return
}
}
func (h *httpStatus) Fail(w http.ResponseWriter, msg string, body any) {
status := &Status{
Code: http.StatusBadRequest,
Message: msg,
Body: body,
}
marshal, err := json.Marshal(status)
if err != nil {
fmt.Println(err)
return
}
w.Header().Set("Content-Type", jsonType)
w.WriteHeader(http.StatusBadRequest)
_, err = w.Write(marshal)
if err != nil {
fmt.Println(err)
return
}
}
func (h *httpStatus) Error(w http.ResponseWriter, msg string, body any) {
status := &Status{
Code: http.StatusInternalServerError,
Message: msg,
Body: body,
}
marshal, err := json.Marshal(status)
if err != nil {
fmt.Println(err)
return
}
w.Header().Set("Content-Type", jsonType)
w.WriteHeader(http.StatusInternalServerError)
_, err = w.Write(marshal)
if err != nil {
fmt.Println(err)
return
}
}
日志
使用slog,生成带色彩的日志
package dogapm
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"os"
"runtime"
"strings"
)
// 颜色代码
const (
colorRed = "\033[31m"
colorBlue = "\033[34m"
colorGreen = "\033[32m"
colorYellow = "\033[33m"
colorReset = "\033[0m"
)
// Logger 定义日志接口
type Logger interface {
Debug(ctx context.Context, action string, msg string, args ...any)
Info(ctx context.Context, action string, msg string, args ...any)
Warn(ctx context.Context, action string, msg string, args ...any)
Error(ctx context.Context, action string, msg string, err error, args ...any)
}
// ColorfulJSONHandler 自定义支持彩色输出的 JSON 处理器
type ColorfulJSONHandler struct {
slog.Handler
}
// NewColorfulJSONHandler 创建一个新的 ColorfulJSONHandler 实例
func NewColorfulJSONHandler() *ColorfulJSONHandler {
return &ColorfulJSONHandler{
Handler: slog.NewJSONHandler(os.Stdout, nil),
}
}
// 处理日志记录
func (h *ColorfulJSONHandler) Handle(ctx context.Context, r slog.Record) error {
var color string
switch r.Level {
case slog.LevelDebug:
color = colorBlue
case slog.LevelInfo:
color = colorGreen
case slog.LevelWarn:
color = colorYellow
case slog.LevelError:
color = colorRed
default:
color = colorReset
}
var buf []byte
r.Attrs(func(a slog.Attr) bool {
if a.Value.Kind() == slog.KindGroup {
gv := a.Value.Group()
for _, ga := range gv {
buf = appendAttr(buf, ga)
}
} else {
buf = appendAttr(buf, a)
}
return true
})
record := map[string]any{
"time": r.Time,
"level": r.Level.String(),
"message": r.Message,
}
if err := json.Unmarshal(buf, &record); err != nil {
return err
}
jsonData, err := json.Marshal(record)
if err != nil {
return err
}
fmt.Fprint(os.Stdout, color, string(jsonData), colorReset, "\n")
return nil
}
// 追加属性到缓冲区
func appendAttr(buf []byte, a slog.Attr) []byte {
if len(buf) > 0 {
buf = append(buf, ',')
}
keyBytes := []byte(jsonEscape(a.Key))
buf = append(buf, '"')
buf = append(buf, keyBytes...)
buf = append(buf, '"', ':')
switch a.Value.Kind() {
case slog.KindString:
valBytes := []byte(jsonEscape(a.Value.String()))
buf = append(buf, '"')
buf = append(buf, valBytes...)
buf = append(buf, '"')
default:
valBytes, _ := json.Marshal(a.Value.Any())
buf = append(buf, valBytes...)
}
return buf
}
// 转义 JSON 字符串
func jsonEscape(s string) string {
return strings.ReplaceAll(strings.ReplaceAll(s, `\`, `\\`), `"`, `\"`)
}
// 获取调用者信息
func getCallerInfo() (file string, line int) {
_, file, line, _ = runtime.Caller(3)
return
}
// SlogColorfulJSONLogger 实现 Logger 接口
type SlogColorfulJSONLogger struct {
logger *slog.Logger
}
// NewSlogColorfulJSONLogger 创建一个新的 SlogColorfulJSONLogger 实例
func NewJSONLogger() *SlogColorfulJSONLogger {
handler := NewColorfulJSONHandler()
logger := slog.New(handler)
return &SlogColorfulJSONLogger{
logger: logger,
}
}
// Debug 记录调试级别的日志
func (l *SlogColorfulJSONLogger) Debug(ctx context.Context, action string, msg string, args ...any) {
file, line := getCallerInfo()
args = append([]any{"action", action, "file", file, "line", line}, args...)
l.logger.Debug(msg, args...)
}
// Info 记录信息级别的日志
func (l *SlogColorfulJSONLogger) Info(ctx context.Context, action string, msg string, args ...any) {
file, line := getCallerInfo()
args = append([]any{"action", action, "file", file, "line", line}, args...)
l.logger.Info(msg, args...)
}
// Warn 记录警告级别的日志
func (l *SlogColorfulJSONLogger) Warn(ctx context.Context, action string, msg string, args ...any) {
file, line := getCallerInfo()
args = append([]any{"action", action, "file", file, "line", line}, args...)
l.logger.Warn(msg, args...)
}
// Error 记录错误级别的日志
func (l *SlogColorfulJSONLogger) Error(ctx context.Context, action string, msg string, err error, args ...any) {
file, line := getCallerInfo()
args = append([]any{"action", action, "file", file, "line", line, "error", err.Error()}, args...)
l.logger.Error(msg, args...)
}
// NewCustomError 自定义错误类型
func NewCustomError(message string) error {
return &customError{message: message}
}
// customError 自定义错误结构体
type customError struct {
message string
}
// Error 实现 error 接口
func (e *customError) Error() string {
return e.message
}
sql
package dogapm
import "database/sql"
type dbUtil struct {
}
var DBUtil = &dbUtil{}
func (d *dbUtil) Query(rows *sql.Rows, err any) []map[string]any {
if err != nil {
return nil
}
if rows == nil {
return []map[string]any{}
}
defer rows.Close()
columns, _ := rows.Columns()
scanArgs := make([]any, len(columns))
values := make([]any, len(columns))
for j := range values {
scanArgs[j] = &values[j]
}
res := make([]map[string]any, 0, 5)
for rows.Next() {
record := make(map[string]any)
rows.Scan(scanArgs...)
for i, col := range values {
if col != nil {
switch col.(type) {
case []byte:
record[columns[i]] = string(col.([]byte))
default:
record[columns[i]] = col
}
}
}
res = append(res, record)
}
return res
}
func (d *dbUtil) QueryFirst(rows *sql.Rows, err any) map[string]any {
res := d.Query(rows, err)
if len(res) > 0 {
return res[0]
}
return nil
}
测试框架
package dogapm
import (
"context"
"fmt"
"net/http"
pb "proto/hello"
"testing"
"time"
)
// mysql db
func TestFrame(t *testing.T) {
mysqlDB := FrameMysqlDBOption("root:1230123@tcp(127.0.0.1:3306)/ordersvc")
redisDB := FrameRedisDBOption("127.0.0.1:6379", "1230123")
Frame.Init(mysqlDB, redisDB)
}
// http
func TestHttp(t *testing.T) {
server := NewHttpServer(":8080")
server.HandleFunc("/test", func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte("ok"))
})
server.Start()
time.Sleep(3600 * time.Second)
}
// grpc
func TestGrpc(t *testing.T) {
go func() {
server := NewGrpcServer(":8080")
pb.RegisterHelloServiceServer(server, &Hello{})
server.Start()
}()
grpcClient := NewGrpcClient("127.0.0.1:8080")
receive, err := pb.NewHelloServiceClient(grpcClient).Receive(context.TODO(), &pb.HelloMsg{Msg: "hello world"})
if err != nil {
t.Fatal(err)
}
fmt.Println(receive.Msg)
}
type Hello struct {
pb.UnimplementedHelloServiceServer
}
func (h *Hello) Receive(ctx context.Context, req *pb.HelloMsg) (*pb.HelloMsg, error) {
return req, nil
}
链路追踪的基本概念
类比:链路追踪就像给一次用户请求办一张"登机牌 + 全程行程单"。从网关进门起,每经过一个服务就在行程单上盖一个章(Span),章与章之间有父子关系(谁调谁)。哪一段飞得慢、在哪坠机(panic),一眼就能从行程单上定位——这就是 Traces 相对 Logs / Metrics 最独特的价值:还原一次请求横跨多个服务的完整路径。
下面这张图把"一次请求跨三个服务、各自访问存储、最后汇成一条 Trace 上报"讲清楚:
flowchart LR
A[客户端请求] --> B[网关 Span 根]
B --> C[订单服务 Span]
C --> D[用户服务 Span]
C --> E[商品服务 Span]
D --> F[(MySQL)]
E --> G[(Redis)]
C --> H[上报 Trace 到 Jaeger / OTel Collector]监控三件套:Metrics / Logs / Traces
| 类型 | 回答的问题 | 形态举例 |
|---|---|---|
| Metrics(指标) | “发生了多少次 / 多频繁 / 多慢” | QPS、P99 延迟、panic_total 计数 |
| Logs(日志) | “具体发生了什么事件” | 结构化 ERROR 日志、panic 堆栈 |
| Traces(链路) | “一次请求慢在哪、错在哪一段” | 跨服务的调用链 Span 树 |
OpenTelemetry
创建http程序
在同一目录下创建 main.go 文件,并添加以下代码。
package main
import (
"fmt"
"log"
"math/rand"
"net/http"
)
func main() {
handler := http.NewServeMux()
handler.HandleFunc("/roll", func(writer http.ResponseWriter, request *http.Request) {
number := 1 + rand.Intn(6)
_, _ = fmt.Fprintln(writer, number)
})
log.Fatal(http.ListenAndServe(":8080", handler))
}
添加 OpenTelemetry 测量仪器
接下来,我们将展示如何在示例应用程序中添加 OpenTelemetry 测量仪器。
1. 引入依赖
在你的Go项目中安装以下依赖包。
go get "go.opentelemetry.io/otel" \
"go.opentelemetry.io/otel/exporters/stdout/stdoutmetric" \
"go.opentelemetry.io/otel/exporters/stdout/stdouttrace" \
"go.opentelemetry.io/otel/propagation" \
"go.opentelemetry.io/otel/sdk/metric" \
"go.opentelemetry.io/otel/sdk/resource" \
"go.opentelemetry.io/otel/sdk/trace" \
"go.opentelemetry.io/otel/semconv/v1.24.0" \
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
这里安装的是 OpenTelemety SDK 组件和 net/http 测量仪器。如果要对不同的库进行网络请求检测,则需要安装相应的仪器库。
2. 初始化OpenTelemetry SDK
首先,我们将初始化OpenTelemetry SDK。任何想导出追踪数据的应用程序都必需完成这一步初始化。
新建一个otel.go文件,并在其中添加以下代码。
package main
import (
"context"
"errors"
"time"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/exporters/stdout/stdoutmetric"
"go.opentelemetry.io/otel/exporters/stdout/stdouttrace"
"go.opentelemetry.io/otel/propagation"
"go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/trace"
)
// setupOTelSDK 引导 OpenTelemetry pipeline。
// 如果没有返回错误,请确保调用 shutdown 进行适当清理。
func setupOTelSDK(ctx context.Context) (shutdown func(context.Context) error, err error) {
var shutdownFuncs []func(context.Context) error
// shutdown 会调用通过 shutdownFuncs 注册的清理函数。
// 调用产生的错误会被合并。
// 每个注册的清理函数将被调用一次。
shutdown = func(ctx context.Context) error {
var err error
for _, fn := range shutdownFuncs {
err = errors.Join(err, fn(ctx))
}
shutdownFuncs = nil
return err
}
// handleErr 调用 shutdown 进行清理,并确保返回所有错误信息。
handleErr := func(inErr error) {
err = errors.Join(inErr, shutdown(ctx))
}
// 设置传播器
prop := newPropagator()
otel.SetTextMapPropagator(prop)
// 设置 trace provider.
tracerProvider, err := newTraceProvider()
if err != nil {
handleErr(err)
return
}
shutdownFuncs = append(shutdownFuncs, tracerProvider.Shutdown)
otel.SetTracerProvider(tracerProvider)
// 设置 meter provider.
meterProvider, err := newMeterProvider()
if err != nil {
handleErr(err)
return
}
shutdownFuncs = append(shutdownFuncs, meterProvider.Shutdown)
otel.SetMeterProvider(meterProvider)
return
}
func newPropagator() propagation.TextMapPropagator {
return propagation.NewCompositeTextMapPropagator(
propagation.TraceContext{},
propagation.Baggage{},
)
}
func newTraceProvider() (*trace.TracerProvider, error) {
traceExporter, err := stdouttrace.New(
stdouttrace.WithPrettyPrint())
if err != nil {
return nil, err
}
traceProvider := trace.NewTracerProvider(
trace.WithBatcher(traceExporter,
// 默认为 5s。为便于演示,设置为 1s。
trace.WithBatchTimeout(time.Second)),
)
return traceProvider, nil
}
func newMeterProvider() (*metric.MeterProvider, error) {
metricExporter, err := stdoutmetric.New()
if err != nil {
return nil, err
}
meterProvider := metric.NewMeterProvider(
metric.WithReader(metric.NewPeriodicReader(metricExporter,
// 默认为 1m。为便于演示,设置为 3s。
metric.WithInterval(3*time.Second))),
)
return meterProvider, nil
}
如果不使用 tracing ,则可以省略相应的 TracerProvider 的初始化代码;
如果不使用 metrics ,则可以省略 MeterProvider 的初始化代码。
类比:OpenTelemetry SDK 就像一家"快递集散中心"。你的业务代码(Tracer / Meter)把包裹(Span / Metric)交给中心,中心先按批次打包(Batcher,默认 5s / 演示 1s),再交给具体的快递公司(Exporter)——stdout 是"自己家门口扔信箱看一眼",Jaeger 是"发往专业分拣仓库做全局查询"。
下面这张图把 SDK 初始化到数据导出的完整 pipeline 画清楚:
flowchart LR
A[业务代码 tracer.Meter] --> B[TracerProvider / MeterProvider]
B --> C[Batcher 批量打包]
C --> D[Exporter]
D --> E{导出目标}
E -- 本地演示 --> F[stdout 控制台]
E -- 生产 --> G[OTLP HTTP 4318 发往 Jaeger]
G --> H[Jaeger UI 查询]3. 测量 HTTP server
现在,我们已经初始化了OpenTelemetry SDK,可以测量HTTP服务器了。
按如下代码修改 main.go,加入设置 OpenTelemetry SDK 的代码,并使用 otelhttp 仪器库测量 HTTP 服务器:
package main
import (
"context"
"errors"
"log"
"net"
"net/http"
"os"
"os/signal"
"time"
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
)
func newHTTPHandler() http.Handler {
mux := http.NewServeMux()
// handleFunc 是 mux.HandleFunc 的替代品,。
// 它使用 http.route 模式丰富了 handler 的 HTTP 测量
handleFunc := func(pattern string, handlerFunc func(http.ResponseWriter, *http.Request)) {
// 为 HTTP 测量配置 "http.route"。
handler := otelhttp.WithRouteTag(pattern, http.HandlerFunc(handlerFunc))
mux.Handle(pattern, handler)
}
// Register handlers.
handleFunc("/roll", roll)
// 为整个服务器添加 HTTP 测量。
handler := otelhttp.NewHandler(mux, "/")
return handler
}
func main() {
if err := run(); err != nil {
log.Fatalln(err)
}
}
func run() (err error) {
// 平滑处理 SIGINT (CTRL+C) .
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt)
defer stop()
// 设置 OpenTelemetry.
otelShutdown, err := setupOTelSDK(ctx)
if err != nil {
return
}
// 妥善处理停机,确保无泄漏
defer func() {
err = errors.Join(err, otelShutdown(context.Background()))
}()
// 启动 HTTP server.
srv := &http.Server{
Addr: ":8080",
BaseContext: func(_ net.Listener) context.Context { return ctx },
ReadTimeout: time.Second,
WriteTimeout: 10 * time.Second,
Handler: newHTTPHandler(),
}
srvErr := make(chan error, 1)
go func() {
srvErr <- srv.ListenAndServe()
}()
// 等待中断.
select {
case err = <-srvErr:
// 启动 HTTP 服务器时出错.
return
case <-ctx.Done():
// 等待第一个 CTRL+C.
// 尽快停止接收信号通知.
stop()
}
// 调用 Shutdown 时,ListenAndServe 会立即返回 ErrServerClosed。
err = srv.Shutdown(context.Background())
return
}
4. 添加自定义测量
测量库可以捕捉系统边缘的遥测数据,例如入站和出站 HTTP 请求,但无法捕捉应用程序中的情况。因此需要编写一些自定义的手动仪器。
修改 roll.go,使用 OpenTelemetry API 包含定制的测量仪器:
package main
import (
"fmt"
"math/rand"
"net/http"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"
)
var (
tracer = otel.Tracer("roll")
meter = otel.Meter("roll")
rollCnt metric.Int64Counter
)
func init() {
var err error
rollCnt, err = meter.Int64Counter("dice.rolls",
metric.WithDescription("The number of rolls by roll value"),
metric.WithUnit("{roll}"))
if err != nil {
panic(err)
}
}
func roll(w http.ResponseWriter, r *http.Request) {
ctx, span := tracer.Start(r.Context(), "roll") // 开始 span
defer span.End() // 结束 span
number := 1 + rand.Intn(6)
rollValueAttr := attribute.Int("roll.value", number)
span.SetAttributes(rollValueAttr) // span 添加属性
// 摇骰子次数的指标 +1
rollCnt.Add(ctx, 1, metric.WithAttributes(rollValueAttr))
_, _ = fmt.Fprintln(w, number)
}
5. 运行应用程序
使用以下命令构建并运行应用程序:
go mod tidy
export OTEL_RESOURCE_ATTRIBUTES="service.name=dice,service.version=0.1.0"
go run .
使用浏览器中打开 http://127.0.0.1:8080/roll。向服务器发送请求时,你会在控制台显示的链路跟踪中看到两个 span。由仪器库生成的 span 跟踪向 /roll 路由发出请求的生命周期。名为 roll 的 span 是手动创建的,它是前面提到的 span 的子 span。
将链路追踪数据发送至 Jaeger
如果觉着控制台看的 span 不够直观,可以选择将链路追踪的数据发送至 Jaeger,通过 Jaeger UI 查看。
1. 启动 Jaeger
Jaeger 官方提供的 all-in-one 是为快速本地测试而设计的可执行文件。它包括 Jaeger UI、jaeger-collector、jaeger-query 和 jaeger-agent,以及一个内存存储组件。
启动 all-in-one 的最简单方法是使用发布到 DockerHub 的预置镜像(只需一条命令行)。
docker run --rm --name jaeger \
-e COLLECTOR_ZIPKIN_HOST_PORT=:9411 \
-p 6831:6831/udp \
-p 6832:6832/udp \
-p 5778:5778 \
-p 16686:16686 \
-p 4317:4317 \
-p 4318:4318 \
-p 14250:14250 \
-p 14268:14268 \
-p 14269:14269 \
-p 9411:9411 \
jaegertracing/all-in-one:1.55
然后你可以使用浏览器打开 http://localhost:16686 访问Jaeger UI。
容器公开以下端口:

我们这里使用 HTTP 协议的4318 端口上报链路追踪数据。
2. 上报至 Jaeger
安装 otlptracehttp 依赖包。
go get go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp
修改otel.go 代码,新增以下函数。
func newJaegerTraceProvider(ctx context.Context) (*trace.TracerProvider, error) {
// 创建一个使用 HTTP 协议连接本机Jaeger的 Exporter
traceExporter, err := otlptracehttp.New(ctx,
otlptracehttp.WithEndpoint("127.0.0.1:4318"),
otlptracehttp.WithInsecure())
if err != nil {
return nil, err
}
traceProvider := trace.NewTracerProvider(
trace.WithBatcher(traceExporter,
// 默认为 5s。为便于演示,设置为 1s。
trace.WithBatchTimeout(time.Second)),
)
return traceProvider, nil
}
并且按如下代码修改设置 trace provider 部分。
// 设置 trace provider.
//tracerProvider, err := newTraceProvider()
tracerProvider, err := newJaegerTraceProvider(ctx)
if err != nil {
handleErr(err)
return
}
shutdownFuncs = append(shutdownFuncs, tracerProvider.Shutdown)
otel.SetTracerProvider(tracerProvider)
再次构建并启动程序。
go run .
尝试访问一次 http://127.0.0.1:8080/roll ,确保修改后的服务能够正常运行。
3. 使用 Jaeger UI
使用浏览器打开 http://127.0.0.1:16686
的Jaeger UI界面。 在屏幕左侧的 service 下拉框中选中 dice后查找,即可看到上报的 trace 数据。

点击右侧的 trace 数据,即可查看详情。

完整代码请查看 https://github.com/Q1mi/dice
自测题与动手练习
自测题(合上书能答出来,才算懂):
dogapm的Frame用「函数式 option」模式初始化(如FrameMysqlDBOption),相比直接传一个塞了十几个字段的 config 结构体,好处是什么?globalStart/globalClose两个全局切片在NewHttpServer/NewGrpcServer里被 append。endPoint.Start和Shutdown分别怎么用它们?如果某次新增服务忘了 append 会怎样?- OpenTelemetry 里的
Propagator(TraceContext + Baggage)是干什么的?一次 HTTP 请求跨进程调另一个服务时,trace_id是怎么跟着请求"传过去"的? - 手写 Span 时
tracer.Start(ctx, "roll")返回的两个值分别是什么?如果忘了defer span.End(),这条 Span 还会正常上报吗? - 控制台 exporter(stdout)和 Jaeger exporter 分别在什么场景合适?OTLP 的
4318端口走的是什么协议?
动手练习(建议真做一遍):
- 用
docker-compose起 MySQL + Redis,把FrameMysqlDBOption/FrameRedisDBOption的地址换成你的,跑TestFrame验证Ping成功、连接初始化不 panic。 - 跑通
dice示例,go run .后浏览器访问/roll,数控制台里是不是出现了两个 Span(otelhttp 自动生成的根 Span + 你手写的roll子 Span)。 docker run起 Jaegerall-in-one,把newJaegerTraceProvider接上127.0.0.1:4318,再访问一次/roll,打开 http://localhost:16686 在service下拉框选dice找到上报的 trace 并点开看详情。
本章小结
dogapm用「函数式 option + 全局 Start/Close 注册表 +endPoint优雅启停」把 MySQL / Redis / HTTP / gRPC 初始化收口成一个可复用框架;拦截器(Interceptor)是做监控埋点的天然切入点。- 可观测性三件套各司其职:Metrics 回答"多少次/多频繁",Logs 回答"发生了什么",Traces 回答"一次请求慢在哪、错在哪一段"。
- OpenTelemetry 落地三步:初始化 SDK(Propagator + Tracer/Meter Provider + Exporter)→ 用仪器库(otelhttp)自动测边缘 → 手写 Span / Metric 测业务内部;批量导出靠 Batcher,目标是后端(Jaeger 等)。
trace_id靠 Propagator 注入请求头实现跨进程传播;本地开发用 stdout exporter 看结构,生产用 OTLP(HTTP4318/ gRPC4317)上报做可视化。
下一章可以把这套 OTel 接入 Kitex 的 middleware:把 panic 恢复、RPC 调用都变成可观测的 Span,与前面讲的 Panic 中间件真正串起来。