Idempotency + Dead Letter Queue: защита от дублей и ям

Idempotency + Dead Letter Queue: защита от дублей и ям

Повторная доставка сообщений - не баг, а нормальное явление. Consumer упал после обработки, но до ack (как именно работают ack, nack и prefetch)? Сообщение доставится снова. Значит, обработчики обязаны быть идемпотентными - см. идемпотентность в event-driven.

Exactly-once vs At-least-once

Гарантия доставки      Реальность
──────────────────     ─────────────────────────────────────
At-most-once           Сообщение может потеряться (auto-ack)
At-least-once          Сообщение может прийти дважды (manual ack) ← обычно выбираем
Exactly-once           RabbitMQ такого режима не даёт

Формулировка «exactly-once невозможен» ходит по статьям, но она неточная - и из неё делают неверный вывод, будто транзакции Kafka маркетинг. Точнее так:

End-to-end exactly-once для произвольных внешних эффектов нельзя получить одной транспортной гарантией. Kafka даёт exactly-once в границах своей транзакции - чтение из топика, запись в топик и коммит offset атомарны (как это выглядит в коде). RabbitMQ подобного механизма не предоставляет вовсе. Но как только эффект уходит за пределы брокера - запись в PostgreSQL, письмо, чужой HTTP-API, - гарантия заканчивается, и «ровно один раз» обеспечивает уже идемпотентность обработчика.

Для RabbitMQ это значит: берём at-least-once + idempotency и получаем «effectively once»:

At-least-once delivery + Idempotent handler = Effectively-once processing

Idempotency: inbox-таблица

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

Как делать нельзя

НЕПРАВИЛЬНО - разбор ниже

1. INSERT INTO inbox (event_id, handler) ... ON CONFLICT DO NOTHING
2. если вставки не было  → «уже обработано», выходим
3. вызываем handler       ← а если здесь ошибка?

Отметка «обработано» ставится до самой обработки. Между этими двумя строками есть окно: если handler упадёт - сеть, паника, OOM killer, перезапуск пода, - в базе останется запись, что событие обработано, а бизнес-операции не будет.

RabbitMQ честно доставит сообщение снова. Consumer посмотрит в inbox, увидит 0 affected rows и скажет «дубль, пропускаю». Событие потеряно навсегда, и никакой алерт об этом не сработает: с точки зрения брокера всё в порядке, сообщение подтверждено.

Ирония в том, что такой код выглядит как защита от потерь, а сам их создаёт.

Как правильно

Отметка и бизнес-операция должны быть в одной транзакции базы - тогда при ошибке откатятся вместе.

CREATE TABLE inbox (
    event_id     TEXT NOT NULL,
    handler      TEXT NOT NULL, -- какой обработчик выполнил
    processed_at TIMESTAMPTZ NOT NULL DEFAULT now(),

    -- Ключ составной, и это принципиально: одно событие обрабатывают
    -- несколько подписчиков. При PRIMARY KEY (event_id) второй обработчик
    -- увидел бы конфликт и решил, что это дубль первого - хотя свою работу
    -- ещё не делал.
    PRIMARY KEY (event_id, handler)
);
// Handler получает транзакцию, а не открывает своё соединение.
//
// Это не украшение сигнатуры, а условие корректности: обработчик со своим
// соединением коммитится отдельно от отметки в inbox, и «одна транзакция»
// становится неправдой при живом `BeginTx` в коде рядом.
type Handler func(ctx context.Context, tx *sql.Tx, event []byte) error

// IdempotentHandler отмечает событие обработанным и выполняет работу
// в одной транзакции.
type IdempotentHandler struct {
	db     *sql.DB
	name   string // имя обработчика: часть ключа дедупликации
	handle Handler
	log    *slog.Logger
}

func New(db *sql.DB, name string, handle Handler, log *slog.Logger) *IdempotentHandler {
	return &IdempotentHandler{db: db, name: name, handle: handle, log: log}
}

// Handle обрабатывает событие ровно один раз на обработчика.
//
// Порядок шагов важен: сначала транзакция, потом отметка, потом работа.
// Отметка служит и проверкой - отдельный SELECT перед INSERT прошли бы
// два воркера одновременно, оба увидев «не обработано».
func (h *IdempotentHandler) Handle(ctx context.Context, eventID string, body []byte) error {
	tx, err := h.db.BeginTx(ctx, nil)
	if err != nil {
		return fmt.Errorf("begin tx: %w", err)
	}
	// После успешного Commit это no-op, а при любом раннем return
	// транзакция откатится - вместе с отметкой в inbox.
	defer tx.Rollback()

	res, err := tx.ExecContext(ctx,
		`INSERT INTO inbox (event_id, handler)
		 VALUES ($1, $2)
		 ON CONFLICT (event_id, handler) DO NOTHING`,
		eventID, h.name,
	)
	if err != nil {
		return fmt.Errorf("inbox insert: %w", err)
	}

	rows, err := res.RowsAffected()
	if err != nil {
		return fmt.Errorf("rows affected: %w", err)
	}
	if rows == 0 {
		// Этот обработчик событие уже проводил - подтверждаем и уходим.
		h.log.Info("duplicate event, skipping",
			slog.String("event_id", eventID),
			slog.String("handler", h.name),
		)
		return nil
	}

	// Бизнес-операция в той же транзакции. Ошибка здесь откатит и отметку.
	if err := h.handle(ctx, tx, body); err != nil {
		return fmt.Errorf("handler: %w", err)
	}

	return tx.Commit()
}

<?php
declare(strict_types=1);

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

final class IdempotentHandler
{
    /** @param callable(Connection $tx, string $body): void $handler */
    public function __construct(
        private readonly Connection $db,
        private readonly LoggerInterface $logger,
        private readonly string $name,
        private readonly Closure $handler,
    ) {}

    public function handle(string $eventId, string $body): void
    {
        // transactional откатит всё при исключении - и отметку, и работу
        // обработчика. Ключ в том, что обработчик пишет через тот же $tx.
        $this->db->transactional(function (Connection $tx) use ($eventId, $body): void {
            $inserted = $tx->executeStatement(
                'INSERT INTO inbox (event_id, handler) VALUES (:id, :name)
                 ON CONFLICT (event_id, handler) DO NOTHING',
                ['id' => $eventId, 'name' => $this->name],
            );

            if ($inserted === 0) {
                $this->logger->info('duplicate event {id} for {handler}, skipping', [
                    'id' => $eventId,
                    'handler' => $this->name,
                ]);
                return;
            }

            ($this->handler)($tx, $body);
        });
    }
}
Обработчик принимает транзакцию, а не пул соединений, и это не стилистика. Если внутри он возьмёт соединение сам - через глобальный пул, другой репозиторий или отдельный `db.Exec` - его запись окажется **вне** транзакции с отметкой, и мы вернёмся к тому же окну отказа, только менее заметному.

Правило простое: отметка и всё, что она защищает, пишутся одним и тем же tx. Если бизнес-операция физически не может быть в этой транзакции - например, это отправка письма или вызов чужого API - inbox не поможет, нужен outbox и повторная попытка на стороне получателя.

Из того же правила: inbox лежит в базе, куда пишет обработчик. Отдельная «база под идемпотентность» ломает всю схему - две базы нельзя изменить атомарно без распределённых транзакций, которые здесь никто не разворачивает.

Idempotency через бизнес-логику

Иногда inbox не нужен - операция сама по себе идемпотентна:

// Идемпотентно: UPSERT не создаст дубль
func MarkLessonCompleted(ctx context.Context, db *sql.DB, userID, lessonID int64) error {
    _, err := db.ExecContext(ctx,
        `INSERT INTO progress (user_id, lesson_id, completed, completed_at)
         VALUES ($1, $2, true, now())
         ON CONFLICT (user_id, lesson_id)
         DO UPDATE SET completed = true, completed_at = COALESCE(progress.completed_at, now())`,
        userID, lessonID,
    )
    return err
}

// НЕ идемпотентно: каждый вызов добавляет запись
func AddPoints(ctx context.Context, db *sql.DB, userID int64, points int) error {
    _, err := db.ExecContext(ctx,
        `UPDATE users SET points = points + $1 WHERE id = $2`,
        points, userID,
    )
    return err // повторный вызов удвоит очки!
}
<?php
declare(strict_types=1);

use Doctrine\DBAL\Connection;

final class ProgressService
{
    public function __construct(
        private readonly Connection $db,
    ) {}

    // Идемпотентно: UPSERT не создаст дубль
    public function markLessonCompleted(int $userId, int $lessonId): void
    {
        $this->db->executeStatement(
            'INSERT INTO progress (user_id, lesson_id, completed, completed_at)
             VALUES (:uid, :lid, TRUE, NOW())
             ON CONFLICT (user_id, lesson_id)
             DO UPDATE SET completed = TRUE,
                           completed_at = COALESCE(progress.completed_at, NOW())',
            ['uid' => $userId, 'lid' => $lessonId],
        );
    }

    // НЕ идемпотентно: каждый вызов прибавит очки
    public function addPoints(int $userId, int $points): void
    {
        $this->db->executeStatement(
            'UPDATE users SET points = points + :p WHERE id = :id',
            ['p' => $points, 'id' => $userId],
        );
        // повторный вызов удвоит очки!
    }
}

Если операция не идемпотентна по природе - используй inbox.

Dead Letter Queue (DLQ)

DLQ - карантин для сообщений, которые не удалось обработать. Это не мусорка - это место для анализа и повторной обработки.

Настройка DLQ в RabbitMQ

// 1. Создаём DLQ exchange и queue
ch.ExchangeDeclare("dlx", "topic", true, false, false, false, nil)

ch.QueueDeclare("dlq.progress", true, false, false, false, nil)
ch.QueueBind("dlq.progress", "#", "dlx", false, nil)

// 2. Основная очередь с DLQ настройкой
ch.QueueDeclare("progress-worker", true, false, false, false, amqp.Table{
    "x-dead-letter-exchange":    "dlx",           // куда отправлять отклонённые
    "x-dead-letter-routing-key": "progress.dead",  // routing key в DLX
    "x-message-ttl":             int32(60000),      // TTL: 60 секунд (опционально)
})
<?php
declare(strict_types=1);

use PhpAmqpLib\Channel\AMQPChannel;
use PhpAmqpLib\Wire\AMQPTable;

final class DlqTopology
{
    public function __construct(
        private readonly AMQPChannel $ch,
    ) {}

    public function declare(): void
    {
        // 1. Создаём DLQ exchange и queue
        $this->ch->exchange_declare('dlx', 'topic', false, true, false);

        $this->ch->queue_declare('dlq.progress', false, true, false, false);
        $this->ch->queue_bind('dlq.progress', 'dlx', '#');

        // 2. Основная очередь с DLQ настройкой
        $this->ch->queue_declare(
            queue: 'progress-worker',
            durable: true,
            auto_delete: false,
            arguments: new AMQPTable([
                'x-dead-letter-exchange'    => 'dlx',            // куда отправлять отклонённые
                'x-dead-letter-routing-key' => 'progress.dead',  // routing key в DLX
                'x-message-ttl'             => 60000,            // TTL: 60 секунд (опционально)
            ]),
        );
    }
}

Теперь msg.Nack(false, false) (без requeue) отправит сообщение в DLQ вместо удаления.

                    ack
Consumer ◄──── Queue ◄──── Exchange
    │              │
    │ error        │ nack(requeue=false)
    │              ▼
    └────────► DLQ Queue ◄──── DLX Exchange

Что должно быть в DLQ-сообщении

RabbitMQ добавляет заголовок x-death при dead-lettering - то есть когда сообщение отвергнуто без requeue, протухло по TTL или вытеснено по длине очереди:

x-death: [
  {
    count: 3                    # сколько раз через эту очередь с этой причиной
    reason: rejected            # rejected | expired | maxlen | delivery_limit
    queue: progress-worker      # из какой очереди выпало
    time: 2026-05-05 12:00:00
    exchange: events
    routing-keys: [lesson.completed]
  },
  { count: 3, reason: expired, queue: retry.progress, ... }
]
Записи группируются по паре **очередь + причина**, поэтому в retry-цикле их будет несколько: одна про `rejected` из рабочей очереди, другая про `expired` из retry-очереди. Читать `x-death[0].count` и считать это «числом попыток» - ошибка, которая проявится при первом же усложнении топологии.

