Реализация паттерна Saga в микросервисах на Go: оркестрация и хореография транзакций с PostgreSQL и Redis
Введение: почему распределённые транзакции сложны
В монолитном приложении транзакция — это атомарная операция в рамках одной базы данных. ACID-гарантии PostgreSQL делают её надёжной и предсказуемой. Но когда система разбивается на микросервисы, каждый со своей БД, классические транзакции становятся невозможны: нельзя вызвать BEGIN ... COMMIT одновременно в двух изолированных хранилищах.
Двухфазный коммит (2PC) формально решает эту проблему, но на практике создаёт жёсткую связанность сервисов, блокирующие ожидания и единую точку отказа в координаторе. Для высоконагруженных систем это неприемлемо.
Паттерн Saga — альтернативный подход: вместо одной распределённой транзакции выполняется последовательность локальных транзакций в каждом сервисе. Если на каком-то шаге происходит сбой, запускаются компенсирующие транзакции для отмены уже выполненных шагов. Это не rollback в привычном смысле — это явные обратные действия, реализованные в коде.
Два подхода: оркестрация vs хореография
Saga-паттерн реализуется двумя фундаментально разными способами. Выбор между ними определяет архитектуру всей системы.
Хореография (Choreography)
Каждый сервис публикует события после выполнения локальной транзакции. Другие сервисы подписываются на эти события и реагируют на них. Центрального координатора нет — сервисы сами знают, что делать в ответ на конкретное событие.
- Плюсы: полная децентрализация, низкая связанность, простота добавления новых участников.
- Минусы: сложно отследить состояние всей саги целиком, логика размазана по сервисам, отладка превращается в детективное расследование.
Оркестрация (Orchestration)
Центральный оркестратор (Saga Orchestrator) управляет всеми шагами: вызывает сервисы, ждёт ответов, принимает решения о компенсации. Он хранит состояние саги — обычно в базе данных.
- Плюсы: полная видимость состояния саги, централизованная логика, понятный порядок шагов, простая отладка.
- Минусы: оркестратор — потенциальная точка отказа, более высокая связанность с координатором, риск превращения оркестратора в «толстый» сервис.
Сравнение подходов
- Состояние саги: при хореографии — распределено по событиям и сервисам; при оркестрации — централизовано в БД оркестратора.
- Связанность: хореография даёт слабую связанность между сервисами; оркестрация создаёт зависимость всех сервисов от оркестратора.
- Наблюдаемость: хореография требует агрегации логов и трейсов; оркестрация предоставляет полное состояние из одного места.
- Когда выбирать хореографию: простые саги (2–3 шага), команда уже работает с event-driven архитектурой, сервисы слабо связаны.
- Когда выбирать оркестрацию: сложные бизнес-процессы (5+ шагов), требуется явный мониторинг состояния, важна централизованная компенсация.
Реализация хореографической Saga на Go с Redis Streams
Для хореографии хорошо подходят Redis Streams — они обеспечивают персистентную очередь сообщений с группами потребителей. В отличие от простого Pub/Sub, сообщения сохраняются и могут быть перечитаны при сбое.
Структура событий
// events.go
package saga
import "time"
type EventType string
const (
OrderCreated EventType = "order.created"
PaymentProcessed EventType = "payment.processed"
PaymentFailed EventType = "payment.failed"
InventoryReserved EventType = "inventory.reserved"
InventoryFailed EventType = "inventory.failed"
OrderCompleted EventType = "order.completed"
OrderCancelled EventType = "order.cancelled"
)
type SagaEvent struct {
SagaID string `json:"saga_id"`
EventType EventType `json:"event_type"`
Payload map[string]string `json:"payload"`
Timestamp time.Time `json:"timestamp"`
}
Публикация событий в Redis Streams
// publisher.go
package saga
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
type EventPublisher struct {
rdb *redis.Client
stream string
}
func NewEventPublisher(rdb *redis.Client, stream string) *EventPublisher {
return &EventPublisher{rdb: rdb, stream: stream}
}
func (p *EventPublisher) Publish(ctx context.Context, event SagaEvent) error {
event.Timestamp = time.Now()
data, err := json.Marshal(event)
if err != nil {
return fmt.Errorf("marshal event: %w", err)
}
return p.rdb.XAdd(ctx, &redis.XAddArgs{
Stream: p.stream,
Values: map[string]interface{}{
"saga_id": event.SagaID,
"event_type": string(event.EventType),
"data": string(data),
},
}).Err()
}
Потребитель событий с Consumer Group
// consumer.go
package saga
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
"github.com/redis/go-redis/v9"
)
type EventHandler func(ctx context.Context, event SagaEvent) error
type EventConsumer struct {
rdb *redis.Client
stream string
group string
consumer string
handlers map[EventType]EventHandler
}
func NewEventConsumer(rdb *redis.Client, stream, group, consumer string) *EventConsumer {
return &EventConsumer{
rdb: rdb,
stream: stream,
group: group,
consumer: consumer,
handlers: make(map[EventType]EventHandler),
}
}
func (c *EventConsumer) Register(eventType EventType, handler EventHandler) {
c.handlers[eventType] = handler
}
func (c *EventConsumer) Start(ctx context.Context) error {
// Создаём группу потребителей, если не существует
c.rdb.XGroupCreateMkStream(ctx, c.stream, c.group, "0")
for {
select {
case <-ctx.Done():
return ctx.Err()
default:
}
streams, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: c.group,
Consumer: c.consumer,
Streams: []string{c.stream, ">"},
Count: 10,
Block: 5 * time.Second,
}).Result()
if err != nil && err != redis.Nil {
log.Printf("xreadgroup error: %v", err)
continue
}
for _, stream := range streams {
for _, msg := range stream.Messages {
if err := c.processMessage(ctx, msg); err != nil {
log.Printf("process message %s error: %v", msg.ID, err)
continue // сообщение останется в PEL для ретрая
}
// ACK только при успешной обработке
c.rdb.XAck(ctx, c.stream, c.group, msg.ID)
}
}
}
}
func (c *EventConsumer) processMessage(ctx context.Context, msg redis.XMessage) error {
dataStr, ok := msg.Values["data"].(string)
if !ok {
return fmt.Errorf("invalid message format")
}
var event SagaEvent
if err := json.Unmarshal([]byte(dataStr), &event); err != nil {
return fmt.Errorf("unmarshal event: %w", err)
}
handler, ok := c.handlers[event.EventType]
if !ok {
return nil // игнорируем неизвестные события
}
return handler(ctx, event)
}
Пример: сервис оплаты в хореографической саге
// payment_service.go
package payment
import (
"context"
"fmt"
"log"
"github.com/yourorg/saga"
)
type PaymentService struct {
publisher *saga.EventPublisher
repo *PaymentRepository
}
func (s *PaymentService) HandleOrderCreated(ctx context.Context, event saga.SagaEvent) error {
orderID := event.Payload["order_id"]
amount := event.Payload["amount"]
// Идемпотентная проверка — не обрабатывать дважды
if exists, _ := s.repo.ExistsPayment(ctx, event.SagaID); exists {
log.Printf("payment for saga %s already processed, skipping", event.SagaID)
return nil
}
err := s.repo.ProcessPayment(ctx, event.SagaID, orderID, amount)
if err != nil {
// Публикуем событие неудачи — другие сервисы запустят компенсацию
return s.publisher.Publish(ctx, saga.SagaEvent{
SagaID: event.SagaID,
EventType: saga.PaymentFailed,
Payload: map[string]string{"reason": err.Error()},
})
}
return s.publisher.Publish(ctx, saga.SagaEvent{
SagaID: event.SagaID,
EventType: saga.PaymentProcessed,
Payload: map[string]string{"order_id": orderID},
})
}
func (s *PaymentService) HandleOrderCancelled(ctx context.Context, event saga.SagaEvent) error {
// Компенсирующая транзакция: возврат платежа
if err := s.repo.RefundPayment(ctx, event.SagaID); err != nil {
return fmt.Errorf("refund payment: %w", err)
}
log.Printf("payment refunded for saga %s", event.SagaID)
return nil
}
Реализация оркестрационной Saga на Go с PostgreSQL
Оркестратор хранит состояние каждой саги в PostgreSQL. Это позволяет восстановиться после сбоя: при перезапуске оркестратор читает незавершённые саги и продолжает выполнение.
Схема базы данных
-- migrations/001_saga_state.sql
CREATE TYPE saga_status AS ENUM (
'started',
'processing',
'completed',
'compensating',
'failed'
);
CREATE TABLE saga_instances (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
saga_type VARCHAR(100) NOT NULL,
status saga_status NOT NULL DEFAULT 'started',
current_step VARCHAR(100),
payload JSONB NOT NULL DEFAULT '{}',
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE TABLE saga_steps (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
saga_id UUID NOT NULL REFERENCES saga_instances(id),
step_name VARCHAR(100) NOT NULL,
status VARCHAR(50) NOT NULL DEFAULT 'pending',
result JSONB,
executed_at TIMESTAMPTZ,
UNIQUE(saga_id, step_name)
);
CREATE INDEX idx_saga_instances_status ON saga_instances(status);
Модель состояния и репозиторий
// orchestrator/state.go
package orchestrator
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
"github.com/google/uuid"
)
type SagaStatus string
const (
StatusStarted SagaStatus = "started"
StatusProcessing SagaStatus = "processing"
StatusCompleted SagaStatus = "completed"
StatusCompensating SagaStatus = "compensating"
StatusFailed SagaStatus = "failed"
)
type SagaInstance struct {
ID uuid.UUID `db:"id"`
SagaType string `db:"saga_type"`
Status SagaStatus `db:"status"`
CurrentStep string `db:"current_step"`
Payload map[string]any `db:"payload"`
CreatedAt time.Time `db:"created_at"`
UpdatedAt time.Time `db:"updated_at"`
}
type SagaRepository struct {
db *sql.DB
}
func (r *SagaRepository) Create(ctx context.Context, sagaType string, payload map[string]any) (*SagaInstance, error) {
payloadJSON, err := json.Marshal(payload)
if err != nil {
return nil, err
}
instance := &SagaInstance{
ID: uuid.New(),
SagaType: sagaType,
Status: StatusStarted,
Payload: payload,
}
_, err = r.db.ExecContext(ctx,
`INSERT INTO saga_instances (id, saga_type, status, payload)
VALUES ($1, $2, $3, $4)`,
instance.ID, instance.SagaType, instance.Status, payloadJSON,
)
return instance, err
}
func (r *SagaRepository) UpdateStatus(ctx context.Context, sagaID uuid.UUID, status SagaStatus, step string) error {
_, err := r.db.ExecContext(ctx,
`UPDATE saga_instances
SET status = $1, current_step = $2, updated_at = NOW()
WHERE id = $3`,
status, step, sagaID,
)
return err
}
func (r *SagaRepository) MarkStepDone(ctx context.Context, sagaID uuid.UUID, stepName string, result map[string]any) error {
resultJSON, _ := json.Marshal(result)
_, err := r.db.ExecContext(ctx,
`INSERT INTO saga_steps (saga_id, step_name, status, result, executed_at)
VALUES ($1, $2, 'completed', $3, NOW())
ON CONFLICT (saga_id, step_name) DO UPDATE
SET status = 'completed', result = $3, executed_at = NOW()`,
sagaID, stepName, resultJSON,
)
return err
}
func (r *SagaRepository) IsStepDone(ctx context.Context, sagaID uuid.UUID, stepName string) (bool, error) {
var count int
err := r.db.QueryRowContext(ctx,
`SELECT COUNT(*) FROM saga_steps
WHERE saga_id = $1 AND step_name = $2 AND status = 'completed'`,
sagaID, stepName,
).Scan(&count)
return count > 0, err
}
func (r *SagaRepository) GetPendingSagas(ctx context.Context) ([]*SagaInstance, error) {
rows, err := r.db.QueryContext(ctx,
`SELECT id, saga_type, status, current_step, payload, created_at, updated_at
FROM saga_instances
WHERE status IN ('started', 'processing', 'compensating')
AND updated_at < NOW() - INTERVAL '5 minutes'`)
if err != nil {
return nil, fmt.Errorf("query pending sagas: %w", err)
}
defer rows.Close()
var instances []*SagaInstance
for rows.Next() {
var inst SagaInstance
var payloadJSON []byte
if err := rows.Scan(&inst.ID, &inst.SagaType, &inst.Status,
&inst.CurrentStep, &payloadJSON, &inst.CreatedAt, &inst.UpdatedAt); err != nil {
return nil, err
}
json.Unmarshal(payloadJSON, &inst.Payload)
instances = append(instances, &inst)
}
return instances, rows.Err()
}
Центральный оркестратор
// orchestrator/orchestrator.go
package orchestrator
import (
"context"
"fmt"
"log"
"github.com/google/uuid"
)
type Step struct {
Name string
Execute func(ctx context.Context, payload map[string]any) (map[string]any, error)
Compensate func(ctx context.Context, payload map[string]any) error
}
type SagaOrchestrator struct {
repo *SagaRepository
steps []Step
}
func NewSagaOrchestrator(repo *SagaRepository, steps []Step) *SagaOrchestrator {
return &SagaOrchestrator{repo: repo, steps: steps}
}
func (o *SagaOrchestrator) Execute(ctx context.Context, sagaType string, payload map[string]any) error {
instance, err := o.repo.Create(ctx, sagaType, payload)
if err != nil {
return fmt.Errorf("create saga: %w", err)
}
return o.run(ctx, instance)
}
func (o *SagaOrchestrator) run(ctx context.Context, instance *SagaInstance) error {
executedSteps := []int{}
for i, step := range o.steps {
// Идемпотентная проверка: уже выполнен?
done, err := o.repo.IsStepDone(ctx, instance.ID, step.Name)
if err != nil {
return fmt.Errorf("check step %s: %w", step.Name, err)
}
if done {
executedSteps = append(executedSteps, i)
continue
}
o.repo.UpdateStatus(ctx, instance.ID, StatusProcessing, step.Name)
log.Printf("[saga %s] executing step: %s", instance.ID, step.Name)
result, err := step.Execute(ctx, instance.Payload)
if err != nil {
log.Printf("[saga %s] step %s failed: %v — starting compensation", instance.ID, step.Name, err)
o.repo.UpdateStatus(ctx, instance.ID, StatusCompensating, step.Name)
return o.compensate(ctx, instance, executedSteps)
}
// Сохраняем результат шага для идемпотентности
if err := o.repo.MarkStepDone(ctx, instance.ID, step.Name, result); err != nil {
return fmt.Errorf("mark step done: %w", err)
}
executedSteps = append(executedSteps, i)
// Обновляем payload результатами шага
for k, v := range result {
instance.Payload[k] = v
}
}
return o.repo.UpdateStatus(ctx, instance.ID, StatusCompleted, "")
}
func (o *SagaOrchestrator) compensate(ctx context.Context, instance *SagaInstance, executedSteps []int) error {
// Компенсируем в обратном порядке
for i := len(executedSteps) - 1; i >= 0; i-- {
stepIdx := executedSteps[i]
step := o.steps[stepIdx]
if step.Compensate == nil {
continue
}
log.Printf("[saga %s] compensating step: %s", instance.ID, step.Name)
if err := step.Compensate(ctx, instance.Payload); err != nil {
// Логируем ошибку, но продолжаем компенсацию других шагов
log.Printf("[saga %s] compensation of %s failed: %v", instance.ID, step.Name, err)
}
}
return o.repo.UpdateStatus(ctx, instance.ID, StatusFailed, "")
}
Использование оркестратора
// main.go (пример сборки)
package main
import (
"context"
"database/sql"
"log"
_ "github.com/lib/pq"
"github.com/yourorg/orchestrator"
)
func main() {
db, _ := sql.Open("postgres", "postgres://user:pass@localhost/orders?sslmode=disable")
repo := &orchestrator.SagaRepository{} // инициализация с db
orderSaga := orchestrator.NewSagaOrchestrator(repo, []orchestrator.Step{
{
Name: "reserve_inventory",
Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
// вызов inventory service
return map[string]any{"reservation_id": "res-123"}, nil
},
Compensate: func(ctx context.Context, payload map[string]any) error {
// отмена резервации
return nil
},
},
{
Name: "process_payment",
Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
// вызов payment service
return map[string]any{"payment_id": "pay-456"}, nil
},
Compensate: func(ctx context.Context, payload map[string]any) error {
// возврат платежа
return nil
},
},
{
Name: "confirm_order",
Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
return map[string]any{"status": "confirmed"}, nil
},
Compensate: nil, // финальный шаг, компенсация не нужна
},
})
ctx := context.Background()
if err := orderSaga.Execute(ctx, "order_saga", map[string]any{
"order_id": "ord-789",
"amount": "150.00",
"user_id": "usr-001",
}); err != nil {
log.Printf("saga failed: %v", err)
}
}
Обработка компенсирующих транзакций и идемпотентность
Компенсирующие транзакции — не rollback. Это отдельные бизнес-операции, которые должны быть реализованы явно. Несколько критически важных правил:
Идемпотентность — обязательное требование
Любая операция в саге (и прямая, и компенсирующая) может быть вызвана более одного раза из-за ретраев после сбоя сети. Каждый обработчик должен корректно работать при повторном вызове.
Стандартный паттерн — хранить idempotency_key (обычно saga_id + step_name) в таблице результатов. Перед выполнением операции проверяем наличие записи с таким ключом.
// idempotency.go
func (r *PaymentRepository) ProcessPaymentIdempotent(
ctx context.Context,
idempotencyKey string,
amount string,
) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
// Атомарная проверка + вставка через INSERT ... ON CONFLICT DO NOTHING
result, err := tx.ExecContext(ctx,
`INSERT INTO payment_operations (idempotency_key, amount, created_at)
VALUES ($1, $2, NOW())
ON CONFLICT (idempotency_key) DO NOTHING`,
idempotencyKey, amount,
)
if err != nil {
return fmt.Errorf("insert payment operation: %w", err)
}
rows, _ := result.RowsAffected()
if rows == 0 {
// Операция уже выполнена — возвращаем успех
return tx.Commit()
}
// Выполняем реальный платёж
if err := r.chargeCard(ctx, tx, amount); err != nil {
return err
}
return tx.Commit()
}
Ретраи с экспоненциальным backoff
При временных сбоях (сеть, перегрузка БД) используйте ретраи с jitter, чтобы не создавать «грозовые стада»:
// retry.go
package saga
import (
"context"
"math"
"math/rand"
"time"
)
func WithRetry(ctx context.Context, maxAttempts int, fn func() error) error {
var lastErr error
for attempt := 0; attempt < maxAttempts; attempt++ {
if err := fn(); err == nil {
return nil
} else {
lastErr = err
}
if attempt == maxAttempts-1 {
break
}
backoff := time.Duration(math.Pow(2, float64(attempt))) * 100 * time.Millisecond
jitter := time.Duration(rand.Intn(100)) * time.Millisecond
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(backoff + jitter):
}
}
return lastErr
}
Мониторинг и отладка Saga-транзакций
Без наблюдаемости саги превращаются в «чёрный ящик». Для production-систем необходимо:
- Structured logging: каждое событие саги логируется с полями
saga_id,step,status,duration_ms. Используйтеlog/slogилиzerolog. - Метрики: счётчики запущенных, завершённых и скомпенсированных саг через Prometheus. Отдельные гистограммы для длительности каждого шага.
- Distributed tracing: пробрасывайте
trace_idчерез все сервисы через контекст. OpenTelemetry + Jaeger дают полную картину прохождения саги через сервисы. - Dead Letter Queue (DLQ): в Redis Streams сообщения с многократными неудачами переносить в отдельный стрим для ручного анализа.
- Административный интерфейс: простой HTTP-эндпоинт, отдающий состояние саги по её ID из PostgreSQL — неоценим при инцидентах.
// monitoring.go
package orchestrator
import (
"log/slog"
"time"
"github.com/prometheus/client_golang/prometheus"
)
var (
sagaTotal = prometheus.NewCounterVec(
prometheus.CounterOpts{Name: "saga_total", Help: "Total saga executions"},
[]string{"type", "status"},
)
sagaDuration = prometheus.NewHistogramVec(
prometheus.HistogramOpts{
Name: "saga_duration_seconds",
Help: "Saga execution duration",
Buckets: prometheus.DefBuckets,
},
[]string{"type"},
)
)
func init() {
prometheus.MustRegister(sagaTotal, sagaDuration)
}
func recordSagaMetrics(sagaType string, status SagaStatus, start time.Time) {
sagaTotal.WithLabelValues(sagaType, string(status)).Inc()
sagaDuration.WithLabelValues(sagaType).Observe(time.Since(start).Seconds())
slog.Info("saga finished",
"type", sagaType,
"status", status,
"duration_ms", time.Since(start).Milliseconds(),
)
}
Типичные ошибки и антипаттерны
- Неидемпотентные обработчики. Самая частая ошибка. Если шаг саги не идемпотентен, повторный вызов создаст дублирующие операции — двойное списание, двойное резервирование. Всегда используйте
idempotency_key. - Отсутствие таймаутов. Шаг саги может зависнуть, ожидая ответа от недоступного сервиса. Каждый вызов должен иметь контекст с дедлайном. Используйте
context.WithTimeout. - Компенсация без ретраев. Если компенсирующая транзакция провалилась — это критическая ситуация. Нужны ретраи, алерты и, возможно, ручное вмешательство. Нельзя молча игнорировать ошибки компенсации.
- Слишком длинные саги. Сага из 10+ шагов — сигнал о проблемах в декомпозиции сервисов. Длинная сага означает долгое окно несогласованности данных и огромный граф компенсации.
- Читать uncommitted данные. Пока сага выполняется, данные в промежуточном состоянии. Другие сервисы, читающие их, могут получить некорректную картину. Проектируйте для eventual consistency явно.
- Хранить состояние саги только в памяти. После перезапуска сервиса все незавершённые саги теряются. Состояние должно персистироваться в PostgreSQL до начала каждого шага.
- Игнорировать порядок компенсации. Компенсировать нужно строго в обратном порядке выполненных шагов. Иначе можно создать новые несогласованности.
Заключение: как выбрать подход
Паттерн Saga решает реальную проблему распределённых транзакций в микросервисных системах на Go, но требует тщательного проектирования и дисциплины при реализации.
Saga не делает систему проще — она делает сложность явной и управляемой.
Выбирайте хореографию, если у вас 2–4 участника в саге, команда уже работает с event-driven подходом и Redis Streams, и вы готовы инвестировать в хорошую трассировку для отладки. Хореография хорошо масштабируется горизонтально.
Выбирайте оркестрацию, если сага содержит сложную бизнес-логику с ветвлениями, нужна полная видимость состояния в PostgreSQL, команда предпочитает явный контроль потока, а скорость отладки инцидентов критична.
Независимо от выбора: реализуйте идемпотентность с первого дня, добавьте распределённую трассировку до выхода в production, предусмотрите механизм восстановления зависших саг через периодическое сканирование БД и никогда не игнорируйте ошибки компенсирующих транзакций. Go с его явной обработкой ошибок и мощными примитивами конкурентности — отличный выбор для реализации обоих подходов.
Технологии
Теги
Руслан Исмаилов
Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →