Idempotency: как не выполнить одно и то же дважды

Idempotency: как не выполнить одно и то же дважды

В мире событий доставка работает по принципу at-least-once: брокер гарантирует, что сообщение дойдёт хотя бы один раз, но может доставить его дважды, трижды или десять раз. Сеть моргнула, consumer перезапустился, offset не закоммитился - и событие приходит повторно. Если обработчик не готов к этому, пользователь получит два списания, две ачивки или два письма.

Idempotency - это свойство операции: повторный вызов с теми же входными данными не меняет результат. f(x) = f(f(x)). Для event-driven архитектуры это не оптимизация и не nice-to-have. Это обязательное требование к каждому подписчику.

Почему at-least-once - это норма

Гарантия exactly-once на уровне транспорта невозможна в распределённых системах (доказано в теории - Two Generals' Problem). Брокеры вроде Kafka, NATS и RabbitMQ дают at-least-once: если consumer не подтвердил обработку, сообщение будет доставлено снова. Exactly-once достигается на уровне приложения - через идемпотентную обработку.

Ситуации, когда событие приходит повторно:

  • Consumer упал после обработки, но до коммита offset
  • Сетевой таймаут между consumer и брокером
  • Rebalancing партиций в Kafka
  • Ручной replay событий при инциденте

Idempotency Key: event ID как натуральный ключ

Каждое событие несёт уникальный event_id (UUID или ULID). Этот идентификатор и есть idempotency key. Перед обработкой проверяем: видели ли мы это событие раньше? Если да - пропускаем. Если нет - обрабатываем и запоминаем.

Inbox-таблица: хранилище обработанных событий

Для надёжной дедупликации нужна таблица в базе данных:

CREATE TABLE processed_events (
    event_id    VARCHAR(64) PRIMARY KEY,
    event_type  VARCHAR(128) NOT NULL,
    handler     VARCHAR(128) NOT NULL,
    processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
);

-- Индекс для очистки старых записей
CREATE INDEX idx_processed_events_processed_at
    ON processed_events (processed_at);

Поле handler важно: одно и то же событие могут обрабатывать несколько подписчиков. Achievement-сервис и Analytics-сервис - это разные обработчики, и каждый должен выполнить свою работу ровно один раз. Поэтому составной ключ (event_id, handler) точнее, но для простоты - если у вас один обработчик на таблицу - хватит event_id.

Check-Process-Mark: паттерн обработки

Первый приём: check inbox not found → process → mark; повторный: check found → skip

Основной паттерн: проверить - обработать - отметить.

// ProcessedEventsRepo отвечает за дедупликацию событий.
type ProcessedEventsRepo struct {
    db *sql.DB
}

// AlreadyProcessed проверяет, было ли событие обработано.
func (r *ProcessedEventsRepo) AlreadyProcessed(ctx context.Context, eventID, handler string) (bool, error) {
    var exists bool
    err := r.db.QueryRowContext(ctx,
        `SELECT EXISTS(SELECT 1 FROM processed_events WHERE event_id = $1 AND handler = $2)`,
        eventID, handler,
    ).Scan(&exists)
    return exists, err
}

// MarkProcessed отмечает событие как обработанное.
func (r *ProcessedEventsRepo) MarkProcessed(ctx context.Context, eventID, handler, eventType string) error {
    _, err := r.db.ExecContext(ctx,
        `INSERT INTO processed_events (event_id, event_type, handler)
         VALUES ($1, $2, $3)
         ON CONFLICT (event_id) DO NOTHING`,
        eventID, eventType, handler,
    )
    return err
}
<?php
// src/Infrastructure/Inbox/ProcessedEventsRepository.php
declare(strict_types=1);

namespace App\Infrastructure\Inbox;

use App\Application\Port\ProcessedEventsRepository as ProcessedEventsRepositoryPort;
use Doctrine\DBAL\Connection;

// Doctrine DBAL - тонкая прослойка над PDO. Симметрична database/sql из Go.
final class ProcessedEventsRepository implements ProcessedEventsRepositoryPort
{
    public function __construct(
        private readonly Connection $connection,
    ) {}

    public function alreadyProcessed(string $eventId, string $handler): bool
    {
        $row = $this->connection->fetchOne(
            'SELECT EXISTS(SELECT 1 FROM processed_events WHERE event_id = :eid AND handler = :h)',
            ['eid' => $eventId, 'h' => $handler],
        );
        return (bool) $row;
    }

    public function markProcessed(string $eventId, string $handler, string $eventType): void
    {
        // ON CONFLICT DO NOTHING - Postgres-specific.
        // На MySQL: INSERT IGNORE INTO ...
        $this->connection->executeStatement(
            'INSERT INTO processed_events (event_id, event_type, handler)
             VALUES (:eid, :etype, :h)
             ON CONFLICT (event_id) DO NOTHING',
            ['eid' => $eventId, 'etype' => $eventType, 'h' => $handler],
        );
    }
}

Обработчик события использует этот репозиторий:

// HandleLessonCompleted обрабатывает событие завершения урока.
// Идемпотентен: повторный вызов с тем же event_id не даёт эффекта.
func (h *AchievementHandler) HandleLessonCompleted(ctx context.Context, evt LessonCompleted) error {
    const handlerName = "achievement_checker"

    // 1. Check: уже обработано?
    seen, err := h.inbox.AlreadyProcessed(ctx, evt.EventID, handlerName)
    if err != nil {
        return fmt.Errorf("inbox check: %w", err)
    }
    if seen {
        slog.InfoContext(ctx, "event already processed, skipping",
            slog.String("event_id", evt.EventID),
            slog.String("handler", handlerName),
        )
        return nil
    }

    // 2. Process: выполняем бизнес-логику
    if err := h.checkAndGrantAchievements(ctx, evt.UserID, evt.LessonID); err != nil {
        return fmt.Errorf("grant achievements: %w", err)
    }

    // 3. Mark: запоминаем
    if err := h.inbox.MarkProcessed(ctx, evt.EventID, handlerName, evt.Type); err != nil {
        return fmt.Errorf("inbox mark: %w", err)
    }

    return nil
}
<?php
// src/Application/Subscriber/AchievementSubscriber.php
declare(strict_types=1);

namespace App\Application\Subscriber;

use App\Application\Port\ProcessedEventsRepository;
use App\Domain\Event\LessonCompleted;
use Psr\Log\LoggerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;

#[AsMessageHandler]
final class AchievementSubscriber
{
    private const HANDLER_NAME = 'achievement_checker';

    public function __construct(
        private readonly ProcessedEventsRepository $inbox,
        private readonly AchievementGranter $granter,
        private readonly LoggerInterface $logger,
    ) {}

    public function __invoke(LessonCompleted $event): void
    {
        // 1. Check: уже обработано?
        if ($this->inbox->alreadyProcessed($event->id, self::HANDLER_NAME)) {
            $this->logger->info('event already processed, skipping', [
                'event_id' => $event->id,
                'handler' => self::HANDLER_NAME,
            ]);
            return;
        }

        // 2. Process: выполняем бизнес-логику
        $this->granter->grantAchievements($event->userId, $event->lessonId);

        // 3. Mark: запоминаем
        $this->inbox->markProcessed($event->id, self::HANDLER_NAME, $event->eventType());
    }
}

Атомарная отметка: INSERT ON CONFLICT DO NOTHING

Обратите внимание на ON CONFLICT DO NOTHING в MarkProcessed. Это критически важно. Если два воркера одновременно получили одно и то же событие, оба пройдут проверку AlreadyProcessed (оба получат false) и оба начнут обработку. ON CONFLICT DO NOTHING гарантирует, что вставка пройдёт только у первого. Но бизнес-логика всё равно выполнится дважды.

Для полной защиты от гонки нужна транзакция с блокировкой:

// HandleAtomically обрабатывает событие в одной транзакции с отметкой.
func (h *AchievementHandler) HandleAtomically(ctx context.Context, evt LessonCompleted) error {
    tx, err := h.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    // Пытаемся вставить - если уже есть, rows == 0
    res, err := tx.ExecContext(ctx,
        `INSERT INTO processed_events (event_id, event_type, handler)
         VALUES ($1, $2, $3)
         ON CONFLICT (event_id) DO NOTHING`,
        evt.EventID, evt.Type, "achievement_checker",
    )
    if err != nil {
        return fmt.Errorf("inbox insert: %w", err)
    }

    rows, _ := res.RowsAffected()
    if rows == 0 {
        // Событие уже обработано другим воркером
        return nil
    }

    // Бизнес-логика внутри той же транзакции
    if err := h.grantAchievementsInTx(ctx, tx, evt.UserID, evt.LessonID); err != nil {
        return fmt.Errorf("grant: %w", err)
    }

    return tx.Commit()
}
<?php
// src/Application/Subscriber/AchievementHandlerAtomic.php
declare(strict_types=1);

namespace App\Application\Subscriber;

use App\Domain\Event\LessonCompleted;
use Doctrine\DBAL\Connection;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;

#[AsMessageHandler]
final class AchievementHandlerAtomic
{
    private const HANDLER_NAME = 'achievement_checker';

    public function __construct(
        private readonly Connection $connection,
        private readonly AchievementGranter $granter,
    ) {}

    public function __invoke(LessonCompleted $event): void
    {
        $this->connection->transactional(function (Connection $tx) use ($event): void {
            // Пытаемся вставить - если уже есть, affectedRows == 0
            $affected = $tx->executeStatement(
                'INSERT INTO processed_events (event_id, event_type, handler)
                 VALUES (:eid, :etype, :h)
                 ON CONFLICT (event_id) DO NOTHING',
                ['eid' => $event->id, 'etype' => $event->eventType(), 'h' => self::HANDLER_NAME],
            );

            if ($affected === 0) {
                // Событие уже обработано другим воркером
                return;
            }

            // Бизнес-логика внутри той же транзакции
            $this->granter->grantInTransaction($tx, $event->userId, $event->lessonId);
        });
    }
}

Теперь отметка и бизнес-логика атомарны. Если бизнес-логика упала - откатится и отметка, и событие будет обработано заново.

Естественная идемпотентность

Не все операции требуют inbox-таблицы. Некоторые идемпотентны по своей природе:

UPSERT вместо INSERT. Если подписчик записывает результат в таблицу, используйте INSERT ... ON CONFLICT DO UPDATE:

// RecordLessonCompletion записывает факт завершения урока.
// Идемпотентна: повторный вызов обновляет completed_at, но не создаёт дубль.
func (r *AnalyticsRepo) RecordLessonCompletion(ctx context.Context, userID int64, lessonID int64) error {
    _, err := r.db.ExecContext(ctx,
        `INSERT INTO lesson_completions (user_id, lesson_id, completed_at)
         VALUES ($1, $2, now())
         ON CONFLICT (user_id, lesson_id)
         DO UPDATE SET completed_at = EXCLUDED.completed_at`,
        userID, lessonID,
    )
    return err
}
<?php
// src/Infrastructure/Analytics/AnalyticsRepository.php
declare(strict_types=1);

namespace App\Infrastructure\Analytics;

use Doctrine\DBAL\Connection;

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

    // Идемпотентна: повторный вызов обновляет completed_at, но не создаёт дубль.
    public function recordLessonCompletion(int $userId, int $lessonId): void
    {
        $this->connection->executeStatement(
            'INSERT INTO lesson_completions (user_id, lesson_id, completed_at)
             VALUES (:uid, :lid, now())
             ON CONFLICT (user_id, lesson_id)
             DO UPDATE SET completed_at = EXCLUDED.completed_at',
            ['uid' => $userId, 'lid' => $lessonId],
        );
    }
}

