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: паттерн обработки
Основной паттерн: проверить - обработать - отметить.
// 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());
}
}
Типичные ошибки
- Проверка «не обработано ли» отдельно от обработки - две 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 и убедись, что бизнес-логика выполнилась один раз