Нужную запись ищут по имени очереди и причине - как в retryCount из раздела про poison messages ниже.

И главное: x-death появляется только при dead-lettering. Если отвергать сообщение с requeue=true, оно вернётся в очередь напрямую, минуя DLX, - и счётчик не изменится никогда. Про эту ловушку - в разделе про poison messages ниже.

Retry с DLQ

Паттерн «retry через DLQ» с задержкой:

Retry-цикл через DLQ с TTL: сообщение зреет в DLQ и возвращается в основную очередь

1. Сообщение не обработано → nack(requeue=false) → DLQ
2. DLQ имеет TTL → через N секунд сообщение «протухает»
3. DLQ имеет свой DLX → «протухшее» сообщение возвращается в основную очередь
4. Consumer пробует снова
5. После N попыток → финальная DLQ (без retry)
// Retry queue с TTL - сообщение вернётся через 30 секунд
ch.QueueDeclare("retry.progress", true, false, false, false, amqp.Table{
    "x-dead-letter-exchange":    "events",           // вернуть в основной exchange
    "x-dead-letter-routing-key": "lesson.completed",
    "x-message-ttl":             int32(30000),        // 30 секунд задержка
})

// Финальная DLQ - сюда попадают после исчерпания ретраев
ch.QueueDeclare("dlq.progress.final", true, false, false, false, nil)
<?php
declare(strict_types=1);

use PhpAmqpLib\Channel\AMQPChannel;
use PhpAmqpLib\Wire\AMQPTable;

final class RetryTopology
{
    public function __construct(
        private readonly AMQPChannel $ch,
    ) {}

    public function declare(): void
    {
        // Retry queue с TTL - сообщение вернётся через 30 секунд
        $this->ch->queue_declare(
            queue: 'retry.progress',
            durable: true,
            arguments: new AMQPTable([
                'x-dead-letter-exchange'    => 'events',           // вернуть в основной exchange
                'x-dead-letter-routing-key' => 'lesson.completed',
                'x-message-ttl'             => 30000,              // 30 секунд задержка
            ]),
        );

        // Финальная DLQ - сюда попадают после исчерпания ретраев
        $this->ch->queue_declare('dlq.progress.final', durable: true);
    }
}

Poison Messages

Poison message - сообщение, которое всегда вызывает ошибку: невалидный JSON, неизвестный тип события, баг в обработчике. Его нельзя ни обработать, ни бесконечно возвращать в очередь.

Ловушка: requeue=true не считается попыткой

Естественное желание - «не получилось, вернём в очередь и попробуем ещё»:

НЕПРАВИЛЬНО - разбор ниже

handler вернул ошибку:
  если retryCount(msg) >= 3   ← счётчик читается из x-death
      nack(requeue=false)      → сдаёмся, в DLQ
  иначе
      nack(requeue=true)       → «попробуем ещё»

Этот код не работает, и ломается он тихо. requeue=true возвращает сообщение в ту же очередь напрямую, не через dead-letter exchange. Значит x-death не появляется и retryCount навсегда остаётся нулём:

handler → ошибка → nack(requeue=true) → та же очередь
       → handler → ошибка → retryCount = 0 → nack(requeue=true)
       → ... и так пока кто-нибудь не заметит

Вместо трёх попыток получается бесконечный цикл, который жжёт CPU и забивает логи одним и тем же сообщением. Причём consumer выглядит живым и работающим.

Как правильно: только dead-lettering

Отвергаем всегда без requeue, а задержку и повторы обеспечивает retry-очередь с TTL - та самая топология из раздела выше. Разница лишь в том, куда сообщение уйдёт: в retry-очередь или в финальную DLQ.

const maxRetries = 3

