Saga: как оркестрировать распределённые процессы

Saga: как оркестрировать распределённые процессы

Когда бизнес-процесс затрагивает несколько сервисов, возникает вопрос: как обеспечить согласованность? В монолите ответ простой - транзакция БД. В микросервисной архитектуре каждый сервис владеет своей базой, и общей транзакции нет. Two-Phase Commit (2PC) технически возможен, но на практике не масштабируется, держит блокировки слишком долго и создаёт single point of failure в координаторе.

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

Анатомия Saga

Saga состоит из шагов. Каждый шаг - это пара: действие (Execute) и компенсация (Compensate). Если шаг N упал, мы последовательно откатываем шаги N-1, N-2, ..., 1, вызывая их компенсации.

Пример - оформление заказа с оплатой и выдачей доступа:

Шаг 1: Создать заказ (status=pending)
  └─ Компенсация: Отменить заказ (status=cancelled)

Шаг 2: Списать оплату
  └─ Компенсация: Вернуть деньги (refund)

Шаг 3: Выдать доступ к курсу
  └─ Компенсация: Отозвать доступ

Если выдача доступа (шаг 3) провалилась - возвращаем деньги (компенсация шага 2) и отменяем заказ (компенсация шага 1).

Saga: happy path - все шаги; failure path - откат назад через компенсации

Orchestration Saga: центральный координатор

В оркестрационной саге один компонент (оркестратор) управляет всеми шагами. Он знает порядок выполнения и запускает компенсации при сбое.

Структура шага и оркестратор

// SagaStep описывает один шаг саги: действие и компенсацию.
type SagaStep struct {
    Name       string
    Execute    func(ctx context.Context) error
    Compensate func(ctx context.Context) error
}

// SagaOrchestrator выполняет шаги по порядку.
// При ошибке откатывает уже выполненные шаги в обратном порядке.
type SagaOrchestrator struct {
    steps []SagaStep
}

// NewSaga создаёт оркестратор с заданными шагами.
func NewSaga(steps ...SagaStep) *SagaOrchestrator {
    return &SagaOrchestrator{steps: steps}
}

// Run выполняет все шаги. При сбое вызывает компенсации.
func (s *SagaOrchestrator) Run(ctx context.Context) error {
    executed := make([]SagaStep, 0, len(s.steps))

    for _, step := range s.steps {
        slog.InfoContext(ctx, "saga step executing",
            slog.String("step", step.Name),
        )

        if err := step.Execute(ctx); err != nil {
            slog.ErrorContext(ctx, "saga step failed, starting compensation",
                slog.String("step", step.Name),
                slog.String("err", err.Error()),
            )
            s.compensate(ctx, executed)
            return fmt.Errorf("saga failed at step %q: %w", step.Name, err)
        }

        executed = append(executed, step)
    }

    slog.InfoContext(ctx, "saga completed successfully",
        slog.Int("steps", len(s.steps)),
    )
    return nil
}

// compensate откатывает выполненные шаги в обратном порядке.
func (s *SagaOrchestrator) compensate(ctx context.Context, executed []SagaStep) {
    for i := len(executed) - 1; i >= 0; i-- {
        step := executed[i]
        slog.InfoContext(ctx, "compensating step",
            slog.String("step", step.Name),
        )

        if err := step.Compensate(ctx); err != nil {
            // Компенсация не должна падать, но если упала - логируем и продолжаем
            slog.ErrorContext(ctx, "compensation failed",
                slog.String("step", step.Name),
                slog.String("err", err.Error()),
            )
        }
    }
}
<?php
// src/Application/Saga/SagaStep.php
declare(strict_types=1);

namespace App\Application\Saga;

// Один шаг саги: действие + компенсация. Immutable - readonly properties.
final readonly class SagaStep
{
    /**
     * @param callable():void $execute
     * @param callable():void $compensate
     */
    public function __construct(
        public string $name,
        public \Closure $execute,
        public \Closure $compensate,
    ) {}
}

// src/Application/Saga/SagaOrchestrator.php
final class SagaOrchestrator
{
    public function __construct(
        private readonly LoggerInterface $logger,
    ) {}

    /** @param list<SagaStep> $steps */
    public function run(array $steps): void
    {
        /** @var list<SagaStep> $executed */
        $executed = [];

        foreach ($steps as $step) {
            $this->logger->info('saga step executing', ['step' => $step->name]);

            try {
                ($step->execute)();
            } catch (\Throwable $e) {
                $this->logger->error('saga step failed, starting compensation', [
                    'step' => $step->name,
                    'err' => $e->getMessage(),
                ]);
                $this->compensate($executed);
                throw new SagaFailedException(sprintf('saga failed at step %s', $step->name), 0, $e);
            }

            $executed[] = $step;
        }

        $this->logger->info('saga completed successfully', ['steps' => count($steps)]);
    }