State machine transitions. Если статус может двигаться только вперёд (pending -> paid -> shipped), повторная попытка перевести из pending в paid ничего не сломает, а попытка перевести из shipped в paid просто не пройдёт:

// AdvanceOrderStatus переводит заказ на следующий статус.
// Идемпотентна: WHERE clause гарантирует, что переход произойдёт только один раз.
func (r *OrderRepo) AdvanceOrderStatus(ctx context.Context, orderID int64, from, to string) error {
    res, err := r.db.ExecContext(ctx,
        `UPDATE orders SET status = $1, updated_at = now()
         WHERE id = $2 AND status = $3`,
        to, orderID, from,
    )
    if err != nil {
        return err
    }
    rows, _ := res.RowsAffected()
    if rows == 0 {
        slog.InfoContext(ctx, "order already advanced past this status",
            slog.Int64("order_id", orderID),
            slog.String("expected_from", from),
        )
    }
    return nil
}
<?php
// src/Infrastructure/Order/OrderRepository.php
declare(strict_types=1);

namespace App\Infrastructure\Order;

use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;

final class OrderRepository
{
    public function __construct(
        private readonly Connection $connection,
        private readonly LoggerInterface $logger,
    ) {}

    // Идемпотентна: WHERE clause гарантирует, что переход произойдёт только один раз.
    public function advanceStatus(int $orderId, string $from, string $to): void
    {
        $affected = $this->connection->executeStatement(
            'UPDATE orders SET status = :to, updated_at = now()
             WHERE id = :oid AND status = :from',
            ['to' => $to, 'oid' => $orderId, 'from' => $from],
        );

        if ($affected === 0) {
            $this->logger->info('order already advanced past this status', [
                'order_id' => $orderId,
                'expected_from' => $from,
            ]);
        }
    }
}

