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

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

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

Exactly-once vs At-least-once

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

Exactly-once - миф. На практике мы используем at-least-once + idempotency и получаем «effectively once»:

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

Idempotency: inbox-таблица

Самый надёжный способ - сохранять event_id обработанных сообщений:

Inbox-таблица: повторный INSERT по event_id падает по PK, handler пропускает обработку

CREATE TABLE inbox (
    event_id     TEXT PRIMARY KEY,
    processed_at TIMESTAMPTZ NOT NULL DEFAULT now(),
    handler      TEXT NOT NULL - какой обработчик выполнил
);
type IdempotentHandler struct {
    db      *sql.DB
    handler func(ctx context.Context, event []byte) error
    name    string
}

func (h *IdempotentHandler) Handle(ctx context.Context, eventID string, body []byte) error {
    // Пытаемся вставить - если уже есть, ON CONFLICT ничего не делает
    result, err := h.db.ExecContext(ctx,
        `INSERT INTO inbox (event_id, handler) VALUES ($1, $2) ON CONFLICT DO NOTHING`,
        eventID, h.name,
    )
    if err != nil {
        return fmt.Errorf("inbox insert: %w", err)
    }

    rows, _ := result.RowsAffected()
    if rows == 0 {
        // Уже обработано - пропускаем
        log.Printf("duplicate event %s, skipping", eventID)
        return nil
    }

    // Первый раз - обрабатываем
    return h.handler(ctx, body)
}
<?php
declare(strict_types=1);

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

final class IdempotentHandler
{
    /** @param callable(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
    {
        // Пытаемся вставить - если уже есть, ON CONFLICT ничего не делает
        $inserted = $this->db->executeStatement(
            'INSERT INTO inbox (event_id, handler) VALUES (:id, :name) ON CONFLICT DO NOTHING',
            ['id' => $eventId, 'name' => $this->name],
        );

        if ($inserted === 0) {
            // Уже обработано - пропускаем
            $this->logger->info('duplicate event {id}, skipping', ['id' => $eventId]);
            return;
        }

        // Первый раз - обрабатываем
        ($this->handler)($body);
    }
}
Inbox-таблица должна быть в той же базе, что и бизнес-данные. Тогда вставка в 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 с метаданными:

x-death[0]:
  count: 3                    # сколько раз сообщение попадало в DLQ
  reason: rejected            # почему: rejected, expired, maxlen
  queue: progress-worker      # откуда пришло
  time: 2026-05-05 12:00:00   # когда
  exchange: events
  routing-keys: [lesson.completed]

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, неизвестный тип, баг в обработчике).

func (c *Consumer) handleSafely(ctx context.Context, msg amqp.Delivery, handler func(context.Context, []byte) error) {
    // Защита от паники
    defer func() {
        if r := recover(); r != nil {
            log.Printf("panic processing message %s: %v", msg.MessageId, r)
            msg.Nack(false, false) // в DLQ
        }
    }()

    // Проверяем базовую валидность
    if len(msg.Body) == 0 {
        log.Printf("empty message body, sending to DLQ")
        msg.Nack(false, false)
        return
    }

    // Пробуем обработать
    if err := handler(ctx, msg.Body); err != nil {
        retries := getRetryCount(msg)
        if retries >= 3 {
            log.Printf("message %s failed after %d retries: %v", msg.MessageId, retries, err)
            msg.Nack(false, false) // в DLQ
            return
        }
        msg.Nack(false, true) // попробуем ещё
        return
    }

    msg.Ack(false)
}
<?php
declare(strict_types=1);

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

final class PoisonSafeConsumer
{
    public function __construct(
        private readonly LoggerInterface $logger,
        private readonly RetryConsumer $retryHelper, // вычисляет x-death count
    ) {}

    /** @param callable(string): void $handler */
    public function handleSafely(AMQPMessage $msg, callable $handler): void
    {
        try {
            // Проверяем базовую валидность
            if ($msg->getBody() === '') {
                $this->logger->warning('empty message body, sending to DLQ');
                $msg->nack(requeue: false);
                return;
            }

            // Пробуем обработать
            $handler($msg->getBody());
            $msg->ack();
        } catch (Throwable $e) {
            $retries = $this->retryHelper->getRetryCount($msg);
            if ($retries >= 3) {
                $this->logger->error('message {id} failed after {n} retries: {err}', [
                    'id'  => $msg->get('message_id') ?? '',
                    'n'   => $retries,
                    'err' => $e->getMessage(),
                ]);
                $msg->nack(requeue: false); // в DLQ
                return;
            }
            $msg->nack(requeue: true); // попробуем ещё
        }
    }
}
DLQ без мониторинга - бесполезна. Настрой алерт, если в DLQ накопилось > N сообщений. Периодически проверяй DLQ: часть сообщений можно перепроцессить после фикса бага, часть - удалить как невалидные.

Мониторинг DLQ

// Проверяем размер 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);
    }
}

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

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

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