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 обработанных сообщений:
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);
}
}
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» с задержкой:
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 через 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