HardПрактика12 min

Saga pattern

Orchestration vs Choreography, компенсирующие транзакции, state machine, идемпотентность

Проблема: распределённая транзакция

Оформление заказа затрагивает 4 сервиса:

  1. Order Service -- создать заказ в БД
  2. Payment Service -- списать деньги (Stripe)
  3. Inventory Service -- зарезервировать товар
  4. Shipping Service -- создать задачу на доставку

Каждый сервис имеет свою БД. 2PC (two-phase commit) здесь не годится:

  • 2PC блокирует ресурсы на время всего протокола
  • Координатор -- single point of failure
  • Не работает через HTTP/gRPC без специальной инфраструктуры
  • Stripe API вообще не поддерживает 2PC

Saga -- это последовательность локальных транзакций. Если шаг падает -- выполняем компенсирующие транзакции для откатки предыдущих шагов.

Forward:    T1 -> T2 -> T3 -> T4  (success)
Compensate: T1 -> T2 -> T3 -> X (T4 failed)
            C1 <- C2 <- C3      (rollback)

Два стиля: Orchestration vs Choreography

Orchestration

Центральный координатор (orchestrator) командует всеми шагами. Каждый сервис просто выполняет команду и отвечает.

[Orchestrator] --> CreatePayment --> [Payment]
     <----------- PaymentCreated ----------
[Orchestrator] --> ReserveStock  --> [Inventory]
     <----------- StockReserved  ----------
[Orchestrator] --> CreateShipping --> [Shipping]

Плюсы: явный flow, легко дебажить, явные таймауты на каждый шаг. Минусы: orchestrator = single point of coupling, его надо HA.

Choreography

Сервисы общаются через события. Каждый сам решает, что делать.

[Order] -- OrderCreated --> [Payment]
[Payment] -- PaymentCreated --> [Inventory]
[Inventory] -- StockReserved --> [Shipping]

Плюсы: декомпозиция, нет центральной точки. Минусы: flow неявный, сложно отследить, легко создать циклические зависимости.

Сравнение

Критерий Orchestration Choreography
Видимость flow Высокая (в одном месте) Низкая (размазан)
Связность Высокая (все знают orchestrator) Низкая (только события)
Отладка Проще (trace в одном сервисе) Сложнее (трассировка через шину)
Изменение flow В одном месте Надо менять всех
Количество сервисов Хорошо при >4 Хорошо при 2-3
Риск инфинит-лупов Низкий Высокий

Правило: при >4 участниках -- orchestration. При 2-3 -- choreography.

Saga vs 2PC

Свойство 2PC Saga
Атомарность ACID Eventual
Блокировки Ресурсы заблокированы Нет блокировок
Откат Автоматический Через компенсации
Латентность Высокая (2 round trips) Низкая (pipeline)
Доступность Низкая (все живы) Высокая
Сложность кода Низкая Высокая (компенсации)
Работает через HTTP Нет Да

State machine саги

Orchestrator -- по сути конечный автомат:

           CreatePayment              ReserveStock
  START ─────────────► PAID ────────────────► RESERVED
    │                   │                        │
    │                   │                        │ CreateShipping
    │                   │                        ▼
    │                   │                   SHIPPED
    │                   │                        │
    │                   │                        ▼
    │                   │                    COMPLETED
    │                   │
    │           RefundPayment
    │                   │
    ▼                   ▼
 FAILED  <----------- COMPENSATING
              ReleaseStock

Реализация: асинхронный оркестратор на шине сообщений

<?php

declare(strict_types=1);

namespace App\Saga\Order;

/**
 * Persistent saga state stored in DB.
 * Each step result is recorded for idempotency and recovery.
 */
enum OrderSagaStatus: string
{
    case Started = 'started';
    case PaymentCompleted = 'payment_completed';
    case StockReserved = 'stock_reserved';
    case ShippingCreated = 'shipping_created';
    case Completed = 'completed';
    case Compensating = 'compensating';
    case Failed = 'failed';
}

