一、踩坑:异步任务里的TraceId丢失惨案

做过分布式系统的朋友应该都有过这种体验:线上报了个错误,日志里翻来翻去找不到完整链路——只看到某个服务报错了,却不知道这个请求是从哪个前端页面来的、中间经过了哪些服务、触发了哪些异步操作。这时候全链路追踪(就是把一个请求从开始到结束的所有环节串成一条线)就显得特别重要,而串起这条线的核心,就是一个叫TraceId的唯一标识,每个请求的所有环节都带着这个ID,就能把分散在不同服务、不同存储里的日志拼起来。

我之前就踩过一个特别闹心的坑:业务里有个“用户下单后自动触发物流推送”的逻辑,下单服务是同步的,推物流是异步的(用消息队列解耦,怕下单卡)。结果线上偶尔会出现:下单日志里有TraceId,但物流推送的日志里没这个ID,导致出问题时查不到是哪个下单请求触发的推送,只能盲查,效率极低。后来定位到原因:下单服务把下单成功的消息发到消息队列时,偶尔会漏传TraceId,异步消费的时候自然就丢了。

二、解决方案:用Go-Zero的Context传TraceId

后来我用Go-Zero框架解决了这个问题,核心思路是:利用Go-Zero自带的Context传播机制,把TraceId从同步请求(比如下单的HTTP请求)开始,一直传到异步任务(比如消息队列的消费),全程不丢。

先给大家解释下几个容易混的概念,怕新手看不懂:

  • TraceId:每个请求的唯一身份证,从请求进来就生成,所有环节都带它;
  • Context:Go语言里用来传请求相关数据的容器,比如请求ID、用户信息这些,不会乱;
  • Go-Zero:一个Go写的微服务框架,自带了很多好用的功能,比如自动生成代码、日志、全链路追踪的基础支持;
  • RPC:远程过程调用,简单说就是A服务调用B服务的方法,像调用自己的方法一样;
  • 消息队列:比如RabbitMQ、Kafka,用来解耦同步逻辑,比如下单后发消息到队列,消费端异步处理。

解决这个问题的核心逻辑分三步:

  1. 同步请求进来时,Go-Zero自动生成TraceId,放到Context里;
  2. 调用RPC时,Go-Zero自动把Context里的TraceId带到下游服务;
  3. 发消息到消息队列时,手动把Context里的TraceId取出来,放到消息的元数据里;消费消息时,再把TraceId塞回新的Context里,后面的操作就都能拿到了。

三、具体实现:从同步到异步的完整流程

我会用Go-Zero写一个完整的例子,包含HTTP下单服务、RPC库存服务、消息队列(用RabbitMQ)消费的物流推送服务,全程演示TraceId怎么传。

3.1 技术栈说明

所有例子统一用这个技术栈:

  • 语言:Go 1.20+
  • 框架:Go-Zero 1.5+
  • RPC:Go-Zero自带的gRPC实现
  • 消息队列:RabbitMQ 3.11+
  • 日志:Go-Zero自带的日志组件

3.2 第一步:同步请求的TraceId生成

首先写一个HTTP下单服务,请求进来时Go-Zero会自动生成TraceId,放到Context里。我们要做的就是在日志里打出来,确认存在。

// api/order/order.go 下单服务的HTTP处理逻辑
package order

import (
	"context"
	"github.com/zeromicro/go-zero/core/logx"
	"github.com/zeromicro/go-zero/rest/httpx"
	"net/http"
)

// OrderReq 下单请求参数
type OrderReq struct {
	GoodsID int64 `json:"goods_id"`
	UserID  int64 `json:"user_id"`
}

// OrderResp 下单响应
type OrderResp struct {
	TraceID string `json:"trace_id"` // 把TraceId返回给前端,方便调试
	Msg     string `json:"msg"`
}

// OrderHandler 下单的HTTP处理函数
func OrderHandler(ctx context.Context, req OrderReq) (*OrderResp, error) {
	// 从Context里拿TraceId,Go-Zero已经自动放到这里了
	traceID := logx.TraceID(ctx)
	// 打日志,确认TraceId存在
	logx.WithContext(ctx).Infof("收到下单请求,用户ID:%d,商品ID:%d,TraceId:%s", req.UserID, req.GoodsID, traceID)

	// 这里会调用库存RPC服务,后面再写
	// err := callInventoryRPC(ctx, req.GoodsID)
	// if err != nil {
	// 	return nil, err
	// }

	// 这里会发消息到RabbitMQ,后面再写
	// err := sendMQMessage(ctx, req.UserID, req.GoodsID, traceID)
	// if err != nil {
	// 	return nil, err
	// }

	return &OrderResp{
		TraceID: traceID,
		Msg:     "下单成功",
	}, nil
}