TTL для processed_events

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

// CleanupOldEvents удаляет записи старше ttl.
// Запускайте по cron раз в сутки.
func (r *ProcessedEventsRepo) CleanupOldEvents(ctx context.Context, ttl time.Duration) (int64, error) {
    res, err := r.db.ExecContext(ctx,
        `DELETE FROM processed_events WHERE processed_at < $1`,
        time.Now().Add(-ttl),
    )
    if err != nil {
        return 0, err
    }
    return res.RowsAffected()
}
<?php
// src/Application/Cron/CleanupProcessedEventsCommand.php
declare(strict_types=1);

namespace App\Application\Cron;

use Doctrine\DBAL\Connection;
use Symfony\Component\Console\Attribute\AsCommand;
use Symfony\Component\Console\Command\Command;
use Symfony\Component\Console\Input\InputInterface;
use Symfony\Component\Console\Output\OutputInterface;

// Запускается из cron: `bin/console app:cleanup-processed-events --ttl-days=14`.
// В Symfony альтернатива - Scheduler component (symfony/scheduler).
#[AsCommand(name: 'app:cleanup-processed-events')]
final class CleanupProcessedEventsCommand extends Command
{
    public function __construct(
        private readonly Connection $connection,
    ) {
        parent::__construct();
    }

    protected function execute(InputInterface $input, OutputInterface $output): int
    {
        $ttlDays = 14;
        $threshold = (new \DateTimeImmutable())->modify(sprintf('-%d days', $ttlDays));

        $affected = $this->connection->executeStatement(
            'DELETE FROM processed_events WHERE processed_at < :threshold',
            ['threshold' => $threshold->format('Y-m-d H:i:s')],
        );

        $output->writeln(sprintf('deleted %d rows', $affected));
        return Command::SUCCESS;
    }
}