final class OrderSagaState
{
    public function __construct(
        public readonly string $sagaId,
        public readonly string $orderId,
        public OrderSagaStatus $status,
        public ?string $paymentId = null,
        public ?string $reservationId = null,
        public ?string $shipmentId = null,
        public ?string $failureReason = null,
        public int $version = 0,
    ) {}
}
## Реализация: синхронный конечный автомат с resume
package saga

import (
	"context"
	"errors"
	"fmt"
	"time"
)

// Status represents the saga state machine.
type Status string

const (
	StatusStarted         Status = "started"
	StatusPaid            Status = "paid"
	StatusStockReserved   Status = "stock_reserved"
	StatusShippingCreated Status = "shipping_created"
	StatusCompleted       Status = "completed"
	StatusCompensating    Status = "compensating"
	StatusFailed          Status = "failed"
)

// State persists saga progress in DB.
type State struct {
	SagaID         string
	OrderID        string
	Status         Status
	PaymentID      string
	ReservationID  string
	ShipmentID     string
	FailureReason  string
	Version        int
	UpdatedAt      time.Time
}

type Orchestrator struct {
	repo      StateRepo
	payment   PaymentClient
	inventory InventoryClient
	shipping  ShippingClient
}

// Run executes the saga forward. On failure, compensations are triggered.
// All steps are idempotent via keys derived from sagaID.
func (o *Orchestrator) Run(ctx context.Context, sagaID, orderID string, amount int64) error {
	state, err := o.repo.Load(ctx, sagaID)
	if err != nil {
		return fmt.Errorf("load state: %w", err)
	}

	// Resume from wherever we left off (recovery after crash)
	switch state.Status {
	case StatusStarted:
		if err := o.stepCharge(ctx, state, amount); err != nil {
			return o.compensate(ctx, state, err)
		}
		fallthrough
	case StatusPaid:
		if err := o.stepReserve(ctx, state, orderID); err != nil {
			return o.compensate(ctx, state, err)
		}
		fallthrough
	case StatusStockReserved:
		if err := o.stepShipping(ctx, state, orderID); err != nil {
			return o.compensate(ctx, state, err)
		}
		state.Status = StatusCompleted
		return o.repo.Save(ctx, state)
	case StatusCompleted:
		return nil
	case StatusFailed:
		return errors.New("saga already failed")
	}
	return nil
}

func (o *Orchestrator) stepCharge(ctx context.Context, s *State, amount int64) error {
	pid, err := o.payment.Charge(ctx, s.SagaID+":payment", amount)
	if err != nil {
		return fmt.Errorf("charge: %w", err)
	}
	s.PaymentID = pid
	s.Status = StatusPaid
	return o.repo.Save(ctx, s)
}

func (o *Orchestrator) stepReserve(ctx context.Context, s *State, orderID string) error {
	rid, err := o.inventory.Reserve(ctx, s.SagaID+":stock", orderID)
	if err != nil {
		return fmt.Errorf("reserve: %w", err)
	}
	s.ReservationID = rid
	s.Status = StatusStockReserved
	return o.repo.Save(ctx, s)
}

func (o *Orchestrator) stepShipping(ctx context.Context, s *State, orderID string) error {
	sid, err := o.shipping.Create(ctx, s.SagaID+":shipping", orderID)
	if err != nil {
		return fmt.Errorf("shipping: %w", err)
	}
	s.ShipmentID = sid
	s.Status = StatusShippingCreated
	return o.repo.Save(ctx, s)
}

// compensate runs reverse actions for completed steps.
// Each compensation must be idempotent via idempotency keys.
func (o *Orchestrator) compensate(ctx context.Context, s *State, cause error) error {
	s.Status = StatusCompensating
	s.FailureReason = cause.Error()
	_ = o.repo.Save(ctx, s)

	if s.ReservationID != "" {
		if err := o.inventory.Release(ctx, s.SagaID+":release", s.ReservationID); err != nil {
			// Log and retry via outbox; do not abort compensation chain
			_ = err
		}
	}
	if s.PaymentID != "" {
		if err := o.payment.Refund(ctx, s.SagaID+":refund", s.PaymentID); err != nil {
			_ = err
		}
	}

	s.Status = StatusFailed
	return o.repo.Save(ctx, s)
}
## Recovery после краша

