返回首页

Golang项目实践:4.4 微服务架构实战之构建可观测、可扩展的服务体系

在单体应用发展到一定规模后,团队往往会面临几个典型问题:模块边界不清晰、发布相互影响、扩容粒度粗、故障容易级联扩散。微服务架构的价值,就在于把复杂系统拆成一组职责明确、可独立演进的服务,并通过标准化通信、治理与观测手段,把系统重新组织起来。

本文以一个“用户服务 + 订单服务 + API 网关”的小型案例为主线,完整讲解 Golang 微服务落地中的关键技术点,包括:

  • gRPC 服务开发:Proto 定义、服务实现、拦截器
  • 服务注册与发现(Consul / etcd)
  • 负载均衡与熔断降级
  • 链路追踪(OpenTelemetry)
  • 微服务间通信:同步 vs 异步
  • API 网关设计

文章中的示例代码尽量保持“拿来就能跑”的风格。为了保证可读性,本文会把不同文件分块展示,你可以按目录结构保存后直接运行。

一、项目目标与整体架构

我们实现一个简化版的微服务系统:

  • user-service:提供用户查询能力
  • order-service:创建订单时同步调用用户服务校验用户
  • event-consumer:异步消费订单创建事件
  • api-gateway:统一对外提供 HTTP API
  • registry:采用 etcd 或 Consul 做服务注册与发现
  • tracing:使用 OpenTelemetry 输出链路追踪数据
  • resilience:客户端侧实现轮询负载均衡与熔断降级

架构关系如下:

  1. 浏览器 / 前端请求到达 API Gateway。
  2. Gateway 通过服务发现拿到 order-service 地址。
  3. order-service 在创建订单时同步调用 user-service
  4. 订单创建成功后,再异步发布一条事件给消息系统。
  5. event-consumer 订阅事件,执行通知、积分、审计等后续逻辑。
  6. 所有请求链路都被 OpenTelemetry 采集。

二、项目目录

建议目录结构如下:

microservices-demo/
├── go.mod
├── docker-compose.yml
├── proto/
│   ├── user.proto
│   └── order.proto
├── cmd/
│   ├── user-service/
│   │   └── main.go
│   ├── order-service/
│   │   └── main.go
│   ├── event-consumer/
│   │   └── main.go
│   └── api-gateway/
│       └── main.go
├── internal/
│   ├── config/
│   │   └── config.go
│   ├── discovery/
│   │   ├── registry.go
│   │   ├── etcd_registry.go
│   │   └── consul_registry.go
│   ├── lb/
│   │   └── picker.go
│   ├── resilience/
│   │   └── breaker.go
│   ├── transport/
│   │   └── grpc_interceptor.go
│   ├── tracing/
│   │   └── otel.go
│   └── event/
│       └── nats.go
└── gen/
    └── pb/

三、准备运行环境

先初始化 go.mod

module microservices-demo

go 1.22

require (
	github.com/gin-gonic/gin v1.10.0
	github.com/nats-io/nats.go v1.39.1
	github.com/hashicorp/consul/api v1.31.0
	github.com/sony/gobreaker v1.0.0
	go.etcd.io/etcd/client/v3 v3.5.18
	go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.57.0
	go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.57.0
	go.opentelemetry.io/otel v1.32.0
	go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.32.0
	go.opentelemetry.io/otel/sdk v1.32.0
	go.opentelemetry.io/otel/trace v1.32.0
	google.golang.org/grpc v1.68.1
	google.golang.org/protobuf v1.35.2
)

为了方便本地演示,可以直接用 Docker 启动 etcd、Consul、NATS 和 Jaeger。

version: "3.9"
services:
  etcd:
    image: quay.io/coreos/etcd:v3.5.18
    command:
      - /usr/local/bin/etcd
      - --name=node1
      - --advertise-client-urls=http://0.0.0.0:2379
      - --listen-client-urls=http://0.0.0.0:2379
    ports:
      - "2379:2379"

  consul:
    image: hashicorp/consul:1.20
    command: agent -server -bootstrap -ui -client=0.0.0.0
    ports:
      - "8500:8500"

  nats:
    image: nats:2.10
    ports:
      - "4222:4222"

  jaeger:
    image: jaegertracing/all-in-one:1.61
    environment:
      - COLLECTOR_OTLP_ENABLED=true
    ports:
      - "16686:16686"
      - "4317:4317"

启动命令如下:

docker compose up -d

如果你需要生成 Proto 代码,请先安装编译工具:

go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest

四、gRPC 服务开发:Proto 定义、服务实现、拦截器

4.1 Proto 定义

先定义用户服务。

syntax = "proto3";

package user;

option go_package = "microservices-demo/gen/pb;pb";

service UserService {
  rpc GetUser(GetUserRequest) returns (GetUserResponse);
}

message GetUserRequest {
  int64 user_id = 1;
}

