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).
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: сервисы реагируют на события
В хореографической саге нет центрального координатора. Каждый сервис слушает события и реагирует на них, публикуя новые события:
Каждый сервис знает только о своих событиях и о том, на какие чужие события реагировать. Нет единой точки, которая видит весь процесс.
Orchestration vs Choreography
| Критерий | Orchestration | Choreography |
|---|---|---|
| Видимость процесса | Весь 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 секунд на шаг. Для внешних платёжных систем может быть больше. Главное - не ждать бесконечно: застрявшая сага блокирует ресурсы.
Компенсации: правила проектирования
Компенсация - это не просто DELETE или undo. Несколько правил:
- Компенсация должна быть идемпотентной. Она может быть вызвана несколько раз.
- Компенсация не должна падать. Если упала - логируем и продолжаем компенсировать остальные шаги. Иначе получим «застрявшую» сагу.
- Компенсация - это семантический откат. Не обязательно удалять данные. Можно пометить заказ как
cancelled, вместо удаления строки. - Не все шаги требуют компенсации. Если шаг только читает данные (валидация) - компенсировать нечего.
Типичные ошибки
- 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после каждого шага