func (c *Consumer) handleSafely(ctx context.Context, msg amqp.Delivery, handler func(context.Context, []byte) error) {
    // Паника обработчика - тот же случай, что ошибка: сообщение
    // отправляем на разбор, а не роняем consumer.
    defer func() {
        if r := recover(); r != nil {
            log.Printf("panic on message %s: %v", msg.MessageId, r)
            c.reject(msg, fmt.Errorf("panic: %v", r))
        }
    }()

    if len(msg.Body) == 0 {
        // Пустое тело не станет валидным от повторов - сразу в финальную DLQ,
        // без круга по retry-очереди.
        log.Printf("empty body, straight to final DLQ")
        c.publishToFinalDLQ(ctx, msg, "empty body")
        msg.Ack(false)
        return
    }

    if err := handler(ctx, msg.Body); err != nil {
        c.reject(msg, err)
        return
    }

    msg.Ack(false)
}

// retryCount ищет запись x-death по очереди и причине. Брать первую нельзя:
// в retry-цикле их несколько - rejected из рабочей очереди и expired
// из retry-очереди.
func retryCount(msg amqp.Delivery, queue, reason string) int32 {
    deaths, ok := msg.Headers["x-death"].([]interface{})
    if !ok {
        return 0
    }
    for _, d := range deaths {
        entry, ok := d.(amqp.Table)
        if !ok {
            continue
        }
        if entry["queue"] == queue && entry["reason"] == reason {
            if count, ok := entry["count"].(int64); ok {
                return int32(count)
            }
        }
    }
    return 0
}

// reject решает, куда отправить сообщение: на очередной круг retry
// или в финальную DLQ. Счётчик берётся из x-death именно той очереди
// и той причины, которые нас интересуют.
func (c *Consumer) reject(msg amqp.Delivery, cause error) {
    attempts := retryCount(msg, c.queue, "rejected")

    if attempts >= maxRetries {
        log.Printf("message %s failed %d times, final DLQ: %v", msg.MessageId, attempts, cause)
        // Публикуем в финальную DLQ сами: так к сообщению можно приложить
        // причину последней ошибки - при dead-lettering её негде взять.
        c.publishToFinalDLQ(context.Background(), msg, cause.Error())
        msg.Ack(false)
        return
    }

    log.Printf("message %s attempt %d failed, retry later: %v", msg.MessageId, attempts+1, cause)
    // Без requeue: сообщение уходит через DLX в retry-очередь,
    // отлёживается там TTL и возвращается обратно. Именно этот проход
    // увеличивает x-death.
    msg.Nack(false, false)
}
<?php
declare(strict_types=1);

use PhpAmqpLib\Message\AMQPMessage;
use Psr\Log\LoggerInterface;

final class PoisonSafeConsumer
{
    private const MAX_RETRIES = 3;

    public function __construct(
        private readonly LoggerInterface $logger,
        private readonly FinalDlqPublisher $finalDlq,
        private readonly string $queue,
    ) {}

    /** @param callable(string): void $handler */
    public function handleSafely(AMQPMessage $msg, callable $handler): void
    {
        if ($msg->getBody() === '') {
            // Повторы не сделают пустое тело валидным.
            $this->logger->warning('empty body, straight to final DLQ');
            $this->finalDlq->publish($msg, 'empty body');
            $msg->ack();
            return;
        }

        try {
            $handler($msg->getBody());
            $msg->ack();
        } catch (Throwable $e) {
            $this->reject($msg, $e);
        }
    }

    private function reject(AMQPMessage $msg, Throwable $cause): void
    {
        $attempts = $this->retryCount($msg, $this->queue, 'rejected');

        if ($attempts >= self::MAX_RETRIES) {
            $this->logger->error('message {id} failed {n} times, final DLQ: {err}', [
                'id'  => $msg->get_properties()['message_id'] ?? '',
                'n'   => $attempts,
                'err' => $cause->getMessage(),
            ]);
            $this->finalDlq->publish($msg, $cause->getMessage());
            $msg->ack();
            return;
        }

        // requeue: false - обязательно. Только так сообщение пройдёт через DLX
        // в retry-очередь, и только так вырастет x-death.
        $msg->nack(requeue: false);
    }