message GetUserResponse {
  int64 user_id = 1;
  string name = 2;
  bool active = 3;
}

再定义订单服务。

syntax = "proto3";

package order;

option go_package = "microservices-demo/gen/pb;pb";

service OrderService {
  rpc CreateOrder(CreateOrderRequest) returns (CreateOrderResponse);
}

message CreateOrderRequest {
  int64 user_id = 1;
  string product = 2;
  int32 amount = 3;
}

message CreateOrderResponse {
  string order_id = 1;
  string status = 2;
}

生成代码命令如下:

protoc --go_out=. --go-grpc_out=. proto/user.proto
protoc --go_out=. --go-grpc_out=. proto/order.proto

4.2 通用配置

为了让各个服务启动方式一致,我们定义一个配置结构。

package config

type Config struct {
	ServiceName      string
	ServiceAddr      string
	RegistryType     string
	EtcdEndpoints    []string
	ConsulAddr       string
	NATSURL          string
	JaegerEndpoint   string
	GatewayHTTPAddr  string
}

4.3 gRPC 拦截器

在微服务里,拦截器的价值非常大。它和 HTTP 中间件类似,适合统一处理日志、超时、鉴权、Tracing、恢复 panic 等横切逻辑。

下面实现一个基础的一元拦截器。

package transport

import (
	"context"
	"log"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/grpc/codes"
	"google.golang.org/grpc/status"
)

func UnaryServerInterceptor() grpc.UnaryServerInterceptor {
	return func(
		ctx context.Context,
		req any,
		info *grpc.UnaryServerInfo,
		handler grpc.UnaryHandler,
	) (resp any, err error) {
		start := time.Now()

		defer func() {
			cost := time.Since(start)
			if err != nil {
				log.Printf("[grpc] method=%s cost=%s err=%v", info.FullMethod, cost, err)
				return
			}
			log.Printf("[grpc] method=%s cost=%s", info.FullMethod, cost)
		}()

		resp, err = handler(ctx, req)
		if err != nil {
			return nil, status.Errorf(codes.Internal, "internal error: %v", err)
		}
		return resp, nil
	}
}

4.4 UserService 实现

下面是用户服务实现。为了突出重点,示例里直接用内存数据模拟数据库。

package main

import (
	"context"
	"log"
	"net"
	"os"
	"os/signal"
	"syscall"
	"time"

	"microservices-demo/gen/pb"
	"microservices-demo/internal/discovery"
	"microservices-demo/internal/tracing"
	"microservices-demo/internal/transport"

	"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
	"google.golang.org/grpc"
)

type userServer struct {
	pb.UnimplementedUserServiceServer
}

func (s *userServer) GetUser(ctx context.Context, req *pb.GetUserRequest) (*pb.GetUserResponse, error) {
	if req.UserId == 0 {
		return &pb.GetUserResponse{
			UserId: 0,
			Name:   "",
			Active: false,
		}, nil
	}

	return &pb.GetUserResponse{
		UserId: req.UserId,
		Name:   "张三",
		Active: true,
	}, nil
}

func main() {
	cfg := struct {
		Name          string
		Addr          string
		RegistryType  string
		EtcdEndpoints []string
		ConsulAddr    string
		Jaeger        string
	}{
		Name:          "user-service",
		Addr:          ":50051",
		RegistryType:  envOrDefault("REGISTRY_TYPE", "etcd"),
		EtcdEndpoints: []string{"localhost:2379"},
		ConsulAddr:    envOrDefault("CONSUL_ADDR", "127.0.0.1:8500"),
		Jaeger:        envOrDefault("JAEGER_ENDPOINT", "localhost:4317"),
	}

	shutdownTracer, err := tracing.InitTracer(cfg.Name, cfg.Jaeger)
	if err != nil {
		log.Fatalf("init tracer failed: %v", err)
	}
	defer shutdownTracer(context.Background())

	lis, err := net.Listen("tcp", cfg.Addr)
	if err != nil {
		log.Fatalf("listen failed: %v", err)
	}

	server := grpc.NewServer(
		grpc.ChainUnaryInterceptor(
			transport.UnaryServerInterceptor(),
		),
		grpc.StatsHandler(otelgrpc.NewServerHandler()),
	)
	pb.RegisterUserServiceServer(server, &userServer{})

	registry, err := discovery.NewRegistry(
		cfg.RegistryType,
		cfg.EtcdEndpoints,
		cfg.ConsulAddr,
	)
	if err != nil {
		log.Fatalf("create registry failed: %v", err)
	}
	defer registry.Close()

	ctx := context.Background()
	if err := registry.Register(ctx, cfg.Name, "127.0.0.1:50051", 10*time.Second); err != nil {
		log.Fatalf("register service failed: %v", err)
	}
	defer registry.Deregister(context.Background(), cfg.Name, "127.0.0.1:50051")

	go func() {
		log.Printf("%s listening at %s", cfg.Name, cfg.Addr)
		if err := server.Serve(lis); err != nil {
			log.Fatalf("serve failed: %v", err)
		}
	}()

	quit := make(chan os.Signal, 1)
	signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
	<-quit

	server.GracefulStop()
}