Разумный TTL - от 7 до 30 дней, в зависимости от того, как далеко назад вы можете replay-ить события.

HTTP API: заголовок Idempotency-Key

Idempotency нужна не только внутри event-driven системы. HTTP-клиенты тоже могут повторять запросы (таймаут, retry). Stripe популяризировал паттерн с заголовком Idempotency-Key:

// IdempotencyMiddleware проверяет заголовок Idempotency-Key.
// Если запрос с таким ключом уже обработан - возвращает кэшированный ответ.
func IdempotencyMiddleware(repo IdempotencyRepo) func(http.Handler) http.Handler {
    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            key := r.Header.Get("Idempotency-Key")
            if key == "" {
                next.ServeHTTP(w, r)
                return
            }

            // Проверяем кэш
            cached, err := repo.GetCachedResponse(r.Context(), key)
            if err == nil && cached != nil {
                w.WriteHeader(cached.StatusCode)
                w.Write(cached.Body)
                return
            }

            // Оборачиваем ResponseWriter для захвата ответа
            rec := &responseRecorder{ResponseWriter: w}
            next.ServeHTTP(rec, r)

            // Кэшируем ответ
            repo.CacheResponse(r.Context(), key, rec.statusCode, rec.body.Bytes())
        })
    }
}
<?php
// src/Infrastructure/Http/IdempotencyListener.php
declare(strict_types=1);