    /** Ищем запись x-death по очереди и причине, а не берём первую. */
    private function retryCount(AMQPMessage $msg, string $queue, string $reason): int
    {
        $headers = $msg->get_properties()['application_headers'] ?? null;
        $deaths = $headers?->getNativeData()['x-death'] ?? [];

        foreach ($deaths as $entry) {
            if (($entry['queue'] ?? null) === $queue && ($entry['reason'] ?? null) === $reason) {
                return (int) ($entry['count'] ?? 0);
            }
        }
        return 0;
    }
}
Запрет не абсолютный. `requeue=true` подходит там, где ошибка заведомо не связана с самим сообщением и пройдёт сама: например, consumer потерял соединение с базой и знает, что переподключение уже идёт. Тогда возврат в очередь - это «сейчас не я, пусть возьмёт другой».

Но такой возврат не считается попыткой и не имеет задержки, поэтому ограничивать его надо своим счётчиком в памяти consumer-а, а не x-death. Для повторов бизнес-ошибок остаётся схема с retry-очередью.

Мониторинг DLQ

Алерт на сам факт «в DLQ что-то есть» будит дежурного на каждом случайном сбое. Смотреть надо на rate попаданий и на возраст старейшего сообщения; полный набор сигналов и порогов для event-driven системы - в уроке наблюдаемость.

// Проверяем размер DLQ через Management API
func CheckDLQSize(mgmtURL, queue string) (int, error) {
    resp, err := http.Get(fmt.Sprintf("%s/api/queues/%%2F/%s", mgmtURL, queue))
    if err != nil {
        return 0, err
    }
    defer resp.Body.Close()

    var q struct {
        Messages int `json:"messages"`
    }
    json.NewDecoder(resp.Body).Decode(&q)
    return q.Messages, nil
}
<?php
declare(strict_types=1);

use Symfony\Contracts\HttpClient\HttpClientInterface;

final class DlqMonitor
{
    public function __construct(
        private readonly HttpClientInterface $http,
        private readonly string $mgmtUrl,    // например: http://rabbitmq:15672
        private readonly string $user,
        private readonly string $pass,
    ) {}

    // Проверяем размер DLQ через Management API
    public function size(string $queue): int
    {
        $vhost = rawurlencode('/');
        $response = $this->http->request('GET', "{$this->mgmtUrl}/api/queues/{$vhost}/{$queue}", [
            'auth_basic' => [$this->user, $this->pass],
        ]);

        $data = $response->toArray();
        return (int) ($data['messages'] ?? 0);
    }
}

Как это проверено

Оба утверждения этого урока - «отметка и бизнес-операция применяются атомарно» и «после N неудачных попыток сообщение уходит в очередь разбора» - проверяются прогоном, а не чтением: examples/reliability/.

Стенд собран так, чтобы проверять именно их, а не работу моков:

  • inbox гоняет код из этого урока на настоящих BeginTx, Commit и Rollback - под ними драйвер database/sql в памяти, умеющий вернуть ошибку на любой по счёту операции. Тест перебирает все точки отказа (вставка отметки, бизнес-операция, Commit) и после каждой проверяет, что состояние согласовано: либо отметка и эффект вместе, либо ни того, ни другого. Состояние «отметка есть, эффекта нет» запрещено - именно оно означает потерю события навсегда;
  • retry моделирует ту часть RabbitMQ, которая здесь и решает: очереди, dead-letter exchange, TTL и x-death как массив записей по паре (очередь, причина). Тест доводит сообщение до очереди разбора и сверяет число доставок с числом разрешённых попыток.

Рядом с каждой рабочей версией лежит сломанная - та, что была в этом уроке до правок. Без неё тесты ничего не доказывали бы: проверка имеет смысл только если известно, что она способна упасть. Сломанный inbox теряет событие при сбое в окне, сломанный consumer крутит сообщение бесконечно, и оба факта зафиксированы тестами.

Утверждение «в одной транзакции» держалось в этом уроке на том, что рядом стоит `BeginTx`. Обработчик при этом писал в своё соединение - и код, и текст выглядели правильными по отдельности. Такую ошибку не видно чтением: её видно только прогоном, в котором сбой случается **между** двумя побочными эффектами.

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

  • Создай таблицу inbox и оберни свой handler в IdempotentHandler
  • Отправь одно и то же сообщение дважды - убедись, что обработка произошла один раз
  • Настрой DLQ для своей очереди через x-dead-letter-exchange
  • Сделай handler, который всегда возвращает ошибку, и убедись, что после 3 ретраев сообщение попало в DLQ

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