func envOrDefault(key, fallback string) string {
	if v := os.Getenv(key); v != "" {
		return v
	}
	return fallback
}

这里有几个实践要点:

  • 服务实现关注业务逻辑本身。
  • 链路追踪、日志记录等通用逻辑交给拦截器。
  • 服务启动后立即注册到注册中心。
  • 服务退出前要反注册,避免脏地址残留。

五、服务注册与发现(Consul / etcd)

微服务拆开之后,服务地址不应该再写死在配置文件里。因为实例可能随时扩容、缩容、重启、迁移,调用方必须通过注册中心动态获取可用节点。

5.1 抽象统一接口

先定义统一的注册中心接口。

package discovery

import (
	"context"
	"time"
)

type Registry interface {
	Register(ctx context.Context, serviceName, addr string, ttl time.Duration) error
	Deregister(ctx context.Context, serviceName, addr string) error
	Discover(ctx context.Context, serviceName string) ([]string, error)
	Close() error
}

func NewRegistry(kind string, etcdEndpoints []string, consulAddr string) (Registry, error) {
	switch kind {
	case "consul":
		return NewConsulRegistry(consulAddr)
	default:
		return NewEtcdRegistry(etcdEndpoints)
	}
}

5.2 etcd 注册与发现实现

package discovery

import (
	"context"
	"fmt"
	"path"
	"strings"
	"sync"
	"time"

	clientv3 "go.etcd.io/etcd/client/v3"
)

type EtcdRegistry struct {
	client *clientv3.Client
	mu     sync.Mutex
	leases map[string]clientv3.LeaseID
}

func NewEtcdRegistry(endpoints []string) (*EtcdRegistry, error) {
	cli, err := clientv3.New(clientv3.Config{Endpoints: endpoints, DialTimeout: 5 * time.Second})
	if err != nil {
		return nil, err
	}
	return &EtcdRegistry{client: cli, leases: make(map[string]clientv3.LeaseID)}, nil
}

func (r *EtcdRegistry) Register(ctx context.Context, serviceName, addr string, ttl time.Duration) error {
	leaseResp, err := r.client.Grant(ctx, int64(ttl.Seconds()))
	if err != nil {
		return err
	}

	key := path.Join("/services", serviceName, addr)
	if _, err = r.client.Put(ctx, key, addr, clientv3.WithLease(leaseResp.ID)); err != nil {
		return err
	}

	keepAliveCh, err := r.client.KeepAlive(ctx, leaseResp.ID)
	if err != nil {
		return err
	}

	r.mu.Lock()
	r.leases[key] = leaseResp.ID
	r.mu.Unlock()

	go func() {
		for range keepAliveCh {
		}
	}()

	return nil
}

func (r *EtcdRegistry) Deregister(ctx context.Context, serviceName, addr string) error {
	key := path.Join("/services", serviceName, addr)
	_, err := r.client.Delete(ctx, key)
	return err
}

func (r *EtcdRegistry) Discover(ctx context.Context, serviceName string) ([]string, error) {
	prefix := path.Join("/services", serviceName)
	resp, err := r.client.Get(ctx, prefix, clientv3.WithPrefix())
	if err != nil {
		return nil, err
	}

	addrs := make([]string, 0, len(resp.Kvs))
	for _, kv := range resp.Kvs {
		addr := strings.TrimSpace(string(kv.Value))
		if addr != "" {
			addrs = append(addrs, addr)
		}
	}
	if len(addrs) == 0 {
		return nil, fmt.Errorf("no instance found for service %s", serviceName)
	}
	return addrs, nil
}

func (r *EtcdRegistry) Close() error {
	return r.client.Close()
}

5.3 Consul 注册与发现实现

如果你的团队已经在使用 Consul,也可以平滑替换。下面是 Consul 版实现。

package discovery

import (
	"context"
	"fmt"
	"strings"
	"time"

	consulapi "github.com/hashicorp/consul/api"
)

type ConsulRegistry struct {
	client *consulapi.Client
}

func NewConsulRegistry(addr string) (*ConsulRegistry, error) {
	cfg := consulapi.DefaultConfig()
	cfg.Address = addr
	cli, err := consulapi.NewClient(cfg)
	if err != nil {
		return nil, err
	}
	return &ConsulRegistry{client: cli}, nil
}