    /** @param list<SagaStep> $executed */
    private function compensate(array $executed): void
    {
        foreach (array_reverse($executed) as $step) {
            $this->logger->info('compensating step', ['step' => $step->name]);
            try {
                ($step->compensate)();
            } catch (\Throwable $e) {
                $this->logger->error('compensation failed', [
                    'step' => $step->name,
                    'err' => $e->getMessage(),
                ]);
            }
        }
    }
}

Реальный пример - CreateOrder Saga

// CreateOrderSaga оформляет заказ: резервирует товар, списывает оплату, выдаёт доступ.
func CreateOrderSaga(
    orderSvc OrderService,
    paymentSvc PaymentService,
    accessSvc AccessService,
    req CreateOrderRequest,
) *SagaOrchestrator {
    var orderID int64

    return NewSaga(
        SagaStep{
            Name: "create_order",
            Execute: func(ctx context.Context) error {
                id, err := orderSvc.Create(ctx, req.UserID, req.CourseID)
                if err != nil {
                    return err
                }
                orderID = id
                return nil
            },
            Compensate: func(ctx context.Context) error {
                return orderSvc.Cancel(ctx, orderID)
            },
        },
        SagaStep{
            Name: "charge_payment",
            Execute: func(ctx context.Context) error {
                return paymentSvc.Charge(ctx, req.UserID, req.Amount, orderID)
            },
            Compensate: func(ctx context.Context) error {
                return paymentSvc.Refund(ctx, orderID)
            },
        },
        SagaStep{
            Name: "grant_access",
            Execute: func(ctx context.Context) error {
                return accessSvc.GrantCourseAccess(ctx, req.UserID, req.CourseID)
            },
            Compensate: func(ctx context.Context) error {
                return accessSvc.RevokeCourseAccess(ctx, req.UserID, req.CourseID)
            },
        },
    )
}

// Использование:
// saga := CreateOrderSaga(orderSvc, paymentSvc, accessSvc, req)
// if err := saga.Run(ctx); err != nil {
//     // Saga откатилась, все компенсации выполнены
//     return fmt.Errorf("order creation failed: %w", err)
// }
<?php
// src/Application/Saga/CreateOrderSaga.php
declare(strict_types=1);

namespace App\Application\Saga;

use App\Application\Port\AccessService;
use App\Application\Port\OrderService;
use App\Application\Port\PaymentService;

final class CreateOrderSaga
{
    public function __construct(
        private readonly OrderService $orders,
        private readonly PaymentService $payments,
        private readonly AccessService $access,
        private readonly SagaOrchestrator $orchestrator,
    ) {}

    public function run(CreateOrderRequest $req): void
    {
        // Захватываем orderId через by-reference замыкания
        $orderId = null;

        $this->orchestrator->run([
            new SagaStep(
                name: 'create_order',
                execute: function () use ($req, &$orderId): void {
                    $orderId = $this->orders->create($req->userId, $req->courseId);
                },
                compensate: function () use (&$orderId): void {
                    if ($orderId !== null) {
                        $this->orders->cancel($orderId);
                    }
                },
            ),
            new SagaStep(
                name: 'charge_payment',
                execute: function () use ($req, &$orderId): void {
                    $this->payments->charge($req->userId, $req->amount, $orderId);
                },
                compensate: function () use (&$orderId): void {
                    $this->payments->refund($orderId);
                },
            ),
            new SagaStep(
                name: 'grant_access',
                execute: function () use ($req): void {
                    $this->access->grantCourseAccess($req->userId, $req->courseId);
                },
                compensate: function () use ($req): void {
                    $this->access->revokeCourseAccess($req->userId, $req->courseId);
                },
            ),
        ]);
    }
}

Choreography Saga: сервисы реагируют на события

В хореографической саге нет центрального координатора. Каждый сервис слушает события и реагирует на них, публикуя новые события:

Choreography saga: три сервиса обмениваются событиями; happy path - OrderCreated и PaymentCharged идут вперёд, failure path - AccessFailed и PaymentRefunded возвращаются обратно как компенсации

Каждый сервис знает только о своих событиях и о том, на какие чужие события реагировать. Нет единой точки, которая видит весь процесс.

Orchestration vs Choreography

КритерийOrchestrationChoreography
Видимость процессаВесь flow в одном местеРазбросан по сервисам
СвязанностьОркестратор знает обо всехСервисы знают только о соседних событиях
ОтладкаПроще: смотришь оркестраторСложнее: восстанавливаешь по логам
МасштабированиеОркестратор может стать bottleneckЛучше масштабируется
Когда использовать3-7 шагов, сложная логика2-3 шага, простые реакции

