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);
});
}
}
Правило простое: отметка и всё, что она защищает, пишутся одним и тем же
tx. Если бизнес-операция физически не может быть в этой транзакции -
например, это отправка письма или вызов чужого API - inbox не поможет,
нужен outbox и повторная попытка на стороне получателя.
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, ... }
]
Нужную запись ищут по имени очереди и причине - как в retryCount
из раздела про poison messages ниже.
И главное: x-death появляется только при dead-lettering. Если отвергать
сообщение с requeue=true, оно вернётся в очередь напрямую, минуя DLX, -
и счётчик не изменится никогда. Про эту ловушку - в разделе про poison
messages ниже.
Retry с DLQ
Паттерн «retry через 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;
}
}
Но такой возврат не считается попыткой и не имеет задержки, поэтому
ограничивать его надо своим счётчиком в памяти 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 крутит сообщение бесконечно, и оба факта зафиксированы тестами.
Мини-задание
- Создай таблицу
inboxи оберни свой handler вIdempotentHandler - Отправь одно и то же сообщение дважды - убедись, что обработка произошла один раз
- Настрой DLQ для своей очереди через
x-dead-letter-exchange - Сделай handler, который всегда возвращает ошибку, и убедись, что после 3 ретраев сообщение попало в DLQ