func (r *ConsulRegistry) Register(ctx context.Context, serviceName, addr string, ttl time.Duration) error {
	hostPort := strings.Split(addr, ":")
	if len(hostPort) != 2 {
		return fmt.Errorf("invalid addr: %s", addr)
	}

	registration := &consulapi.AgentServiceRegistration{
		ID:      serviceName + "-" + strings.ReplaceAll(addr, ":", "-"),
		Name:    serviceName,
		Address: hostPort[0],
		Port:    mustAtoi(hostPort[1]),
		Check: &consulapi.AgentServiceCheck{
			TCP:                            addr,
			Interval:                       "10s",
			Timeout:                        "3s",
			DeregisterCriticalServiceAfter: "30s",
		},
	}
	return r.client.Agent().ServiceRegister(registration)
}

func (r *ConsulRegistry) Deregister(ctx context.Context, serviceName, addr string) error {
	serviceID := serviceName + "-" + strings.ReplaceAll(addr, ":", "-")
	return r.client.Agent().ServiceDeregister(serviceID)
}

func (r *ConsulRegistry) Discover(ctx context.Context, serviceName string) ([]string, error) {
	services, _, err := r.client.Health().Service(serviceName, "", true, nil)
	if err != nil {
		return nil, err
	}

	addrs := make([]string, 0, len(services))
	for _, svc := range services {
		if svc.Service == nil {
			continue
		}
		addrs = append(addrs, fmt.Sprintf("%s:%d", svc.Service.Address, svc.Service.Port))
	}
	if len(addrs) == 0 {
		return nil, fmt.Errorf("no instance found for service %s", serviceName)
	}
	return addrs, nil
}

func (r *ConsulRegistry) Close() error {
	return nil
}

func mustAtoi(s string) int {
	var n int
	fmt.Sscanf(s, "%d", &n)
	return n
}

5.4 Consul 与 etcd 的选择建议

两者都能胜任注册中心角色,但侧重点略有不同。

维度 etcd Consul
核心定位 强一致 KV 存储 服务发现与服务治理
一致性模型 Raft,强一致 也支持一致性,但偏服务治理场景
使用习惯 Kubernetes 生态常见 传统微服务体系常见
健康检查 需要自行组合实现 内置更成熟
UI 与运营友好度 相对简洁 更直观

如果你偏向云原生、Kubernetes、控制面能力,通常会更容易接受 etcd;如果你更看重服务治理、健康检查与现成运维体验,Consul 也是很好的选择。

六、负载均衡与熔断降级

服务发现解决了“去哪找服务”的问题,负载均衡和熔断解决的是“怎么调用更稳”。

  • 负载均衡:把流量合理分配到多个实例
  • 熔断:当下游持续异常时,快速失败,保护自身线程与连接资源
  • 降级:依赖不可用时,返回兜底结果或简化结果

6.1 简单轮询负载均衡

为了让示例容易理解,我们先实现一个客户端侧轮询选择器。

package lb

import "sync/atomic"

type RoundRobinPicker struct {
	idx uint64
}

func (p *RoundRobinPicker) Pick(addrs []string) string {
	if len(addrs) == 0 {
		return ""
	}
	n := atomic.AddUint64(&p.idx, 1)
	return addrs[(int(n)-1)%len(addrs)]
}

6.2 熔断器封装

这里使用 gobreaker 实现熔断。当连续失败超过阈值时,熔断器会短时间打开,后续请求直接失败,避免把问题放大。

package resilience

import (
	"time"

	"github.com/sony/gobreaker"
)

func NewBreaker(name string) *gobreaker.CircuitBreaker {
	st := gobreaker.Settings{
		Name:        name,
		MaxRequests: 5,
		Interval:    30 * time.Second,
		Timeout:     10 * time.Second,
		ReadyToTrip: func(counts gobreaker.Counts) bool {
			return counts.ConsecutiveFailures >= 3
		},
	}
	return gobreaker.NewCircuitBreaker(st)
}

6.3 OrderService:同步调用 UserService

订单服务在创建订单前,需要先同步调用用户服务,确认用户存在且状态正常。

package main

import (
	"context"
	"fmt"
	"log"
	"net"
	"os"
	"os/signal"
	"syscall"
	"time"

	"microservices-demo/gen/pb"
	"microservices-demo/internal/discovery"
	"microservices-demo/internal/event"
	"microservices-demo/internal/lb"
	"microservices-demo/internal/resilience"
	"microservices-demo/internal/tracing"
	"microservices-demo/internal/transport"

	"github.com/sony/gobreaker"
	"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials/insecure"
)

type orderServer struct {
	pb.UnimplementedOrderServiceServer
	registry discovery.Registry
	picker   *lb.RoundRobinPicker
	breaker  *gobreaker.CircuitBreaker
	publisher *event.NATSPublisher
}

