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: центральный координатор
В оркестрационной саге один компонент (оркестратор) управляет всеми шагами. Он знает порядок выполнения и запускает компенсации при сбое.
Структура шага и оркестратор
// Step - один шаг саги: действие и его компенсация.
//
// Обе функции получают транзакцию: шаг обязан записывать свой эффект и отметку
// о себе одной командой COMMIT. Иначе «состояние саги в базе» разъезжается
// с реальностью ровно так же, как отметка в inbox с бизнес-операцией.
type Step struct {
Name string
Execute func(ctx context.Context, tx *sql.Tx) error
Compensate func(ctx context.Context, tx *sql.Tx) error
}
// Orchestrator выполняет шаги по порядку, при сбое откатывает выполненные.
type Orchestrator struct {
db *sql.DB
log *slog.Logger
steps []Step
}
func New(db *sql.DB, log *slog.Logger, steps ...Step) *Orchestrator {
return &Orchestrator{db: db, log: log, steps: steps}
}
// Run выполняет сагу с идентификатором sagaID.
//
// Идентификатор нужен не для логов: он ключ идемпотентности. Повторный запуск
// той же саги не должен повторять уже выполненные шаги, а восстановление
// после перезапуска процесса - находить их в базе, а не в памяти.
func (o *Orchestrator) Run(ctx context.Context, sagaID string) error {
for _, step := range o.steps {
if err := o.runStep(ctx, sagaID, step); err != nil {
o.log.ErrorContext(ctx, "saga step failed, compensating",
slog.String("saga_id", sagaID),
slog.String("step", step.Name),
slog.String("err", err.Error()),
)
// Компенсации идут отдельной операцией и **возвращают** ошибку.
// Проглотить её значит оставить сагу в промежуточном состоянии
// молча - «оплачено, но не выдано» без единого сигнала.
if cErr := o.Compensate(ctx, sagaID); cErr != nil {
return fmt.Errorf("шаг %q упал (%w), откат не завершён: %w",
step.Name, err, cErr)
}
return fmt.Errorf("сага %s: шаг %q: %w", sagaID, step.Name, err)
}
}
return nil
}
// runStep выполняет один шаг: действие и отметку о нём одной транзакцией.
func (o *Orchestrator) runStep(ctx context.Context, sagaID string, step Step) error {
tx, err := o.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
if err := step.Execute(ctx, tx); err != nil {
return err
}
// Отметка в том же COMMIT, что и эффект шага. Порядок здесь не важен -
// важно, что они в одной транзакции: иначе после сбоя журнал скажет
// «шаг сделан», а эффекта не будет (или наоборот).
if err := markStep(ctx, tx, sagaID, step.Name, faultsql.StepDone); err != nil {
return err
}
return tx.Commit()
}
// Compensate откатывает все шаги саги, помеченные выполненными.
//
// Отдельный экспортируемый метод, а не приватная часть Run: после
// перезапуска процесса откат должен продолжиться, и запускать его будет
// внешний восстановитель, у которого нет ни `executed`, ни стека вызовов -
// только sagaID и база.
func (o *Orchestrator) Compensate(ctx context.Context, sagaID string) error {
done, err := o.doneSteps(ctx, sagaID)
if err != nil {
return err
}
var failed []string
// Обратный порядок: шаги откатываются от последнего к первому.
for i := len(o.steps) - 1; i >= 0; i-- {
step := o.steps[i]
if !done[step.Name] {
continue
}
if err := o.compensateStep(ctx, sagaID, step); err != nil {
// Не прекращаем откат: остальные шаги тоже надо попытаться
// откатить. Но и не забываем - список уйдёт вызывающему.
o.log.ErrorContext(ctx, "compensation failed",
slog.String("saga_id", sagaID),
slog.String("step", step.Name),
slog.String("err", err.Error()),
)
failed = append(failed, step.Name)
}
}
if len(failed) > 0 {
return fmt.Errorf("%w: шаги %v требуют вмешательства", ErrCompensationFailed, failed)
}
return nil
}
// compensateStep откатывает шаг и обновляет журнал одной транзакцией.
func (o *Orchestrator) compensateStep(ctx context.Context, sagaID string, step Step) error {
tx, err := o.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
if err := step.Compensate(ctx, tx); err != nil {
// Пометку о неудаче пишем **отдельной** транзакцией: та, в которой
// упала компенсация, откатится и унесёт пометку с собой.
o.markCompensationFailed(ctx, sagaID, step.Name)
return err
}
if err := markStep(ctx, tx, sagaID, step.Name, faultsql.StepCompensated); err != nil {
return err
}
return tx.Commit()
}
// markCompensationFailed фиксирует, что откат шага не удался.
//
// Ошибку записи здесь только логируем: если и она не прошла, сага уже
// в состоянии, требующем человека, а вернуть мы всё равно должны исходную
// ошибку компенсации - она информативнее.
func (o *Orchestrator) markCompensationFailed(ctx context.Context, sagaID, step string) {
tx, err := o.db.BeginTx(ctx, nil)
if err != nil {
o.log.ErrorContext(ctx, "cannot record compensation failure",
slog.String("saga_id", sagaID), slog.String("step", step),
slog.String("err", err.Error()))
return
}
defer tx.Rollback()
if err := markStep(ctx, tx, sagaID, step, faultsql.StepCompensationFailed); err != nil {
o.log.ErrorContext(ctx, "cannot record compensation failure",
slog.String("saga_id", sagaID), slog.String("step", step),
slog.String("err", err.Error()))
return
}
if err := tx.Commit(); err != nil {
o.log.ErrorContext(ctx, "cannot commit compensation failure",
slog.String("saga_id", sagaID), slog.String("step", step),
slog.String("err", err.Error()))
}
}
// doneSteps читает из базы шаги, которые выполнены и ещё не откачены.
func (o *Orchestrator) doneSteps(ctx context.Context, sagaID string) (map[string]bool, error) {
rows, err := o.db.QueryContext(ctx,
`SELECT step FROM saga_steps WHERE saga_id = $1 AND state = $2`,
sagaID, string(faultsql.StepDone),
)
if err != nil {
return nil, fmt.Errorf("query saga steps: %w", err)
}
defer rows.Close()
done := map[string]bool{}
for rows.Next() {
var step string
if err := rows.Scan(&step); err != nil {
return nil, fmt.Errorf("scan step: %w", err)
}
done[step] = true
}
return done, rows.Err()
}
// markStep пишет состояние шага. UPSERT, потому что шаг может быть записан
// повторно: retry, восстановление после перезапуска, повторная компенсация.
func markStep(ctx context.Context, tx *sql.Tx, sagaID, step string, st faultsql.StepState) error {
_, err := tx.ExecContext(ctx,
`INSERT INTO saga_steps (saga_id, step, state) VALUES ($1, $2, $3)
ON CONFLICT (saga_id, step) DO UPDATE SET state = excluded.state`,
sagaID, step, string(st),
)
if err != nil {
return fmt.Errorf("mark step %q as %s: %w", step, st, err)
}
return nil
}
<?php
// src/Application/Saga/SagaStep.php
declare(strict_types=1);
namespace App\Application\Saga;
use Closure;
use Throwable;
// Один шаг саги: действие + компенсация. 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);
},
),
]);
}
}
1. Состояние саги жило в памяти. Список выполненных шагов был локальной
переменной executed. Процесс перезапустился на середине - и знания «что
надо откатить» не осталось нигде. Раздел типичных ошибок требует «сохраняй
состояние саги в БД после каждого шага», код этого не делал. Теперь
Compensate - отдельный метод, которому нужен только sagaID: он читает
журнал из базы, поэтому откат может продолжить любой другой процесс.
2. Сбой компенсации проглатывался. compensate логировала ошибку
и шла дальше, ничего не возвращая. Снаружи это выглядело как «сага
откатилась», хотя средний шаг застрял в «оплачено, но не выдано». Тот же
урок пишет: «Каждый сбой компенсации = алёрт оператора, не просто запись
в лог». Теперь ошибка возвращается вызывающему вместе со списком застрявших
шагов, а сами шаги помечаются в журнале - по этой пометке оператор их
и находит.
3. Шаг и отметка о нём - в одной транзакции. Отсюда *sql.Tx
в сигнатуре Execute и Compensate. Запиши отметку отдельно, и сбой между
ней и эффектом оставит журнал и реальность в разных состояниях: откат либо
пропустит шаг, либо попытается откатить не сделанное. Это тот же класс
ошибки, что в inbox.
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;
use DateTimeImmutable;
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;
use Throwable;
// В 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. Несколько правил:
- Компенсация должна быть идемпотентной. Она может быть вызвана несколько раз: оркестратор перезапустился, retry сработал дважды - и
refundне должен вернуть деньги два раза. Механика та же, что у подписчиков: дедупликация по идентификатору операции, UPSERT вместо инкремента - см. идемпотентность. - Компенсация не должна падать. Если упала - логируем и продолжаем компенсировать остальные шаги. Иначе получим «застрявшую» сагу.
- Компенсация - это семантический откат. Не обязательно удалять данные. Можно пометить заказ как
cancelled, вместо удаления строки. - Не все шаги требуют компенсации. Если шаг только читает данные (валидация) - компенсировать нечего.
Как это проверено
examples/reliability/saga (go test ./...) проверяет утверждения этого
урока прогоном:
- сбой на каждом шаге - по очереди падает первый, второй, третий. После любого сбоя журнал эффектов должен свернуться в ноль: сага либо прошла целиком, либо откатилась целиком;
- сбой на каждой операции базы, включая
COMMITи запись в журнал шагов. Инвариант: число непогашенных эффектов равно числу шагов в состоянииdone. Расхождение означает, что откат пропустит шаг или попробует откатить не сделанное; - идемпотентность компенсации - три вызова
Compensateподряд дают один откат на шаг, а не три; - откат после «перезапуска» - оркестратор создаётся заново, без списка выполненных шагов, и находит их в базе;
- сбой компенсации виден вызывающему - ошибка содержит имя застрявшего шага, а сам шаг помечен в журнале.
Рядом лежит BrokenOrchestrator - версия из первой редакции урока. Тесты
на ней фиксируют оба дефекта: сбой откатa не виден снаружи, а после
перезапуска откатывать нечего. Без сломанной версии тесты выше ничего
не доказывали бы - проверка имеет смысл, только если известно, что она
способна упасть.
Типичные ошибки
- 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после каждого шага