// 注册HTTP路由的函数
func RegisterOrderRoutes(server *httpx.Server) {
	server.AddRoute(httpx.Route{
		Method:  http.MethodPost,
		Path:    "/order",
		Handler: httpx.Handler(OrderHandler),
	})
}

这里要注意:Go-Zero的logx.TraceID(ctx)是专门用来拿Context里的TraceId的,只要是同步请求(HTTP、RPC),这个Context都是Go-Zero生成的,自带TraceId,不会丢。

3.3 第二步:RPC调用的TraceId传播

接下来写库存服务的RPC接口,演示下单服务调用库存服务时,TraceId怎么自动传过去。

首先写库存服务的RPC定义(proto文件):

// proto/inventory/inventory.proto 库存服务的接口定义
syntax = "proto3";

package inventory;

option go_package = "./inventory";

service InventoryService {
  // 扣减库存的RPC方法
  rpc DeductStock(DeductStockReq) returns (DeductStockResp);
}

message DeductStockReq {
  int64 goods_id = 1;
}

message DeductStockResp {
  bool success = 1;
}

然后写库存服务的RPC实现:

// rpc/inventory/inventory.go 库存服务的RPC实现
package inventory

import (
	"context"
	"github.com/zeromicro/go-zero/core/logx"
)

// InventoryService 库存服务的结构体
type InventoryService struct {
}

// NewInventoryService 初始化库存服务
func NewInventoryService() *InventoryService {
	return &InventoryService{}
}

// DeductStock 扣减库存的RPC方法
func (s *InventoryService) DeductStock(ctx context.Context, req *DeductStockReq) (*DeductStockResp, error) {
	// 从Context里拿TraceId,这里的Context是下单服务传过来的,自动带TraceId
	traceID := logx.TraceID(ctx)
	// 打日志,确认TraceId和下单服务的一致
	logx.WithContext(ctx).Infof("扣减商品ID:%d的库存,TraceId:%s", req.GoodsId, traceID)

	// 这里写扣减库存的逻辑,比如查数据库、减库存
	// ...

	return &DeductStockResp{
		Success: true,
	}, nil
}

然后回到下单服务,补全调用库存RPC的代码:

// api/order/order.go 补全的调用库存RPC的代码
import (
	"context"
	"github.com/zeromicro/go-zero/core/logx"
	"github.com/zeromicro/go-zero/rest/httpx"
	"net/http"
	// 导入库存服务的客户端
	"your-project-path/rpc/inventory"
)

// 调用库存RPC的函数
func callInventoryRPC(ctx context.Context, goodsID int64) error {
	// 初始化库存服务的客户端,这里的ctx要传进去,Go-Zero会自动把TraceId带到RPC请求里
	client := inventory.NewInventoryServiceClient(httpx.GetClient(ctx))
	resp, err := client.DeductStock(ctx, &inventory.DeductStockReq{
		GoodsId: goodsID,
	})
	if err != nil {
		logx.WithContext(ctx).Errorf("调用库存RPC失败,错误:%v,TraceId:%s", err, logx.TraceID(ctx))
		return err
	}
	logx.WithContext(ctx).Infof("调用库存RPC成功,响应:%v,TraceId:%s", resp, logx.TraceID(ctx))
	return nil
}

这里的关键是:调用RPC时,一定要把当前的Context传进去,Go-Zero的客户端会自动把Context里的TraceId放到RPC请求的元数据(Metadata)里,服务端拿到后会自动塞回自己的Context里,所以库存服务里能拿到和下单服务一模一样的TraceId。

3.4 第三步:消息队列的TraceId传递

接下来是最容易丢TraceId的地方:发消息到消息队列。因为消息队列是异步的,原来的Context会在HTTP请求结束后销毁,所以我们要手动把TraceId取出来,放到消息的元数据里,消费的时候再塞回去。