func (s *orderServer) CreateOrder(ctx context.Context, req *pb.CreateOrderRequest) (*pb.CreateOrderResponse, error) {
	addrs, err := s.registry.Discover(ctx, "user-service")
	if err != nil {
		return &pb.CreateOrderResponse{OrderId: "", Status: "degraded: user service unavailable"}, nil
	}

	addr := s.picker.Pick(addrs)
	if addr == "" {
		return &pb.CreateOrderResponse{OrderId: "", Status: "degraded: no instance available"}, nil
	}

	result, err := s.breaker.Execute(func() (any, error) {
		conn, err := grpc.DialContext(
			ctx,
			addr,
			grpc.WithTransportCredentials(insecure.NewCredentials()),
			grpc.WithStatsHandler(otelgrpc.NewClientHandler()),
		)
		if err != nil {
			return nil, err
		}
		defer conn.Close()

		client := pb.NewUserServiceClient(conn)
		callCtx, cancel := context.WithTimeout(ctx, 2*time.Second)
		defer cancel()

		userResp, err := client.GetUser(callCtx, &pb.GetUserRequest{UserId: req.UserId})
		if err != nil {
			return nil, err
		}
		if !userResp.Active {
			return nil, fmt.Errorf("user %d is inactive", req.UserId)
		}
		return userResp, nil
	})
	if err != nil {
		return &pb.CreateOrderResponse{OrderId: "", Status: "fallback: create order later"}, nil
	}

	_ = result.(*pb.GetUserResponse)
	orderID := fmt.Sprintf("ORD-%d", time.Now().UnixNano())

	if err := s.publisher.PublishOrderCreated(ctx, orderID, req.UserId, req.Product, req.Amount); err != nil {
		log.Printf("publish event failed: %v", err)
	}

	return &pb.CreateOrderResponse{OrderId: orderID, Status: "created"}, nil
}

func main() {
	cfg := struct {
		Name          string
		Addr          string
		RegistryType  string
		EtcdEndpoints []string
		ConsulAddr    string
		NATSURL       string
		Jaeger        string
	}{
		Name:          "order-service",
		Addr:          ":50052",
		RegistryType:  envOrDefault("REGISTRY_TYPE", "etcd"),
		EtcdEndpoints: []string{"localhost:2379"},
		ConsulAddr:    envOrDefault("CONSUL_ADDR", "127.0.0.1:8500"),
		NATSURL:       envOrDefault("NATS_URL", "nats://127.0.0.1:4222"),
		Jaeger:        envOrDefault("JAEGER_ENDPOINT", "localhost:4317"),
	}

	shutdownTracer, err := tracing.InitTracer(cfg.Name, cfg.Jaeger)
	if err != nil {
		log.Fatalf("init tracer failed: %v", err)
	}
	defer shutdownTracer(context.Background())

	registry, err := discovery.NewRegistry(cfg.RegistryType, cfg.EtcdEndpoints, cfg.ConsulAddr)
	if err != nil {
		log.Fatalf("create registry failed: %v", err)
	}
	defer registry.Close()

	publisher, err := event.NewNATSPublisher(cfg.NATSURL)
	if err != nil {
		log.Fatalf("connect nats failed: %v", err)
	}
	defer publisher.Close()

	lis, err := net.Listen("tcp", cfg.Addr)
	if err != nil {
		log.Fatalf("listen failed: %v", err)
	}

	server := grpc.NewServer(
		grpc.ChainUnaryInterceptor(
			transport.UnaryServerInterceptor(),
		),
		grpc.StatsHandler(otelgrpc.NewServerHandler()),
	)

	pb.RegisterOrderServiceServer(server, &orderServer{
		registry: registry,
		picker:   &lb.RoundRobinPicker{},
		breaker:  resilience.NewBreaker("user-service"),
		publisher: publisher,
	})

	ctx := context.Background()
	if err := registry.Register(ctx, cfg.Name, "127.0.0.1:50052", 10*time.Second); err != nil {
		log.Fatalf("register service failed: %v", err)
	}
	defer registry.Deregister(context.Background(), cfg.Name, "127.0.0.1:50052")

	go func() {
		log.Printf("%s listening at %s", cfg.Name, cfg.Addr)
		if err := server.Serve(lis); err != nil {
			log.Fatalf("serve failed: %v", err)
		}
	}()

	quit := make(chan os.Signal, 1)
	signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
	<-quit
	server.GracefulStop()
}

func envOrDefault(key, fallback string) string {
	if v := os.Getenv(key); v != "" {
		return v
	}
	return fallback
}

上面这段代码体现了三个核心点:

  1. 先从注册中心发现 user-service 实例列表。
  2. 通过轮询策略选一个实例。
  3. 用熔断器包裹远程调用,失败时快速返回降级结果。

6.4 常见降级策略

当下游服务不稳定时,不一定都要报错退出,可以按业务场景进行降级。

场景 降级方式
用户信息非关键 返回默认昵称、匿名头像
推荐系统超时 返回热门推荐
统计类接口异常 返回缓存快照
非核心异步任务失败 记录日志,稍后重试
风险控制服务故障 视业务策略决定拒绝或人工介入

微服务治理不是“永不失败”,而是“失败时也能以可控方式继续运行”。

七、链路追踪(OpenTelemetry)