Для большинства проектов начинайте с оркестрации. Переходите к хореографии, когда оркестратор становится bottleneck или когда сервисы принадлежат разным командам.

Persistence: сохранение состояния саги

Если оркестратор упадёт между шагами, нужно знать, на каком шаге остановились. Для этого состояние саги сохраняется в БД:

// SagaState хранит текущее состояние саги в БД.
type SagaState struct {
    ID            string    `json:"id"`
    Type          string    `json:"type"`            // "create_order"
    CurrentStep   int       `json:"current_step"`    // индекс текущего шага
    Status        string    `json:"status"`          // running, completed, compensating, failed
    Payload       []byte    `json:"payload"`         // входные данные саги
    CompletedSteps []string `json:"completed_steps"` // имена выполненных шагов
    CreatedAt     time.Time `json:"created_at"`
    UpdatedAt     time.Time `json:"updated_at"`
}

// SaveState сохраняет состояние саги после каждого шага.
func (r *SagaRepo) SaveState(ctx context.Context, state SagaState) error {
    _, err := r.db.ExecContext(ctx,
        `INSERT INTO saga_states (id, type, current_step, status, payload, completed_steps, created_at, updated_at)
         VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
         ON CONFLICT (id) DO UPDATE SET
            current_step = EXCLUDED.current_step,
            status = EXCLUDED.status,
            completed_steps = EXCLUDED.completed_steps,
            updated_at = EXCLUDED.updated_at`,
        state.ID, state.Type, state.CurrentStep, state.Status,
        state.Payload, pq.Array(state.CompletedSteps),
        state.CreatedAt, time.Now(),
    )
    return err
}
<?php
// src/Infrastructure/Saga/SagaState.php
declare(strict_types=1);

namespace App\Infrastructure\Saga;

// Snapshot состояния саги. Immutable - вместо мутации создаём новую версию.
final readonly class SagaState
{
    /** @param list<string> $completedSteps */
    public function __construct(
        public string $id,
        public string $type,            // 'create_order'
        public int $currentStep,        // индекс текущего шага
        public string $status,          // running, completed, compensating, failed
        public string $payload,         // JSON-строка - входные данные саги
        public array $completedSteps,
        public \DateTimeImmutable $createdAt,
        public \DateTimeImmutable $updatedAt,
    ) {}
}

// src/Infrastructure/Saga/SagaRepository.php
namespace App\Infrastructure\Saga;

use Doctrine\DBAL\Connection;

final class SagaRepository
{
    public function __construct(
        private readonly Connection $connection,
    ) {}

    public function saveState(SagaState $state): void
    {
        // jsonb[]/text[] для completedSteps в Postgres - удобно для запросов.
        $this->connection->executeStatement(
            'INSERT INTO saga_states (id, type, current_step, status, payload, completed_steps, created_at, updated_at)
             VALUES (:id, :type, :step, :status, :payload, :completed, :created, :updated)
             ON CONFLICT (id) DO UPDATE SET
                current_step = EXCLUDED.current_step,
                status = EXCLUDED.status,
                completed_steps = EXCLUDED.completed_steps,
                updated_at = EXCLUDED.updated_at',
            [
                'id' => $state->id,
                'type' => $state->type,
                'step' => $state->currentStep,
                'status' => $state->status,
                'payload' => $state->payload,
                'completed' => '{' . implode(',', $state->completedSteps) . '}',
                'created' => $state->createdAt->format('Y-m-d H:i:s'),
                'updated' => $state->updatedAt->format('Y-m-d H:i:s'),
            ],
        );
    }
}

При старте приложения - находим саги в статусе running или compensating и возобновляем их с нужного шага.

Timeout: что если шаг не отвечает

В распределённой системе шаг может зависнуть. Каждый шаг должен иметь таймаут:

// executeWithTimeout выполняет шаг с таймаутом.
func executeWithTimeout(ctx context.Context, step SagaStep, timeout time.Duration) error {
    ctx, cancel := context.WithTimeout(ctx, timeout)
    defer cancel()

    done := make(chan error, 1)
    go func() {
        done <- step.Execute(ctx)
    }()

    select {
    case err := <-done:
        return err
    case <-ctx.Done():
        return fmt.Errorf("step %q timed out after %v", step.Name, timeout)
    }
}
<?php
// src/Application/Saga/TimeoutExecutor.php
declare(strict_types=1);

namespace App\Application\Saga;

