1. 微服務概述
| 方面 | 巨石 | 微服務 |
|---|---|---|
| 部署 | 全有或全無 | 每個服務獨立 |
| 縮放 | 擴展一切 | 擴展個人化服務 |
| 技術堆疊 | 單一語言/框架 | 多語言 |
| 團隊 | 大型團隊,共享程式碼庫 | 每項服務的小團隊 |
| 複雜性 | 程式碼複雜度 | 操作複雜性 |
| 數據 | 共享資料庫 | 每個服務的資料庫 |
2. 項目結構
microservices/
├── api-gateway/
│ ├── main.go
│ ├── routes.go
│ └── proxy.go
├── services/
│ ├── user-service/
│ │ ├── cmd/
│ │ │ └── main.go
│ │ ├── internal/
│ │ │ ├── handler/
│ │ │ ├── service/
│ │ │ ├── repository/
│ │ │ └── model/
│ │ ├── proto/
│ │ ├── Dockerfile
│ │ └── go.mod
│ ├── order-service/
│ │ └── ... (same structure)
│ ├── product-service/
│ │ └── ...
│ └── notification-service/
│ └── ...
├── shared/
│ ├── proto/ # Shared protobuf definitions
│ ├── middleware/ # Common middleware
│ └── config/ # Shared configuration
├── docker-compose.yml
└── Makefile
3.API網關
package gateway
import (
"net/http"
"net/http/httputil"
"net/url"
"strings"
"github.com/gin-gonic/gin"
)
type ServiceConfig struct {
Name string
URL string
Prefix string
}
type Gateway struct {
services map[string]*httputil.ReverseProxy
}
func NewGateway(configs []ServiceConfig) *Gateway {
gw := &Gateway{
services: make(map[string]*httputil.ReverseProxy),
}
for _, cfg := range configs {
target, _ := url.Parse(cfg.URL)
proxy := httputil.NewSingleHostReverseProxy(target)
// Custom error handler
proxy.ErrorHandler = func(w http.ResponseWriter, r *http.Request, err error) {
w.WriteHeader(http.StatusBadGateway)
w.Write([]byte(`{"error":"service unavailable"}`))
}
gw.services[cfg.Prefix] = proxy
}
return gw
}
func (gw *Gateway) ProxyHandler(c *gin.Context) {
path := c.Request.URL.Path
for prefix, proxy := range gw.services {
if strings.HasPrefix(path, prefix) {
// Strip prefix
c.Request.URL.Path = strings.TrimPrefix(path, prefix)
proxy.ServeHTTP(c.Writer, c.Request)
return
}
}
c.JSON(404, gin.H{"error": "service not found"})
}
func main() {
gw := NewGateway([]ServiceConfig{
{Name: "user", URL: "http://user-service:8081", Prefix: "/api/users"},
{Name: "order", URL: "http://order-service:8082", Prefix: "/api/orders"},
{Name: "product", URL: "http://product-service:8083", Prefix: "/api/products"},
})
r := gin.Default()
r.Use(RateLimitMiddleware())
r.Use(AuthMiddleware())
r.Any("/api/*path", gw.ProxyHandler)
r.Run(":8080")
}
4. 斷路器
go get github.com/sony/gobreaker
package resilience
import (
"fmt"
"time"
"github.com/sony/gobreaker"
)
func NewCircuitBreaker(name string) *gobreaker.CircuitBreaker {
return gobreaker.NewCircuitBreaker(gobreaker.Settings{
Name: name,
MaxRequests: 3, // Half-open: allow 3 requests
Interval: 60 * time.Second, // Reset counts after 60s
Timeout: 30 * time.Second, // Open → Half-open after 30s
ReadyToTrip: func(counts gobreaker.Counts) bool {
// Open circuit after 5 consecutive failures
return counts.ConsecutiveFailures >= 5
},
OnStateChange: func(name string, from, to gobreaker.State) {
log.Printf("Circuit breaker %s: %s → %s", name, from, to)
},
})
}
// Usage trong service client
type UserClient struct {
baseURL string
cb *gobreaker.CircuitBreaker
client *http.Client
}
func (c *UserClient) GetUser(ctx context.Context, id uint) (*User, error) {
result, err := c.cb.Execute(func() (interface{}, error) {
url := fmt.Sprintf("%s/users/%d", c.baseURL, id)
req, _ := http.NewRequestWithContext(ctx, "GET", url, nil)
resp, err := c.client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode >= 500 {
return nil, fmt.Errorf("server error: %d", resp.StatusCode)
}
var user User
json.NewDecoder(resp.Body).Decode(&user)
return &user, nil
})
if err != nil {
return nil, err
}
return result.(*User), nil
}
5. 使用 Consul 進行服務發現
go get github.com/hashicorp/consul/api
package discovery
import (
"fmt"
"log"
"github.com/hashicorp/consul/api"
)
type ServiceRegistry struct {
client *api.Client
}
func NewServiceRegistry(addr string) (*ServiceRegistry, error) {
config := api.DefaultConfig()
config.Address = addr
client, err := api.NewClient(config)
if err != nil {
return nil, err
}
return &ServiceRegistry{client: client}, nil
}
// Register service
func (r *ServiceRegistry) Register(name, host string, port int) error {
return r.client.Agent().ServiceRegister(&api.AgentServiceRegistration{
ID: fmt.Sprintf("%s-%s-%d", name, host, port),
Name: name,
Address: host,
Port: port,
Check: &api.AgentServiceCheck{
HTTP: fmt.Sprintf("http://%s:%d/health", host, port),
Interval: "10s",
Timeout: "5s",
DeregisterCriticalServiceAfter: "30s",
},
})
}
// Discover service
func (r *ServiceRegistry) Discover(name string) (string, error) {
services, _, err := r.client.Health().Service(name, "", true, nil)
if err != nil {
return "", err
}
if len(services) == 0 {
return "", fmt.Errorf("no healthy instances of %s", name)
}
// Simple round-robin (production: dùng load balancer)
svc := services[0]
return fmt.Sprintf("http://%s:%d", svc.Service.Address, svc.Service.Port), nil
}
// Deregister on shutdown
func (r *ServiceRegistry) Deregister(serviceID string) error {
return r.client.Agent().ServiceDeregister(serviceID)
}
6. 使用 OpenTelemetry 進行分散式追蹤
go get go.opentelemetry.io/otel
go get go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp
go get go.opentelemetry.io/contrib/instrumentation/github.com/gin-gonic/gin/otelgin
package tracing
import (
"context"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp"
"go.opentelemetry.io/otel/sdk/resource"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
semconv "go.opentelemetry.io/otel/semconv/v1.24.0"
)
func InitTracer(ctx context.Context, serviceName string) (func(), error) {
exporter, err := otlptracehttp.New(ctx,
otlptracehttp.WithEndpoint("jaeger:4318"),
otlptracehttp.WithInsecure(),
)
if err != nil {
return nil, err
}
tp := sdktrace.NewTracerProvider(
sdktrace.WithBatcher(exporter),
sdktrace.WithResource(resource.NewWithAttributes(
semconv.SchemaURL,
semconv.ServiceNameKey.String(serviceName),
)),
)
otel.SetTracerProvider(tp)
cleanup := func() {
tp.Shutdown(ctx)
}
return cleanup, nil
}
// Usage in Gin
// r.Use(otelgin.Middleware("user-service"))
7.Saga模式(分散式事務)
// Choreography-based saga
type OrderSaga struct {
eventBus EventBus
}
// Mỗi service lắng nghe events và thực hiện compensation nếu cần
// Order flow: CreateOrder → ReserveInventory → ProcessPayment → ConfirmOrder
// Compensation: CancelOrder ← RestoreInventory ← RefundPayment
type SagaStep struct {
Name string
Execute func(ctx context.Context, data interface{}) error
Compensate func(ctx context.Context, data interface{}) error
}
type SagaOrchestrator struct {
steps []SagaStep
}
func (s *SagaOrchestrator) Run(ctx context.Context, data interface{}) error {
completedSteps := make([]int, 0)
for i, step := range s.steps {
if err := step.Execute(ctx, data); err != nil {
// Compensate completed steps in reverse
for j := len(completedSteps) - 1; j >= 0; j-- {
idx := completedSteps[j]
if compErr := s.steps[idx].Compensate(ctx, data); compErr != nil {
log.Printf("Compensation failed for step %s: %v", s.steps[idx].Name, compErr)
}
}
return fmt.Errorf("saga failed at step %s: %w", step.Name, err)
}
completedSteps = append(completedSteps, i)
}
return nil
}
下一篇: 在 Go 中測試 ——單元測試、整合測試、模擬和測試模式。