Если orchestrator упал между шагами, нужен фоновый процесс, который догоняет зависшие саги:

-- Find stuck sagas: not updated for >5 minutes and not in terminal state
SELECT saga_id, status, updated_at
FROM saga_state
WHERE status NOT IN ('completed', 'failed')
  AND updated_at < NOW() - INTERVAL '5 minutes'
ORDER BY updated_at ASC
LIMIT 100;

Для каждой такой саги worker вызывает orchestrator.Run(sagaID, ...) -- state machine резюмирует с последнего сохранённого шага благодаря fallthrough.

Идемпотентность -- абсолютное требование

Каждый шаг может быть вызван дважды из-за retry:

  1. Ключ идемпотентности передаётся во все внешние API: Stripe, inventory, shipping
  2. Локальная запись в БД защищена UNIQUE constraint на (saga_id, step)
  3. Компенсация идемпотентна: release несуществующей резервации -- OK
CREATE TABLE saga_step_log (
    saga_id    UUID NOT NULL,
    step_name  TEXT NOT NULL,
    result_id  TEXT,  -- external ID (payment_id, etc)
    completed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    PRIMARY KEY (saga_id, step_name)
);

Перед каждым шагом проверяем: уже есть запись? Значит шаг уже выполнен, возвращаем сохранённый result_id.

Таймауты и дедлайны

Saga не должна висеть вечно. Настройки:

Параметр Значение Причина
Per-step timeout 10-30 секунд Сохранить latency
Saga overall timeout 5-15 минут Предсказуемость для UX
Retry per step 3 попытки После -- компенсация
Backoff exponential + jitter См. 10.resilience-deep.md

По истечении overall timeout -- принудительная компенсация.

Мониторинг

Ключевые метрики:

  • saga_started_total (counter) -- количество запущенных
  • saga_completed_total / saga_failed_total -- исходы
  • saga_duration_seconds (histogram) -- p50/p95/p99
  • saga_step_duration_seconds -- по каждому шагу
  • saga_stuck_count (gauge) -- висящих >5 мин
  • saga_compensation_total -- частота компенсаций по причинам

Алерты: rate(saga_failed_total[5m]) > 0.05, saga_stuck_count > 10.

Choreography: пример на событиях

Если всё же choreography -- каждый сервис реагирует на событие и публикует следующее:

Order -> OrderCreated
  Payment подписан -> charges -> PaymentSucceeded | PaymentFailed
    Inventory подписан на PaymentSucceeded -> reserve -> StockReserved | OutOfStock
      OutOfStock -> триггерит RefundPaymentCommand
      StockReserved -> Shipping подписан -> create shipment

Каждый сервис хранит свой маленький state machine. Надо явно определять compensation topic'и.

Pitfalls

  • Необратимые операции: отправка email, физическая доставка. Компенсация невозможна -- нужен forward recovery (повторить и зафиксировать).
  • Cycle в choreography: A выпускает событие, которое слушает B, который выпускает событие, которое слушает A. Надо рисовать DAG.
  • Забытая компенсация: добавили шаг 5, а compensation 5 не написали -- в фейле зависнет в compensating.
  • Не-идемпотентный шаг: два charge за одну сагу -- списали дважды. Всегда idempotency key.
  • Компенсация зависит от данных шага: не сохранили payment_id -- не можем сделать refund. Сохраняйте в state каждое внешнее ID.

Выводы

  • Saga -- замена 2PC для распределённых транзакций через eventual consistency
  • Две реализации: orchestration (одна "дирижёрская" функция) и choreography (события)
  • Критично: идемпотентность каждого шага и компенсации, persistent state, recovery после краша
  • При >4 участниках -- orchestration, чтобы flow не размазался
  • Компенсации всегда сложнее прямых шагов: учтите необратимость (email sent, shipment dispatched)
  • Мониторинг saga_stuck и saga_failed обязателен