Modern Go Engineering
Chapitre 9
09 - Microservices avec Go
09 - Microservices avec Go
Cours 09 : Microservices avec Go
1. Fondamentaux des microservices
1.1 Définition et principes
L'architecture microservices est un style d'architecture où une application est composée de petits services indépendants, chacun exécutant un processus unique et communiquant via des mécanismes légers (HTTP, gRPC, messages).
Principes clés :
- Single Responsibility : Chaque service a une responsabilité unique et bien définie
- Déploiement indépendant : Chaque service peut être déployé sans impact sur les autres
- Isolation des pannes : Une défaillance dans un service ne doit pas impacter les autres
- Autonomie des équipes : Chaque service peut être développé par une équipe différente
- Polyglotte : Chaque service peut utiliser la techno la plus adaptée
1.2 Monolithe vs Microservices
// Monolithe : tout dans un seul package
package monolith
type OrderService struct {
db *sql.DB
cache *redis.Client
email *EmailClient
}
func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
// Validation
if err := s.validateOrder(o); err != nil {
return err
}
// Base de données
if err := s.db.SaveOrder(ctx, o); err != nil {
return err
}
// Cache
s.cache.InvalidateOrders(ctx, o.UserID)
// Email
s.email.SendConfirmation(o.UserEmail, o.ID)
return nil
}
// Microservices : découplage en services spécialisés
package orderservice
type OrderService struct {
orderRepo *OrderRepository
eventBus EventBus
userClient UserServiceClient
}
func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
if err := s.orderRepo.Save(ctx, o); err != nil {
return fmt.Errorf("save order: %w", err)
}
return s.eventBus.Publish(ctx, OrderCreatedEvent{Order: o})
}
1.3 Size des services
La taille d'un service se mesure en :
- Responsabilité fonctionnelle : Un bounded context DDD
- Complexité : Pas plus de quelques centaines de lignes par fichier
- Équipe : 2 pizzas rule (Amazon) — une équipe doit pouvoir tenir autour de 2 pizzas
2. API Gateway
2.1 Rôle de l'API Gateway
L'API Gateway est le point d'entrée unique pour tous les clients. Elle gère :
- Routage vers les services internes
- Authentification et autorisation
- Rate limiting
- Aggregation de réponses
- Transformation de protocoles
package gateway
import (
"net/http"
"net/http/httputil"
"net/url"
)
type Gateway struct {
auth AuthMiddleware
router *http.ServeMux
rateLimit *RateLimiter
targets map[string]*url.URL
}
func NewGateway(auth AuthMiddleware) *Gateway {
g := &Gateway{
auth: auth,
router: http.NewServeMux(),
rateLimit: NewRateLimiter(100, time.Minute),
targets: make(map[string]*url.URL),
}
g.registerRoutes()
return g
}
func (g *Gateway) registerRoutes() {
g.router.HandleFunc("/api/v1/users", g.auth.Wrap(g.proxy("user-service")))
g.router.HandleFunc("/api/v1/orders", g.auth.Wrap(g.proxy("order-service")))
g.router.HandleFunc("/api/v1/events", g.auth.Wrap(g.proxy("event-service")))
}
func (g *Gateway) proxy(serviceName string) http.HandlerFunc {
target := g.targets[serviceName]
proxy := httputil.NewSingleHostReverseProxy(target)
return func(w http.ResponseWriter, r *http.Request) {
if !g.rateLimit.Allow(r.RemoteAddr) {
http.Error(w, "rate limit exceeded", http.StatusTooManyRequests)
return
}
proxy.ServeHTTP(w, r)
}
}
2.2 Patterns API Gateway
// Backend for Frontend (BFF)
type BFFGateway struct {
mobileGateway *Gateway // Optimisé mobile : payload réduit
webGateway *Gateway // Optimisé web : HTML + JSON
iotGateway *Gateway // Optimisé IoT : protobuf binaire
}
// Response aggregation
func (g *Gateway) aggregateOrderDetails(ctx context.Context, orderID string) (*OrderDetail, error) {
// Fan-out requests
type result struct {
data any
err error
}
ch := make(chan result, 3)
go func() {
order, err := g.orderClient.GetOrder(ctx, orderID)
ch <- result{order, err}
}()
go func() {
user, err := g.userClient.GetUser(ctx, orderID)
ch <- result{user, err}
}()
go func() {
payment, err := g.paymentClient.GetPayment(ctx, orderID)
ch <- result{payment, err}
}()
// Aggregate responses
detail := &OrderDetail{}
for i := 0; i < 3; i++ {
r := <-ch
if r.err != nil {
return nil, r.err
}
switch v := r.data.(type) {
case *Order:
detail.Order = v
case *User:
detail.User = v
case *Payment:
detail.Payment = v
}
}
return detail, nil
}
3. Service Registry et Discovery
3.1 Service Registry
package registry
import (
"context"
"sync"
"time"
)
type ServiceInstance struct {
ID string
Name string
Address string
Port int
Metadata map[string]string
Healthy bool
LastCheck time.Time
}
type Registry interface {
Register(ctx context.Context, instance ServiceInstance) error
Deregister(ctx context.Context, instanceID string) error
Discover(ctx context.Context, serviceName string) ([]ServiceInstance, error)
Watch(ctx context.Context, serviceName string) (<-chan []ServiceInstance, error)
HealthCheck(ctx context.Context, instanceID string) error
}
type InMemoryRegistry struct {
mu sync.RWMutex
instances map[string]map[string]*ServiceInstance // serviceName -> instanceID -> instance
watchers map[string][]chan []ServiceInstance
}
func NewInMemoryRegistry() *InMemoryRegistry {
r := &InMemoryRegistry{
instances: make(map[string]map[string]*ServiceInstance),
watchers: make(map[string][]chan []ServiceInstance),
}
go r.healthCheckLoop()
return r
}
func (r *InMemoryRegistry) Register(ctx context.Context, inst ServiceInstance) error {
r.mu.Lock()
defer r.mu.Unlock()
inst.Healthy = true
inst.LastCheck = time.Now()
if r.instances[inst.Name] == nil {
r.instances[inst.Name] = make(map[string]*ServiceInstance)
}
r.instances[inst.Name][inst.ID] = &inst
r.notifyWatchers(inst.Name)
return nil
}
func (r *InMemoryRegistry) Discover(ctx context.Context, serviceName string) ([]ServiceInstance, error) {
r.mu.RLock()
defer r.mu.RUnlock()
instances := r.instances[serviceName]
result := make([]ServiceInstance, 0, len(instances))
for _, inst := range instances {
if inst.Healthy {
result = append(result, *inst)
}
}
return result, nil
}
func (r *InMemoryRegistry) healthCheckLoop() {
ticker := time.NewTicker(30 * time.Second)
for range ticker.C {
r.mu.Lock()
for name, instances := range r.instances {
for id, inst := range instances {
if time.Since(inst.LastCheck) > 60*time.Second {
inst.Healthy = false
r.notifyWatchers(name)
_ = r.Deregister(context.Background(), id)
}
}
}
r.mu.Unlock()
}
}
3.2 Client-side Discovery
package discovery
import (
"math/rand"
"sync"
)
type LoadBalancer interface {
Pick(instances []ServiceInstance) (ServiceInstance, error)
}
type RoundRobinLB struct {
mu sync.Mutex
count map[string]int
}
func (lb *RoundRobinLB) Pick(instances []ServiceInstance) (ServiceInstance, error) {
if len(instances) == 0 {
return ServiceInstance{}, ErrNoInstances
}
lb.mu.Lock()
defer lb.mu.Unlock()
if lb.count == nil {
lb.count = make(map[string]int)
}
serviceName := instances[0].Name
lb.count[serviceName]++
idx := lb.count[serviceName] % len(instances)
return instances[idx], nil
}
type RandomLB struct{}
func (lb *RandomLB) Pick(instances []ServiceInstance) (ServiceInstance, error) {
if len(instances) == 0 {
return ServiceInstance{}, ErrNoInstances
}
return instances[rand.Intn(len(instances))], nil
}
4. Résilience
4.1 Circuit Breaker
package circuitbreaker
import (
"errors"
"sync"
"time"
)
type State int
const (
StateClosed State = iota // Normal operation
StateOpen // Failing, reject requests
StateHalfOpen // Testing if service recovered
)
type CircuitBreaker struct {
mu sync.RWMutex
state State
failureCount int
successCount int
failureThreshold int
successThreshold int
timeout time.Duration
lastStateChange time.Time
}
func New(failureThreshold, successThreshold int, timeout time.Duration) *CircuitBreaker {
return &CircuitBreaker{
state: StateClosed,
failureThreshold: failureThreshold,
successThreshold: successThreshold,
timeout: timeout,
}
}
func (cb *CircuitBreaker) Execute(fn func() error) error {
if !cb.allowRequest() {
return ErrCircuitOpen
}
err := fn()
cb.recordResult(err)
return err
}
func (cb *CircuitBreaker) allowRequest() bool {
cb.mu.RLock()
state := cb.state
cb.mu.RUnlock()
switch state {
case StateClosed:
return true
case StateOpen:
if time.Since(cb.lastStateChange) > cb.timeout {
cb.mu.Lock()
cb.state = StateHalfOpen
cb.successCount = 0
cb.mu.Unlock()
return true
}
return false
case StateHalfOpen:
return true
default:
return true
}
}
func (cb *CircuitBreaker) recordResult(err error) {
cb.mu.Lock()
defer cb.mu.Unlock()
if err != nil {
cb.failureCount++
cb.successCount = 0
if cb.failureCount >= cb.failureThreshold {
cb.state = StateOpen
cb.lastStateChange = time.Now()
}
} else {
cb.successCount++
cb.failureCount = 0
if cb.state == StateHalfOpen && cb.successCount >= cb.successThreshold {
cb.state = StateClosed
}
}
}
4.2 Retry avec backoff
package retry
import (
"context"
"math"
"math/rand"
"time"
)
type Strategy int
const (
StrategyConstant Strategy = iota
StrategyLinear
StrategyExponential
)
type Config struct {
MaxAttempts int
Delay time.Duration
MaxDelay time.Duration
Strategy Strategy
Jitter float64 // 0.0 - 1.0
}
func Do(ctx context.Context, cfg Config, fn func(context.Context) error) error {
var lastErr error
for attempt := 0; attempt < cfg.MaxAttempts; attempt++ {
if attempt > 0 {
delay := calculateDelay(cfg, attempt)
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(delay):
}
}
if err := fn(ctx); err != nil {
lastErr = err
continue
}
return nil
}
return lastErr
}
func calculateDelay(cfg Config, attempt int) time.Duration {
var delay time.Duration
switch cfg.Strategy {
case StrategyConstant:
delay = cfg.Delay
case StrategyLinear:
delay = cfg.Delay * time.Duration(attempt)
case StrategyExponential:
delay = cfg.Delay * time.Duration(math.Pow(2, float64(attempt)))
}
if cfg.MaxDelay > 0 && delay > cfg.MaxDelay {
delay = cfg.MaxDelay
}
if cfg.Jitter > 0 {
jitter := time.Duration(float64(delay) * cfg.Jitter * rand.Float64())
delay += jitter
}
return delay
}
4.3 Timeout avec context
package timeout
import (
"context"
"net/http"
"time"
)
func TimeoutMiddleware(timeout time.Duration) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithTimeout(r.Context(), timeout)
defer cancel()
r = r.WithContext(ctx)
done := make(chan bool)
go func() {
next.ServeHTTP(w, r)
done <- true
}()
select {
case <-done:
return
case <-ctx.Done():
switch ctx.Err() {
case context.DeadlineExceeded:
http.Error(w, "request timeout", http.StatusGatewayTimeout)
case context.Canceled:
http.Error(w, "request cancelled", http.StatusRequestTimeout)
}
}
})
}
}
5. Communication synchrone
5.1 Communication HTTP
package httpclient
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"time"
)
type HTTPClient struct {
client *http.Client
baseURL string
headers map[string]string
}
func NewHTTPClient(baseURL string, timeout time.Duration) *HTTPClient {
return &HTTPClient{
client: &http.Client{
Timeout: timeout,
Transport: &http.Transport{
MaxIdleConns: 100,
MaxIdleConnsPerHost: 100,
IdleConnTimeout: 90 * time.Second,
},
},
baseURL: baseURL,
headers: make(map[string]string),
}
}
func (c *HTTPClient) Get(ctx context.Context, path string, result any) error {
return c.doRequest(ctx, http.MethodGet, path, nil, result)
}
func (c *HTTPClient) Post(ctx context.Context, path string, body, result any) error {
return c.doRequest(ctx, http.MethodPost, path, body, result)
}
func (c *HTTPClient) doRequest(ctx context.Context, method, path string, body, result any) error {
var reqBody io.Reader
if body != nil {
data, err := json.Marshal(body)
if err != nil {
return fmt.Errorf("marshal request: %w", err)
}
reqBody = bytes.NewReader(data)
}
req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, reqBody)
if err != nil {
return fmt.Errorf("create request: %w", err)
}
req.Header.Set("Content-Type", "application/json")
for k, v := range c.headers {
req.Header.Set(k, v)
}
resp, err := c.client.Do(req)
if err != nil {
return fmt.Errorf("do request: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
respBody, _ := io.ReadAll(resp.Body)
return fmt.Errorf("unexpected status %d: %s", resp.StatusCode, string(respBody))
}
if result != nil {
if err := json.NewDecoder(resp.Body).Decode(result); err != nil {
return fmt.Errorf("decode response: %w", err)
}
}
return nil
}
5.2 Communication gRPC
package grpcclient
import (
"context"
"crypto/tls"
"crypto/x509"
"fmt"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
"google.golang.org/grpc/keepalive"
)
type GRPCClient struct {
conn *grpc.ClientConn
}
func NewGRPCClient(address string, opts ...grpc.DialOption) (*GRPCClient, error) {
defaultOpts := []grpc.DialOption{
grpc.WithInsecure(),
grpc.WithDefaultCallOptions(
grpc.MaxCallRecvMsgSize(10 * 1024 * 1024),
grpc.MaxCallSendMsgSize(10 * 1024 * 1024),
),
grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: 30 * time.Second,
Timeout: 10 * time.Second,
PermitWithoutStream: true,
}),
}
opts = append(defaultOpts, opts...)
conn, err := grpc.Dial(address, opts...)
if err != nil {
return nil, fmt.Errorf("dial %s: %w", address, err)
}
return &GRPCClient{conn: conn}, nil
}
func (c *GRPCClient) Close() error {
return c.conn.Close()
}
func WithTLS(caCert, clientCert, clientKey []byte) grpc.DialOption {
certPool := x509.NewCertPool()
certPool.AppendCertsFromPEM(caCert)
cert, err := tls.X509KeyPair(clientCert, clientKey)
if err != nil {
return grpc.WithTransportCredentials(nil) // will fail
}
creds := credentials.NewTLS(&tls.Config{
Certificates: []tls.Certificate{cert},
RootCAs: certPool,
})
return grpc.WithTransportCredentials(creds)
}
6. Communication asynchrone (Events)
6.1 Event Bus
package events
import (
"context"
"encoding/json"
"fmt"
"time"
)
type Event struct {
ID string `json:"id"`
Type string `json:"type"`
Source string `json:"source"`
Time time.Time `json:"time"`
Data any `json:"data"`
SchemaURL string `json:"schema_url,omitempty"`
}
type EventBus interface {
Publish(ctx context.Context, event Event) error
Subscribe(ctx context.Context, eventType string, handler EventHandler) error
Close() error
}
type EventHandler func(context.Context, Event) error
type InMemoryEventBus struct {
subscribers map[string][]EventHandler
buffer chan Event
}
func NewInMemoryEventBus(bufferSize int) *InMemoryEventBus {
bus := &InMemoryEventBus{
subscribers: make(map[string][]EventHandler),
buffer: make(chan Event, bufferSize),
}
go bus.processEvents()
return bus
}
func (bus *InMemoryEventBus) Publish(ctx context.Context, event Event) error {
select {
case bus.buffer <- event:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (bus *InMemoryEventBus) processEvents() {
for event := range bus.buffer {
handlers := bus.subscribers[event.Type]
for _, handler := range handlers {
go func(h EventHandler, e Event) {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
if err := h(ctx, e); err != nil {
fmt.Printf("handler failed for event %s: %v\n", e.Type, err)
}
}(handler, event)
}
}
}
func (bus *InMemoryEventBus) Subscribe(ctx context.Context, eventType string, handler EventHandler) error {
bus.subscribers[eventType] = append(bus.subscribers[eventType], handler)
return nil
}
6.2 Kafka Producer/Consumer
package kafka
import (
"context"
"fmt"
"github.com/segmentio/kafka-go"
)
type Producer struct {
writer *kafka.Writer
}
func NewProducer(brokers []string, topic string) *Producer {
w := &kafka.Writer{
Addr: kafka.TCP(brokers...),
Topic: topic,
Balancer: &kafka.LeastBytes{},
BatchTimeout: 10,
Async: true,
}
return &Producer{writer: w}
}
func (p *Producer) Publish(ctx context.Context, key string, value []byte) error {
msg := kafka.Message{
Key: []byte(key),
Value: value,
}
if err := p.writer.WriteMessages(ctx, msg); err != nil {
return fmt.Errorf("kafka write: %w", err)
}
return nil
}
func (p *Producer) Close() error {
return p.writer.Close()
}
type Consumer struct {
reader *kafka.Reader
}
func NewConsumer(brokers []string, topic, groupID string) *Consumer {
r := kafka.NewReader(kafka.ReaderConfig{
Brokers: brokers,
Topic: topic,
GroupID: groupID,
MinBytes: 10e3,
MaxBytes: 10e6,
MaxWait: 1,
StartOffset: kafka.LastOffset,
})
return &Consumer{reader: r}
}
func (c *Consumer) Consume(ctx context.Context, handler func(key string, value []byte) error) error {
for {
msg, err := c.reader.ReadMessage(ctx)
if err != nil {
return fmt.Errorf("read message: %w", err)
}
if err := handler(string(msg.Key), msg.Value); err != nil {
return fmt.Errorf("handler: %w", err)
}
}
}
func (c *Consumer) Close() error {
return c.reader.Close()
}
7. Event Sourcing
7.1 Aggregate et Events
package eventsourcing
import (
"time"
)
type DomainEvent struct {
ID string `json:"id"`
AggregateID string `json:"aggregate_id"`
Type string `json:"type"`
Data []byte `json:"data"`
Version int `json:"version"`
Timestamp time.Time `json:"timestamp"`
}
type Aggregate interface {
AggregateID() string
ApplyEvent(event DomainEvent)
UncommittedEvents() []DomainEvent
ClearEvents()
Version() int
}
type OrderAggregate struct {
id string
status string
items []OrderItem
total float64
version int
events []DomainEvent
}
func (a *OrderAggregate) CreateOrder(id string, items []OrderItem) {
event := DomainEvent{
AggregateID: id,
Type: "OrderCreated",
Data: mustMarshal(OrderCreatedData{Items: items}),
Version: 1,
Timestamp: time.Now(),
}
a.ApplyEvent(event)
}
func (a *OrderAggregate) ApplyEvent(event DomainEvent) {
switch event.Type {
case "OrderCreated":
var data OrderCreatedData
mustUnmarshal(event.Data, &data)
a.status = "created"
a.items = data.Items
a.version = event.Version
case "OrderPaid":
a.status = "paid"
a.version = event.Version
case "OrderShipped":
a.status = "shipped"
a.version = event.Version
}
a.events = append(a.events, event)
}
type EventStore interface {
SaveEvents(aggregateID string, events []DomainEvent, expectedVersion int) error
GetEvents(aggregateID string) ([]DomainEvent, error)
}
type PostgresEventStore struct {
db *sql.DB
}
func (s *PostgresEventStore) SaveEvents(aggregateID string, events []DomainEvent, expectedVersion int) error {
tx, err := s.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
for _, event := range events {
_, err := tx.Exec(
`INSERT INTO events (aggregate_id, type, data, version, timestamp)
VALUES ($1, $2, $3, $4, $5)`,
aggregateID, event.Type, event.Data, event.Version, event.Timestamp,
)
if err != nil {
return err
}
}
return tx.Commit()
}
8. Saga Pattern
8.1 Choreography Saga
package saga
import (
"context"
"fmt"
)
type SagaStep struct {
Name string
Execute func(context.Context) error
Compensate func(context.Context) error
}
type Saga struct {
steps []SagaStep
}
func (s *Saga) Execute(ctx context.Context) error {
executed := make([]int, 0, len(s.steps))
for i, step := range s.steps {
if err := step.Execute(ctx); err != nil {
// Compensate all executed steps in reverse order
return s.compensate(ctx, executed)
}
executed = append(executed, i)
}
return nil
}
func (s *Saga) compensate(ctx context.Context, executed []int) error {
var errs []error
for i := len(executed) - 1; i >= 0; i-- {
step := s.steps[executed[i]]
if err := step.Compensate(ctx); err != nil {
errs = append(errs, fmt.Errorf("compensate step %s: %w", step.Name, err))
}
}
if len(errs) > 0 {
return fmt.Errorf("saga compensation errors: %v", errs)
}
return nil
}
// OrderSaga example
func NewOrderSaga(orderService, paymentService, inventoryService ServiceClient) *Saga {
return &Saga{
steps: []SagaStep{
{
Name: "create-order",
Execute: func(ctx context.Context) error {
return orderService.CreateOrder(ctx)
},
Compensate: func(ctx context.Context) error {
return orderService.CancelOrder(ctx)
},
},
{
Name: "reserve-inventory",
Execute: func(ctx context.Context) error {
return inventoryService.Reserve(ctx)
},
Compensate: func(ctx context.Context) error {
return inventoryService.Release(ctx)
},
},
{
Name: "process-payment",
Execute: func(ctx context.Context) error {
return paymentService.Charge(ctx)
},
Compensate: func(ctx context.Context) error {
return paymentService.Refund(ctx)
},
},
},
}
}
8.2 Orchestrator Pattern
package sagaorchestrator
import (
"context"
"fmt"
)
type SagaState int
const (
SagaPending SagaState = iota
SagaCompleted
SagaCompensating
SagaFailed
)
type Orchestrator struct {
sagaID string
state SagaState
steps []SagaStep
currentStep int
}
func (o *Orchestrator) ProcessMessage(ctx context.Context, msg Message) error {
switch msg.Type {
case "StepCompleted":
return o.onStepCompleted(ctx, msg)
case "StepFailed":
return o.onStepFailed(ctx, msg)
case "CompensationCompleted":
return o.onCompensationCompleted(ctx, msg)
}
return nil
}
func (o *Orchestrator) onStepCompleted(ctx context.Context, msg Message) error {
o.currentStep++
if o.currentStep >= len(o.steps) {
o.state = SagaCompleted
return nil
}
return o.executeStep(ctx, o.currentStep)
}
func (o *Orchestrator) onStepFailed(ctx context.Context, msg Message) error {
o.state = SagaCompensating
return o.compensate(ctx, o.currentStep-1)
}
func (o *Orchestrator) compensate(ctx context.Context, from int) error {
for i := from; i >= 0; i-- {
step := o.steps[i]
if err := step.Compensate(ctx); err != nil {
return fmt.Errorf("compensation failed at step %d: %w", i, err)
}
}
o.state = SagaFailed
return nil
}
9. Outbox Pattern
package outbox
import (
"context"
"encoding/json"
"fmt"
"time"
)
type OutboxMessage struct {
ID string `json:"id"`
AggregateType string `json:"aggregate_type"`
AggregateID string `json:"aggregate_id"`
EventType string `json:"event_type"`
Payload []byte `json:"payload"`
Status string `json:"status"` // pending, published, failed
CreatedAt time.Time `json:"created_at"`
PublishedAt *time.Time `json:"published_at,omitempty"`
}
type OutboxRepository interface {
Insert(ctx context.Context, msg OutboxMessage) error
GetPending(ctx context.Context, limit int) ([]OutboxMessage, error)
MarkPublished(ctx context.Context, id string) error
MarkFailed(ctx context.Context, id string) error
}
type OutboxRelay struct {
repo OutboxRepository
publisher EventPublisher
ticker *time.Ticker
}
func NewOutboxRelay(repo OutboxRepository, publisher EventPublisher, interval time.Duration) *OutboxRelay {
return &OutboxRelay{
repo: repo,
publisher: publisher,
ticker: time.NewTicker(interval),
}
}
func (r *OutboxRelay) Start(ctx context.Context) {
for {
select {
case <-r.ticker.C:
r.processPending(ctx)
case <-ctx.Done():
return
}
}
}
func (r *OutboxRelay) processPending(ctx context.Context) {
messages, err := r.repo.GetPending(ctx, 100)
if err != nil {
fmt.Printf("get pending messages: %v\n", err)
return
}
for _, msg := range messages {
if err := r.publisher.Publish(ctx, msg.EventType, msg.Payload); err != nil {
fmt.Printf("publish message %s: %v\n", msg.ID, err)
r.repo.MarkFailed(ctx, msg.ID)
continue
}
r.repo.MarkPublished(ctx, msg.ID)
}
}
// Usage in service
func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
tx, err := s.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
// 1. Save order
if err := s.orderRepo.Save(ctx, o); err != nil {
return err
}
// 2. Save outbox message in same transaction
payload, _ := json.Marshal(OrderCreatedEvent{OrderID: o.ID})
msg := OutboxMessage{
AggregateID: o.ID,
EventType: "order.created",
Payload: payload,
Status: "pending",
CreatedAt: time.Now(),
}
if err := s.outboxRepo.Insert(ctx, msg); err != nil {
return err
}
return tx.Commit() // Atomic: both saved or none
}
10. Service Mesh
10.1 Linkerd
# linkerd-profile.yaml
apiVersion: linkerd.io/v1alpha2
kind: ServiceProfile
metadata:
name: order-service.default.svc.cluster.local
spec:
routes:
- name: GET /api/orders
condition:
method: GET
pathRegex: /api/orders.*
isRetryable: true
timeout: 500ms
- name: POST /api/orders
condition:
method: POST
pathRegex: /api/orders
timeout: 2s
retryBudget:
retryRatio: 0.2
minRetriesPerSecond: 10
ttl: 10s
10.2 Istio
# istio-virtualservice.yaml
apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
name: order-service
spec:
hosts:
- order-service
http:
- match:
- uri:
prefix: /api/orders
route:
- destination:
host: order-service
subset: v1
weight: 90
- destination:
host: order-service
subset: v2
weight: 10
retries:
attempts: 3
perTryTimeout: 500ms
timeout: 2s
fault:
delay:
percentage:
value: 0.1
fixedDelay: 5s
---
apiVersion: networking.istio.io/v1beta1
kind: DestinationRule
metadata:
name: order-service
spec:
host: order-service
trafficPolicy:
connectionPool:
tcp:
maxConnections: 100
http:
http1MaxPendingRequests: 10
http2MaxRequests: 1000
loadBalancer:
simple: LEAST_CONN
outlierDetection:
consecutive5xxErrors: 5
interval: 30s
baseEjectionTime: 30s
subsets:
- name: v1
labels:
version: v1
- name: v2
labels:
version: v2
11. Observabilité (OpenTelemetry)
11.1 Tracing
package telemetry
import (
"context"
"fmt"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/exporters/otlp/otlptrace"
"go.opentelemetry.io/otel/sdk/resource"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
semconv "go.opentelemetry.io/otel/semconv/v1.17.0"
"go.opentelemetry.io/otel/trace"
)
var tracer trace.Tracer
func InitTracing(serviceName, endpoint string) (*sdktrace.TracerProvider, error) {
exporter, err := otlptrace.New(context.Background(),
otlptrace.WithEndpoint(endpoint),
otlptrace.WithInsecure(),
)
if err != nil {
return nil, fmt.Errorf("create exporter: %w", err)
}
tp := sdktrace.NewTracerProvider(
sdktrace.WithBatcher(exporter),
sdktrace.WithResource(resource.NewWithAttributes(
semconv.SchemaURL,
semconv.ServiceName(serviceName),
attribute.String("environment", "production"),
)),
)
otel.SetTracerProvider(tp)
tracer = tp.Tracer(serviceName)
return tp, nil
}
// Span helpers
func StartSpan(ctx context.Context, name string, opts ...trace.SpanStartOption) (context.Context, trace.Span) {
return tracer.Start(ctx, name, opts...)
}
// HTTP middleware
func TracingMiddleware(serviceName string) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
ctx, span := StartSpan(r.Context(), fmt.Sprintf("%s %s", r.Method, r.URL.Path))
defer span.End()
span.SetAttributes(
attribute.String("http.method", r.Method),
attribute.String("http.url", r.URL.String()),
attribute.String("http.host", r.Host),
)
next.ServeHTTP(w, r.WithContext(ctx))
})
}
}
11.2 Metrics
package metrics
import (
"net/http"
"time"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
type Metrics struct {
requestDuration *prometheus.HistogramVec
requestTotal *prometheus.CounterVec
activeRequests *prometheus.GaugeVec
errorsTotal *prometheus.CounterVec
}
func NewMetrics(namespace, subsystem string) *Metrics {
m := &Metrics{
requestDuration: prometheus.NewHistogramVec(
prometheus.HistogramOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "request_duration_seconds",
Help: "Request duration in seconds",
Buckets: prometheus.DefBuckets,
},
[]string{"method", "path", "status"},
),
requestTotal: prometheus.NewCounterVec(
prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "requests_total",
Help: "Total number of requests",
},
[]string{"method", "path", "status"},
),
activeRequests: prometheus.NewGaugeVec(
prometheus.GaugeOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "active_requests",
Help: "Active requests",
},
[]string{"method"},
),
errorsTotal: prometheus.NewCounterVec(
prometheus.CounterOpts{
Namespace: namespace,
Subsystem: subsystem,
Name: "errors_total",
Help: "Total number of errors",
},
[]string{"type"},
),
}
prometheus.MustRegister(
m.requestDuration,
m.requestTotal,
m.activeRequests,
m.errorsTotal,
)
return m
}
func (m *Metrics) Handler() http.Handler {
return promhttp.Handler()
}
func (m *Metrics) ObserveRequest(method, path string, status int, duration time.Duration) {
m.requestDuration.WithLabelValues(method, path, fmt.Sprintf("%d", status)).Observe(duration.Seconds())
m.requestTotal.WithLabelValues(method, path, fmt.Sprintf("%d", status)).Inc()
}
12. Health Checks
12.1 Health Check Handler
package health
import (
"context"
"encoding/json"
"net/http"
"time"
)
type Checker interface {
Name() string
Check(ctx context.Context) error
}
type HealthHandler struct {
checkers []Checker
}
func NewHealthHandler(checkers ...Checker) *HealthHandler {
return &HealthHandler{checkers: checkers}
}
func (h *HealthHandler) Liveness(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(map[string]string{"status": "alive"})
}
func (h *HealthHandler) Readiness(w http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithTimeout(r.Context(), 5*time.Second)
defer cancel()
status := http.StatusOK
result := make(map[string]string)
for _, checker := range h.checkers {
if err := checker.Check(ctx); err != nil {
result[checker.Name()] = fmt.Sprintf("unhealthy: %v", err)
status = http.StatusServiceUnavailable
} else {
result[checker.Name()] = "healthy"
}
}
w.WriteHeader(status)
json.NewEncoder(w).Encode(result)
}
// Database checker
type DBChecker struct {
db *sql.DB
}
func (c *DBChecker) Name() string { return "database" }
func (c *DBChecker) Check(ctx context.Context) error {
return c.db.PingContext(ctx)
}
// Redis checker
type RedisChecker struct {
client *redis.Client
}
func (c *RedisChecker) Name() string { return "redis" }
func (c *RedisChecker) Check(ctx context.Context) error {
return c.client.Ping(ctx).Err()
}
// Kafka checker
type KafkaChecker struct {
brokers []string
}
func (c *KafkaChecker) Name() string { return "kafka" }
func (c *KafkaChecker) Check(ctx context.Context) error {
conn, err := kafka.Dial("tcp", c.brokers[0])
if err != nil {
return err
}
return conn.Close()
}
12.2 Graceful Degradation
package degradation
import (
"context"
"net/http"
)
type DegradationState struct {
CacheAvailable bool
DatabaseAvailable bool
KafkaAvailable bool
mode string // "full", "cache_only", "read_only"
}
type DegradableHandler struct {
state *DegradationState
fallback http.Handler
}
func (h *DegradableHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
switch h.state.mode {
case "full":
// Normal operation
h.fallback.ServeHTTP(w, r)
case "cache_only":
// Serve stale cache data
w.Header().Set("X-Cache", "HIT")
w.Header().Set("X-Degradation-Mode", "cache_only")
h.serveFromCache(w, r)
case "read_only":
// Block writes, allow reads
if r.Method != http.MethodGet && r.Method != http.MethodHead {
http.Error(w, "service in read-only mode", http.StatusServiceUnavailable)
return
}
h.fallback.ServeHTTP(w, r)
}
}
// Graceful shutdown
func GracefulShutdown(srv *http.Server, timeout time.Duration) {
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
if err := srv.Shutdown(ctx); err != nil {
log.Printf("server shutdown: %v", err)
}
}
Diagrammes Mermaid
Architecture microservices
Diagramme en cours de génération...
Saga Pattern
Diagramme en cours de génération...
Outbox Pattern
Diagramme en cours de génération...
Circuit Breaker States
Diagramme en cours de génération...
Points Clés à Retenir
- Microservices ≠ petits monolithes — Chaque service doit avoir une responsabilité unique
- Communication — Préférer async (events) pour le couplage faible, sync (gRPC) pour le temps réel
- Résilience — Toujours implémenter circuit breaker, retry, timeout
- Observabilité — Metrics, tracing, logging sont obligatoires, pas optionnels
- Saga — Pour les transactions distribuées, toujours prévoir la compensation
- Outbox — Garantit la delivery des événements sans 2PC
- Health Checks — Liveness (process) ≠ Readiness (dépendances)
- Graceful Degradation — Un service doit continuer à fonctionner partiellement
Références
- "Building Microservices" by Sam Newman
- "Microservices Patterns" by Chris Richardson
- Go kit microservices toolkit
- OpenTelemetry Go SDK documentation