微服务一旦拆分,调用链就不再局限于单进程。一个请求可能跨越网关、订单服务、用户服务、消息消费者。没有链路追踪时,定位问题会非常痛苦。

OpenTelemetry 是当前主流的可观测性标准。它解决的关键问题包括:

  • 一个请求经过了哪些服务
  • 哪一跳最慢
  • 哪个服务报错
  • TraceID 能否贯穿日志、指标和告警

7.1 初始化 Tracer

package tracing

import (
	"context"
	"fmt"
	"time"

	"go.opentelemetry.io/otel"
	"go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc"
	"go.opentelemetry.io/otel/propagation"
	"go.opentelemetry.io/otel/sdk/resource"
	tracesdk "go.opentelemetry.io/otel/sdk/trace"
	semconv "go.opentelemetry.io/otel/semconv/v1.26.0"
	"google.golang.org/grpc/credentials/insecure"
)

func InitTracer(serviceName, endpoint string) (func(context.Context) error, error) {
	ctx := context.Background()

	exporter, err := otlptracegrpc.New(
		ctx,
		otlptracegrpc.WithEndpoint(endpoint),
		otlptracegrpc.WithTransportCredentials(insecure.NewCredentials()),
	)
	if err != nil {
		return nil, fmt.Errorf("create otlp exporter failed: %w", err)
	}

	res, err := resource.New(
		ctx,
		resource.WithAttributes(
			semconv.ServiceName(serviceName),
		),
	)
	if err != nil {
		return nil, fmt.Errorf("create resource failed: %w", err)
	}

	tp := tracesdk.NewTracerProvider(
		tracesdk.WithBatcher(exporter),
		tracesdk.WithResource(res),
	)

	otel.SetTracerProvider(tp)
	otel.SetTextMapPropagator(propagation.TraceContext{})

	return func(ctx context.Context) error {
		shutdownCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
		defer cancel()
		return tp.Shutdown(shutdownCtx)
	}, nil
}

7.2 为什么拦截器 + OTel 是最佳组合

在 gRPC 中,我们已经把 grpc.StatsHandler(otelgrpc.NewServerHandler())grpc.WithStatsHandler(otelgrpc.NewClientHandler()) 接入服务端和客户端。这意味着:

  • 服务端收到请求时,自动创建 Span
  • 服务端调用下游时,自动透传上下文
  • TraceID 会贯穿多个服务
  • 后台 Jaeger 可以直接看到完整调用链

这正是微服务“可观测性内建”的典型实践。

7.3 本地验证链路

启动所有服务后,访问 Jaeger 控制台:

http://localhost:16686

然后调用 Gateway 的下单接口,应该能看到大致如下链路:

  • api-gateway 收到 HTTP 请求
  • order-service/CreateOrder
  • user-service/GetUser
  • event-consumer 消费订单事件(如果你在消息头中透传 Trace Context,也可串起来)

在生产中,建议同时做三件事:

  • TraceID 写入日志
  • 延迟与错误率做指标上报
  • 以服务名、实例名、环境标签区分数据源

八、微服务间通信:同步 vs 异步

这是微服务设计里最常见、也最容易被误用的问题。很多团队一上来就把所有调用都做成同步 RPC,结果链路越来越长、耦合越来越深、雪崩风险越来越大。

8.1 同步通信适用场景

同步通信强调“请求后立即拿结果”,常见方式就是 HTTP 或 gRPC。

适用于:

  • 用户下单前必须校验库存
  • 登录时必须验证账号密码
  • 查询页面时必须拉取实时核心数据

优点:

  • 调用链直接,便于理解
  • 结果立即返回
  • 调试成本较低

缺点:

  • 强依赖下游可用性
  • 链路变长后时延叠加明显
  • 容易级联失败

8.2 异步通信适用场景

异步通信强调“先接收请求,再延后处理后续动作”,常见方式包括 NATS、Kafka、RabbitMQ。

适用于:

  • 订单创建后发送通知
  • 用户注册后发积分
  • 审计日志、行为埋点、数据同步

优点:

  • 解耦明显
  • 高峰期可以削峰填谷
  • 下游波动不必直接影响主流程

缺点:

  • 最终一致性更复杂
  • 调试难度更高
  • 需要处理重复消费、幂等、消息丢失等问题

8.3 同步 vs 异步对比

维度 同步调用 异步消息
响应时机 立即返回结果 先返回,后处理
耦合程度 较高 较低
适用场景 核心主流程 通知、补偿、派生任务
失败影响 易级联 更容易隔离
开发难度 较低 较高

8.4 NATS 异步发布与消费

先定义事件发布器。

package event

import (
	"context"
	"encoding/json"
	"fmt"

	"github.com/nats-io/nats.go"
)

type OrderCreatedEvent struct {
	OrderID string `json:"order_id"`
	UserID  int64  `json:"user_id"`
	Product string `json:"product"`
	Amount  int32  `json:"amount"`
}