首先补全下单服务发消息到RabbitMQ的代码:

// api/order/order.go 补全的发消息到RabbitMQ的代码
import (
	"context"
	"github.com/zeromicro/go-zero/core/logx"
	"github.com/zeromicro/go-zero/rest/httpx"
	"net/http"
	"github.com/streadway/amqp" // RabbitMQ的Go客户端
)

// 发消息到RabbitMQ的函数
func sendMQMessage(ctx context.Context, userID int64, goodsID int64, traceID string) error {
	// 初始化RabbitMQ连接,这里省略了初始化逻辑,实际项目要复用连接
	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
	if err != nil {
		logx.WithContext(ctx).Errorf("连接RabbitMQ失败,错误:%v,TraceId:%s", err, traceID)
		return err
	}
	defer conn.Close()

	ch, err := conn.Channel()
	if err != nil {
		logx.WithContext(ctx).Errorf("打开RabbitMQ通道失败,错误:%v,TraceId:%s", err, traceID)
		return err
	}
	defer ch.Close()

	// 声明队列,实际项目要提前声明
	queueName := "order_push_logistics"
	_, err = ch.QueueDeclare(
		queueName,
		true,  // 持久化
		false, // 自动删除
		false, // 排他
		false, // 不等待
		nil,   // 参数
	)
	if err != nil {
		logx.WithContext(ctx).Errorf("声明RabbitMQ队列失败,错误:%v,TraceId:%s", err, traceID)
		return err
	}

	// 构造消息,把TraceId放到消息的元数据(Headers)里
	msg := amqp.Publishing{
		ContentType: "application/json",
		Headers: amqp.Table{
			"trace_id": traceID, // 把TraceId放到Headers里,一定要传
		},
		Body: []byte(`{"user_id":` + string(userID) + `,"goods_id":` + string(goodsID) + `}`),
	}

	// 发消息到队列
	err = ch.Publish(
		"",          // 交换机
		queueName,   // 队列名
		false,       // 强制
		false,       // 立即
		msg,         // 消息
	)
	if err != nil {
		logx.WithContext(ctx).Errorf("发消息到RabbitMQ失败,错误:%v,TraceId:%s", err, traceID)
		return err
	}
	logx.WithContext(ctx).Infof("发消息到RabbitMQ成功,TraceId:%s", traceID)
	return nil
}

这里的关键是:发消息时,一定要把TraceId放到消息的Headers(元数据)里,不能只放到消息体里——因为如果消息体解析失败,我们就拿不到TraceId了,放到Headers里更可靠。

接下来写消费RabbitMQ消息的服务,演示怎么把TraceId塞回Context里:

// mq/consumer/logistics_consumer.go 物流推送的MQ消费服务
package consumer

import (
	"context"
	"github.com/zeromicro/go-zero/core/logx"
	"github.com/streadway/amqp"
)

// 初始化MQ消费服务
func InitLogisticsConsumer() {
	// 初始化RabbitMQ连接
	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
	if err != nil {
		logx.Errorf("连接RabbitMQ失败,错误:%v", err)
		return
	}
	defer conn.Close()

	ch, err := conn.Channel()
	if err != nil {
		logx.Errorf("打开RabbitMQ通道失败,错误:%v", err)
		return
	}
	defer ch.Close()

	// 声明队列,和发消息的队列一致
	queueName := "order_push_logistics"
	q, err := ch.QueueDeclare(
		queueName,
		true,
		false,
		false,
		false,
		nil,
	)
	if err != nil {
		logx.Errorf("声明RabbitMQ队列失败,错误:%v", err)
		return
	}

	// 消费消息
	msgs, err := ch.Consume(
		q.Name, // 队列名
		"",      // 消费者标签
		false,   // 自动确认
		false,   // 排他
		false,   // 不等待
		false,   // 参数
		nil,     // 其他
	)
	if err != nil {
		logx.Errorf("注册RabbitMQ消费者失败,错误:%v", err)
		return
	}

	// 处理消息的协程
	go func() {
		for d := range msgs {
			// 从消息的Headers里拿TraceId
			traceID, ok := d.Headers["trace_id"].(string)
			if !ok {
				// 如果拿不到TraceId,就生成一个新的,避免日志没有TraceId
				traceID = logx.NewTraceID()
				logx.Errorf("消息没有TraceId,生成新的TraceId:%s", traceID)
			}

			// 用Go-Zero的logx.ContextWithTraceID函数,把TraceId塞到新的Context里
			ctx := logx.ContextWithTraceID(context.Background(), traceID)

			// 打日志,确认TraceId和下单服务的一致
			logx.WithContext(ctx).Infof("收到MQ消息,消息体:%s,TraceId:%s", string(d.Body), traceID)

			// 这里写物流推送的逻辑,比如调用第三方物流服务
			// ...

			// 确认消息处理成功
			err := d.Ack(false)
			if err != nil {
				logx.WithContext(ctx).Errorf("确认MQ消息失败,错误:%v,TraceId:%s", err, traceID)
			}
		}
	}()

	// 阻塞主线程,让消费服务一直运行
	select {}
}