namespace App\Infrastructure\Http;

use App\Application\Port\IdempotencyRepository;
use Symfony\Component\EventDispatcher\Attribute\AsEventListener;
use Symfony\Component\HttpFoundation\Response;
use Symfony\Component\HttpKernel\Event\RequestEvent;
use Symfony\Component\HttpKernel\Event\ResponseEvent;

// В Symfony роль HTTP-middleware из Go играют kernel-listener'ы.
// onRequest проверяет кэш, onResponse сохраняет ответ.
final class IdempotencyListener
{
    public function __construct(
        private readonly IdempotencyRepository $repo,
    ) {}

    #[AsEventListener(event: RequestEvent::class)]
    public function onRequest(RequestEvent $event): void
    {
        $key = $event->getRequest()->headers->get('Idempotency-Key');
        if ($key === null) {
            return;
        }

        $cached = $this->repo->getCachedResponse($key);
        if ($cached !== null) {
            // Возвращаем кэшированный ответ - дальше kernel не идёт
            $event->setResponse(new Response($cached->body, $cached->statusCode));
        }
    }

    #[AsEventListener(event: ResponseEvent::class)]
    public function onResponse(ResponseEvent $event): void
    {
        $key = $event->getRequest()->headers->get('Idempotency-Key');
        if ($key === null) {
            return;
        }

        $response = $event->getResponse();
        $this->repo->cacheResponse($key, $response->getStatusCode(), (string) $response->getContent());
    }
}
В at-least-once системе каждый подписчик ОБЯЗАН быть идемпотентным. Без этого повторная доставка приведёт к дублям данных, двойным списаниям и повреждению состояния. Проектируйте идемпотентность с самого начала - добавлять её потом значительно дороже.

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

  • Проверка «не обработано ли» отдельно от обработки - две reads + два writes без транзакции. Между check и mark два инстанса видят «не обработано», оба делают работу. Атомарность через INSERT ... ON CONFLICT DO NOTHING + RETURNING в одной транзакции с бизнес-логикой.
  • Idempotency-key = таймстамп клиента - два разных POST от одного юзера в одну миллисекунду дают одинаковый ключ; одно из действий теряется. Используй UUID, генерируемый клиентом до запроса.
  • Idempotency-store разделён с бизнес-БД - отметка о processing идёт в Redis, бизнес-данные в PostgreSQL. Redis перезагрузился без AOF → дубли. Idempotency живёт в той же БД, что и бизнес-операция; одна транзакция.
  • Различные subscriber-ы пишут в одну processed_events без поля handler - subscriber A обработал, отметил; subscriber B видит «уже обработано» и пропускает свою работу. Ключ дедупа = (event_id, handler_name), не только event_id.
  • UPSERT вместо «начислить +10 бонусов» - повторное событие → бонусы перезаписываются тем же значением (это корректно). Но если бы было UPDATE balance = balance + 10 - двойное начисление. Понимай: идемпотентность через UPSERT работает для установки значения, не для инкремента.
  • Idempotency только в HTTP-слое - клиент шлёт повторный POST с тем же Idempotency-Key, HTTP-middleware возвращает кэшированный ответ. Но внутреннее событие тоже опубликовано дважды - а subscriber не имеет своей дедупликации. Дедуп нужен на каждом узле at-least-once-цепочки.
  • Сравнение processed_at для «свежести» - «не обработать события старше 1 часа». При replay из брокера старые события массово пропускаются и не обновляют read-модели. Не путай идемпотентность с фильтрацией.

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

  • Создай таблицу processed_events с полями event_id, handler, event_type, processed_at
  • Реализуй атомарный check-process-mark паттерн с INSERT ON CONFLICT DO NOTHING внутри транзакции
  • Перепиши один из своих обработчиков, чтобы он использовал естественную идемпотентность (UPSERT или state machine)
  • Добавь горутину очистки записей старше 14 дней
  • Напиши unit-тест: вызови обработчик дважды с одним event_id и убедись, что бизнес-логика выполнилась один раз

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