在单体应用发展到一定规模后,团队往往会面临几个典型问题:模块边界不清晰、发布相互影响、扩容粒度粗、故障容易级联扩散。微服务架构的价值,就在于把复杂系统拆成一组职责明确、可独立演进的服务,并通过标准化通信、治理与观测手段,把系统重新组织起来。
本文以一个“用户服务 + 订单服务 + API 网关”的小型案例为主线,完整讲解 Golang 微服务落地中的关键技术点,包括:
- gRPC 服务开发:Proto 定义、服务实现、拦截器
- 服务注册与发现(Consul / etcd)
- 负载均衡与熔断降级
- 链路追踪(OpenTelemetry)
- 微服务间通信:同步 vs 异步
- API 网关设计
文章中的示例代码尽量保持“拿来就能跑”的风格。为了保证可读性,本文会把不同文件分块展示,你可以按目录结构保存后直接运行。
一、项目目标与整体架构
我们实现一个简化版的微服务系统:
user-service:提供用户查询能力order-service:创建订单时同步调用用户服务校验用户event-consumer:异步消费订单创建事件api-gateway:统一对外提供 HTTP APIregistry:采用 etcd 或 Consul 做服务注册与发现tracing:使用 OpenTelemetry 输出链路追踪数据resilience:客户端侧实现轮询负载均衡与熔断降级
架构关系如下:
- 浏览器 / 前端请求到达 API Gateway。
- Gateway 通过服务发现拿到
order-service地址。 order-service在创建订单时同步调用user-service。- 订单创建成功后,再异步发布一条事件给消息系统。
event-consumer订阅事件,执行通知、积分、审计等后续逻辑。- 所有请求链路都被 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
}
上面这段代码体现了三个核心点:
- 先从注册中心发现
user-service实例列表。 - 通过轮询策略选一个实例。
- 用熔断器包裹远程调用,失败时快速返回降级结果。
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/CreateOrderuser-service/GetUserevent-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 网关通常遵循以下原则:
- 轻业务,重治理:网关不要承载复杂领域逻辑。
- 统一入口:前端只记住一个访问域名。
- 可观测性优先:日志、指标、Trace 一开始就接入。
- 可扩展:便于后续增加鉴权、限流、灰度、黑白名单。
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-serviceevent-consumer收到订单创建事件- Jaeger 中出现完整链路
十一、从示例看微服务落地的核心方法论
写到这里,你会发现微服务实战并不是简单地“把服务拆开”。真正关键的是下面几件事一起成立:
- 通信标准化:服务之间最好统一使用 gRPC 或稳定的 HTTP 契约。
- 地址动态化:通过 Consul / etcd 做服务注册与发现,避免写死地址。
- 调用可治理:要有负载均衡、超时、熔断和降级。
- 观测可追踪:OpenTelemetry 让跨服务问题可定位。
- 同步异步边界清晰:主流程用同步,派生流程用异步。
- 统一入口治理:通过 API 网关对外暴露能力,而不是让前端直连内部服务。
如果只做“服务拆分”,不做治理、观测和容错,系统复杂度通常会上升;只有把这些配套能力一起建设起来,微服务架构才真正能支撑业务演进。
十二、总结
本文通过一个可运行的 Golang 示例,把微服务架构中的几个核心环节串了起来:
- 用 Proto 定义服务契约,并基于 gRPC 实现服务
- 使用拦截器统一接入日志与链路能力
- 通过 etcd / Consul 实现服务注册与发现
- 在客户端实现轮询负载均衡与熔断降级
- 借助 OpenTelemetry 构建调用链追踪
- 按业务约束选择同步或异步通信方式
- 使用 API Gateway 统一承接外部请求
对于 Golang 而言,微服务并不只是“框架选型”问题,更是工程化能力建设问题。只要把服务契约、治理策略、可观测性与边界设计做扎实,即使系统继续扩张,也能保持较好的演进能力与稳定性。
📝 版权声明:本文为原创技术博客,转载请注明出处。
如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!