// В PHP нет goroutines, и таймаут принудительно из основного потока выставить
// нельзя. Реальный таймаут прокидывается внутрь HTTP/DB-клиента шага.
// На уровне саги - проверка длительности через Stopwatch + бросок исключения.
final class TimeoutExecutor
{
    public function __construct(
        private readonly int $timeoutSeconds,
    ) {}

    public function execute(SagaStep $step): void
    {
        $startedAt = hrtime(true);

        try {
            ($step->execute)();
        } catch (\Throwable $e) {
            // Прокидываем как есть - external client уже бросает на своём deadline
            throw $e;
        }

        $elapsedSeconds = (hrtime(true) - $startedAt) / 1_000_000_000;
        if ($elapsedSeconds > $this->timeoutSeconds) {
            // Шаг завершился, но за пределами SLA - бросаем для откатов выше.
            // Реальная отмена должна быть прокинута в каждый внешний вызов:
            // - Symfony HttpClient: ['timeout' => 30]
            // - Doctrine: statement_timeout в Postgres
            // - Symfony Lock: ttl
            throw new SagaStepTimedOutException(sprintf(
                'step %s exceeded SLA: %.2fs > %ds',
                $step->name,
                $elapsedSeconds,
                $this->timeoutSeconds,
            ));
        }
    }
}

Таймаут по умолчанию - 30 секунд на шаг. Для внешних платёжных систем может быть больше. Главное - не ждать бесконечно: застрявшая сага блокирует ресурсы.

Между шагами саги другие процессы могут читать промежуточное состояние (заказ создан, но не оплачен). Это нормально - saga обеспечивает eventual consistency, а не ACID. Каждый шаг и каждая компенсация должны быть идемпотентными, потому что при восстановлении после сбоя они могут быть вызваны повторно.

Компенсации: правила проектирования

Компенсация - это не просто DELETE или undo. Несколько правил:

  1. Компенсация должна быть идемпотентной. Она может быть вызвана несколько раз.
  2. Компенсация не должна падать. Если упала - логируем и продолжаем компенсировать остальные шаги. Иначе получим «застрявшую» сагу.
  3. Компенсация - это семантический откат. Не обязательно удалять данные. Можно пометить заказ как cancelled, вместо удаления строки.
  4. Не все шаги требуют компенсации. Если шаг только читает данные (валидация) - компенсировать нечего.

Типичные ошибки

  • Saga вместо распределённой транзакции «на всякий» - для простого случая (БД + Redis) часто хватает Outbox + idempotency без оркестратора. Saga нужна когда 3+ шага, каждый с побочным эффектом во внешней системе.
  • Компенсация = DELETE строк - другие подписчики уже прочитали данные, ссылаются на них. Лучше - пометка status: cancelled, событие OrderCancelled. Логический откат, не физический.
  • Состояние саги только в памяти - процесс перезапустился на середине, нет ни данных «где мы», ни компенсаций. Сохраняй состояние саги в БД после каждого шага + при старте находи незавершённые и продолжай.
  • Компенсации падают тихо - compensate() поймала ошибку, залогировала и пошла дальше; средний шаг застрял в «оплачено, но не выдано». Каждый сбой компенсации = алёрт оператора, не просто запись в лог.
  • Choreography с цепочкой длиннее 3 событий - OrderCreated → InventoryReserved → PaymentProcessed → AccessGranted → ... Без центрального state почти невозможно понять, на каком шаге saga застряла. Длинные потоки лучше делать Orchestration.
  • Шаг не идемпотентный - Saga при retry повторяет шаг, который начисляет деньги - двойное списание. Каждый шаг (как и компенсация) обязан быть идемпотентным.
  • Timeout не пробрасывается во внешний вызов - select { case <-ctx.Done() } сработал, но HTTP-запрос в платёжку всё равно сидит в сокете. Compensation вызвана, но платёж мог пройти. Передавай ctx в HTTP-клиент, чтобы запрос реально отменился.
  • Saga + долгий шаг ожидания (SleepUntil несколько дней) - оркестратор держит горутину/коннект к БД сутками. Используй durable scheduler (Temporal/Cadence) или persisted Saga с wake-up через cron.

Мини-задание

  • Реализуй SagaOrchestrator с методами Run и compensate по примеру из урока
  • Опиши процесс из 3 шагов (создание заказа, оплата, выдача доступа) и компенсацию к каждому шагу
  • Добавь таймаут 10 секунд на каждый шаг с помощью context.WithTimeout
  • Напиши unit-тест: при сбое на шаге 3 компенсации шагов 2 и 1 должны быть вызваны в обратном порядке
  • Добавь сохранение состояния саги в таблицу saga_states после каждого шага

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.