type NATSPublisher struct {
	nc *nats.Conn
}

func NewNATSPublisher(url string) (*NATSPublisher, error) {
	nc, err := nats.Connect(url)
	if err != nil {
		return nil, err
	}
	return &NATSPublisher{nc: nc}, nil
}

func (p *NATSPublisher) PublishOrderCreated(ctx context.Context, orderID string, userID int64, product string, amount int32) error {
	evt := OrderCreatedEvent{
		OrderID: orderID,
		UserID:  userID,
		Product: product,
		Amount:  amount,
	}
	data, err := json.Marshal(evt)
	if err != nil {
		return err
	}
	return p.nc.Publish("order.created", data)
}

func (p *NATSPublisher) Close() {
	if p.nc != nil {
		p.nc.Close()
	}
}

func PrettyEvent(data []byte) string {
	var evt OrderCreatedEvent
	if err := json.Unmarshal(data, &evt); err != nil {
		return string(data)
	}
	return fmt.Sprintf("order_id=%s user_id=%d product=%s amount=%d", evt.OrderID, evt.UserID, evt.Product, evt.Amount)
}

再定义消费者服务。

package main

import (
	"log"
	"os"
	"os/signal"
	"syscall"

	"microservices-demo/internal/event"

	"github.com/nats-io/nats.go"
)

func main() {
	nc, err := nats.Connect(envOrDefault("NATS_URL", "nats://127.0.0.1:4222"))
	if err != nil {
		log.Fatalf("connect nats failed: %v", err)
	}
	defer nc.Close()

	_, err = nc.Subscribe("order.created", func(msg *nats.Msg) {
		log.Printf("receive async event: %s", event.PrettyEvent(msg.Data))
		log.Printf("simulate send SMS / add points / audit log success")
	})
	if err != nil {
		log.Fatalf("subscribe failed: %v", err)
	}

	log.Println("event-consumer is running")
	quit := make(chan os.Signal, 1)
	signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
	<-quit
}

func envOrDefault(key, fallback string) string {
	if v := os.Getenv(key); v != "" {
		return v
	}
	return fallback
}

在这个模型里,订单创建成功只依赖主流程本身;短信、积分、审计等任务由异步消费者接管。这样可以显著降低主流程的响应时间与耦合度。

8.5 实战建议

一个简单但很有效的判断标准是:

  • 必须立刻拿结果的,用同步调用
  • 允许延后完成的,用异步消息

不要把所有流程都做成同步 RPC,也不要为了“架构高级”而强行异步。架构设计的核心从来不是炫技,而是贴合业务约束。

九、API 网关设计

微服务系统对外暴露接口时,通常不会让前端直接访问所有内部服务,而是由 API Gateway 统一承接入口流量。

网关常见职责包括:

  • 统一路由
  • 认证鉴权
  • 限流
  • 请求聚合
  • 灰度发布
  • 日志与 Trace 注入
  • 隐藏内部服务地址

9.1 Gateway 的设计原则

一个好的 API 网关通常遵循以下原则:

  1. 轻业务,重治理:网关不要承载复杂领域逻辑。
  2. 统一入口:前端只记住一个访问域名。
  3. 可观测性优先:日志、指标、Trace 一开始就接入。
  4. 可扩展:便于后续增加鉴权、限流、灰度、黑白名单。

9.2 Gateway 示例实现

下面我们用 Gin 构建一个轻量网关,对外暴露创建订单接口。

package main

import (
	"context"
	"log"
	"net/http"
	"os"
	"time"

	"microservices-demo/gen/pb"
	"microservices-demo/internal/discovery"
	"microservices-demo/internal/lb"
	"microservices-demo/internal/tracing"

	"github.com/gin-gonic/gin"
	"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
	"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials/insecure"
)

type createOrderReq struct {
	UserID  int64  `json:"user_id" binding:"required"`
	Product string `json:"product" binding:"required"`
	Amount  int32  `json:"amount" binding:"required"`
}