这里的关键是:消费消息时,从Headers里拿到TraceId后,一定要用logx.ContextWithTraceID把它塞到新的Context里,后面所有的操作(比如调用第三方服务、打日志)都用这个Context,就能保证TraceId全程不丢。

四、方案分析:优缺点、应用场景、注意事项

4.1 应用场景

这个方案特别适合以下场景:

  1. 分布式系统里的全链路追踪,需要把同步请求(HTTP、RPC)和异步任务(MQ消费、定时任务)串起来;
  2. 异步任务出问题时,需要快速定位触发它的上游请求;
  3. 服务之间的调用链需要可视化(比如用Jaeger、Zipkin展示)。

4.2 技术优缺点

优点:

  1. 完全复用Go-Zero的原生能力,不需要额外写复杂的逻辑,成本低;
  2. TraceId的传递是自动的(HTTP、RPC),只有MQ需要手动处理,不容易出错;
  3. 即使MQ消息的TraceId丢了,也能生成新的,不会导致整个链路断;
  4. 日志里的TraceId统一,方便排查问题。

缺点:

  1. MQ的TraceId传递需要手动处理,如果忘了放到Headers里,还是会丢;
  2. 如果用的是其他消息队列(比如Kafka),需要自己实现Headers的传递,逻辑类似但要调整;
  3. 定时任务这种没有上游请求的场景,需要自己生成TraceId,不能完全自动。

4.3 注意事项

  1. 所有的Context传递都要传Go-Zero生成的,不能自己随便创建一个新的Context传进去,否则会丢TraceId;
  2. MQ的TraceId一定要放到Headers里,不要只放到消息体里,避免消息体解析失败时拿不到;
  3. 消费MQ消息时,一定要用logx.ContextWithTraceID把TraceId塞到Context里,后面的所有操作都用这个Context;
  4. 日志组件一定要用Go-Zero的logx.WithContext(ctx)打日志,不能直接用fmt或者其他日志组件,否则打出来的日志没有TraceId。

五、方案验证:怎么确认TraceId全程不丢

写完代码后,我们可以用这个方法验证:

  1. 启动下单服务、库存服务、MQ消费服务;
  2. 用Postman发一个POST请求到/order,参数是{"goods_id":123,"user_id":456}
  3. 看三个服务的日志:
    • 下单服务的日志里会有一个TraceId,比如abc123
    • 库存服务的日志里也会有abc123
    • MQ消费服务的日志里也会有abc123
  4. 如果三个日志里的TraceId一致,说明全程传递成功了。

如果出现MQ消费服务的TraceId和下单服务的不一致,检查两个地方:

  1. 发消息时有没有把TraceId放到Headers里;
  2. 消费消息时有没有把TraceId塞到Context里。

六、总结

TraceId丢失是分布式系统里很常见的问题,尤其是异步任务场景,一旦丢了,排查问题的成本会特别高。用Go-Zero的Context传播机制解决这个问题,核心思路就是:从请求进来的那一刻开始,就把TraceId放到Context里,然后所有的调用(HTTP、RPC、MQ)都带着这个Context,全程不丢。

这个方案的好处是简单、可靠,不需要额外的组件,只要掌握了Go-Zero的Context用法,就能轻松实现。而且即使出问题,排查的成本也很低,因为所有的日志都有统一的TraceId,能快速定位到问题的环节。

最后再提醒大家一句:做分布式系统,全链路追踪的核心就是TraceId的传递,一定要保证每个环节都有,不然出问题的时候真的会很头疼。