func main() {
	shutdownTracer, err := tracing.InitTracer("api-gateway", envOrDefault("JAEGER_ENDPOINT", "localhost:4317"))
	if err != nil {
		log.Fatalf("init tracer failed: %v", err)
	}
	defer shutdownTracer(context.Background())

	registry, err := discovery.NewRegistry(
		envOrDefault("REGISTRY_TYPE", "etcd"),
		[]string{"localhost:2379"},
		envOrDefault("CONSUL_ADDR", "127.0.0.1:8500"),
	)
	if err != nil {
		log.Fatalf("create registry failed: %v", err)
	}
	defer registry.Close()

	picker := &lb.RoundRobinPicker{}

	r := gin.Default()
	r.Use(func(c *gin.Context) {
		if c.GetHeader("X-Request-Id") == "" {
			c.Header("X-Request-Id", time.Now().Format("20060102150405.000"))
		}
		c.Next()
	})

	r.POST("/api/v1/orders", func(c *gin.Context) {
		var req createOrderReq
		if err := c.ShouldBindJSON(&req); err != nil {
			c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
			return
		}

		addrs, err := registry.Discover(c.Request.Context(), "order-service")
		if err != nil {
			c.JSON(http.StatusServiceUnavailable, gin.H{"error": "order service unavailable"})
			return
		}

		addr := picker.Pick(addrs)
		conn, err := grpc.DialContext(
			c.Request.Context(),
			addr,
			grpc.WithTransportCredentials(insecure.NewCredentials()),
			grpc.WithStatsHandler(otelgrpc.NewClientHandler()),
		)
		if err != nil {
			c.JSON(http.StatusServiceUnavailable, gin.H{"error": err.Error()})
			return
		}
		defer conn.Close()

		client := pb.NewOrderServiceClient(conn)
		resp, err := client.CreateOrder(c.Request.Context(), &pb.CreateOrderRequest{
			UserId:  req.UserID,
			Product: req.Product,
			Amount:  req.Amount,
		})
		if err != nil {
			c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
			return
		}

		c.JSON(http.StatusOK, gin.H{
			"order_id": resp.OrderId,
			"status":   resp.Status,
		})
	})

	handler := otelhttp.NewHandler(r, "gateway-http")
	server := &http.Server{
		Addr:    envOrDefault("GATEWAY_ADDR", ":8080"),
		Handler: handler,
	}

	log.Printf("api-gateway listening at %s", server.Addr)
	if err := server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
		log.Fatalf("gateway start failed: %v", err)
	}
}

func envOrDefault(key, fallback string) string {
	if v := os.Getenv(key); v != "" {
		return v
	}
	return fallback
}

9.3 网关为什么不能写太多业务

这是很多团队在实战中会踩的坑。刚开始觉得“顺手在网关里做一点逻辑没关系”,后来慢慢变成:

  • 聚合逻辑在网关
  • 权限逻辑在网关
  • 数据拼装在网关
  • 业务判断也在网关

最终网关会变成另一个“巨型单体入口”。

正确做法是:

  • 网关负责入口治理
  • 领域逻辑留在领域服务内部
  • 聚合逻辑只做薄封装,不侵入领域规则

十、如何把整个示例跑起来

10.1 启动基础设施

docker compose up -d

10.2 生成 Proto 代码

protoc --go_out=. --go-grpc_out=. proto/user.proto
protoc --go_out=. --go-grpc_out=. proto/order.proto

10.3 启动用户服务

go run cmd/user-service/main.go

10.4 启动订单服务

go run cmd/order-service/main.go

10.5 启动异步消费者

go run cmd/event-consumer/main.go

10.6 启动 API 网关

go run cmd/api-gateway/main.go

10.7 发起请求测试

curl -X POST http://localhost:8080/api/v1/orders \
  -H 'Content-Type: application/json' \
  -d '{"user_id":1,"product":"Go 微服务实战课程","amount":2}'

如果运行顺利,你会看到:

  • Gateway 返回订单号与状态
  • order-service 成功调用 user-service
  • event-consumer 收到订单创建事件
  • Jaeger 中出现完整链路

十一、从示例看微服务落地的核心方法论

写到这里,你会发现微服务实战并不是简单地“把服务拆开”。真正关键的是下面几件事一起成立:

  1. 通信标准化:服务之间最好统一使用 gRPC 或稳定的 HTTP 契约。
  2. 地址动态化:通过 Consul / etcd 做服务注册与发现,避免写死地址。
  3. 调用可治理:要有负载均衡、超时、熔断和降级。
  4. 观测可追踪:OpenTelemetry 让跨服务问题可定位。
  5. 同步异步边界清晰:主流程用同步,派生流程用异步。
  6. 统一入口治理:通过 API 网关对外暴露能力,而不是让前端直连内部服务。

如果只做“服务拆分”,不做治理、观测和容错,系统复杂度通常会上升;只有把这些配套能力一起建设起来,微服务架构才真正能支撑业务演进。

十二、总结

本文通过一个可运行的 Golang 示例,把微服务架构中的几个核心环节串了起来:

  • 用 Proto 定义服务契约,并基于 gRPC 实现服务
  • 使用拦截器统一接入日志与链路能力
  • 通过 etcd / Consul 实现服务注册与发现
  • 在客户端实现轮询负载均衡与熔断降级
  • 借助 OpenTelemetry 构建调用链追踪
  • 按业务约束选择同步或异步通信方式
  • 使用 API Gateway 统一承接外部请求

对于 Golang 而言,微服务并不只是“框架选型”问题,更是工程化能力建设问题。只要把服务契约、治理策略、可观测性与边界设计做扎实,即使系统继续扩张,也能保持较好的演进能力与稳定性。


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

Golang项目实践:4.2 Web 服务开发实战

下一篇

Golang项目实践:4.5 消息队列与异步处理