Гарантии доставки: at-most-once / at-least-once / exactly-once

Гарантии доставки: at-most-once / at-least-once / exactly-once

Когда событие путешествует от publisher к subscriber через сеть, что-то может пойти не так. Сеть может оборваться, consumer может упасть в середине обработки, broker может потерять сообщение. Гарантии доставки - это ответ на вопрос: что произойдёт при сбое?

Этот вопрос критически важен для продакшна. Представь: subscriber начисляет бонусные баллы за пройденный урок. Если событие потеряется - пользователь не получит бонус. Если событие обработается дважды - получит двойной. Выбор гарантии доставки определяет, какой из этих сценариев возможен.

Три уровня гарантий

Три гарантии: at-most-once с потерями, at-least-once с дубликатами, exactly-once = at-least-once + idempotent

At-most-once: отправил и забыл

Producer публикует событие и не ждёт подтверждения. Если broker не получил сообщение или consumer упал - событие потеряно навсегда.

// At-most-once publisher: отправка без подтверждения
type FireAndForgetPublisher struct {
    conn net.Conn
}

func (p *FireAndForgetPublisher) Publish(ctx context.Context, event Event) error {
    data, err := json.Marshal(event)
    if err != nil {
        return fmt.Errorf("marshal event: %w", err)
    }

    // Отправляем и не проверяем, получил ли broker
    _, _ = p.conn.Write(data)
    return nil
}
<?php
// src/Infrastructure/Messaging/FireAndForgetPublisher.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

use App\Domain\Event\Event;
use PhpAmqpLib\Channel\AMQPChannel;
use PhpAmqpLib\Message\AMQPMessage;

// At-most-once publisher: отправка без подтверждения
final class FireAndForgetPublisher
{
    public function __construct(
        private readonly AMQPChannel $channel,
        private readonly string $exchange,
    ) {}

    public function publish(Event $event): void
    {
        $data = json_encode($event, JSON_THROW_ON_ERROR);
        $msg = new AMQPMessage($data);

        // Отправляем без publisher confirms - ack от брокера не ждём
        // (mandatory=false, immediate=false: даже если очередь не привязана, broker молча отбросит)
        $this->channel->basic_publish($msg, $this->exchange);
    }
}

Когда допустимо: метрики, аналитика, логирование - ситуации, где потеря одного события не приводит к бизнес-проблеме. Если потерялся один клик из миллиона - статистика не пострадает.

Когда не подходит: email-уведомления, начисление бонусов, обновление баланса - всё, где потеря имеет последствия для пользователя.

At-least-once: гарантия доставки через retry

Producer отправляет событие и ждёт подтверждения (ACK) от broker. Если ACK не пришёл - повторяет отправку. Consumer обрабатывает событие и отправляет ACK broker-у. Если consumer упал до ACK - broker доставит событие повторно.

// At-least-once publisher с retry и exponential backoff
type ReliablePublisher struct {
    writer     *kafka.Writer
    maxRetries int
    logger     *slog.Logger
}

func NewReliablePublisher(brokers []string, logger *slog.Logger) *ReliablePublisher {
    w := &kafka.Writer{
        Addr:         kafka.TCP(brokers...),
        RequiredAcks: kafka.RequireAll, // ждём подтверждения от всех реплик
        MaxAttempts:  1,                // retry делаем сами (для контроля backoff)
    }
    return &ReliablePublisher{
        writer:     w,
        maxRetries: 5,
        logger:     logger,
    }
}

func (p *ReliablePublisher) Publish(ctx context.Context, topic string, event Event) error {
    data, err := json.Marshal(event)
    if err != nil {
        return fmt.Errorf("marshal event: %w", err)
    }

    msg := kafka.Message{
        Topic: topic,
        Key:   []byte(event.AggregateID), // partition key для ordering
        Value: data,
    }

    var lastErr error
    for attempt := 0; attempt <= p.maxRetries; attempt++ {
        if attempt > 0 {
            delay := backoff(attempt)
            p.logger.Warn("retrying publish",
                slog.String("topic", topic),
                slog.String("event_id", event.ID),
                slog.Int("attempt", attempt),
                slog.Duration("delay", delay),
            )
            select {
            case <-time.After(delay):
            case <-ctx.Done():
                return ctx.Err()
            }
        }

        lastErr = p.writer.WriteMessages(ctx, msg)
        if lastErr == nil {
            return nil // ACK получен - событие доставлено
        }
    }

    return fmt.Errorf("publish failed after %d retries: %w", p.maxRetries, lastErr)
}

// backoff возвращает задержку с экспоненциальным ростом
// 100ms, 200ms, 400ms, 800ms, 1600ms
func backoff(attempt int) time.Duration {
    base := 100 * time.Millisecond
    return base * time.Duration(1<<uint(attempt-1))
}
<?php
// src/Infrastructure/Messaging/ReliablePublisher.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

use App\Domain\Event\Event;
use PhpAmqpLib\Channel\AMQPChannel;
use PhpAmqpLib\Message\AMQPMessage;
use Psr\Log\LoggerInterface;

// At-least-once publisher с retry и exponential backoff
final class ReliablePublisher
{
    public function __construct(
        private readonly AMQPChannel $channel,
        private readonly LoggerInterface $logger,
        private readonly int $maxRetries = 5,
    ) {
        // Включаем publisher confirms - broker отправит basic.ack/basic.nack
        $this->channel->confirm_select();
    }

    public function publish(string $exchange, string $routingKey, Event $event): void
    {
        $data = json_encode($event, JSON_THROW_ON_ERROR);

        $msg = new AMQPMessage($data, [
            'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
            'message_id' => $event->id,
        ]);

        $lastErr = null;
        for ($attempt = 0; $attempt <= $this->maxRetries; $attempt++) {
            if ($attempt > 0) {
                $delayMs = self::backoffMs($attempt);
                $this->logger->warning('retrying publish', [
                    'exchange' => $exchange,
                    'event_id' => $event->id,
                    'attempt' => $attempt,
                    'delay_ms' => $delayMs,
                ]);
                usleep($delayMs * 1000);
            }

            try {
                $this->channel->basic_publish($msg, $exchange, $routingKey);
                // Блокируемся, пока broker не подтвердит публикацию (publisher confirms)
                $this->channel->wait_for_pending_acks(timeout: 5.0);

                return; // ACK получен - событие доставлено
            } catch (\Throwable $e) {
                $lastErr = $e;
            }
        }

        throw new \RuntimeException(
            sprintf('publish failed after %d retries: %s', $this->maxRetries, $lastErr?->getMessage() ?? ''),
            previous: $lastErr,
        );
    }

    // backoffMs возвращает задержку в миллисекундах с экспоненциальным ростом
    // 100, 200, 400, 800, 1600
    private static function backoffMs(int $attempt): int
    {
        return 100 * (1 << ($attempt - 1));
    }
}

В Symfony Messenger вся retry-логика конфигурируется на уровне транспорта (retry_strategy.max_retries, multiplier, delay), а Kafka transport (symfony/messenger-kafka или enqueue/rdkafka) умеет ждать acks=all - ручной цикл retry в коде не нужен.

Побочный эффект at-least-once: дубликаты. Если consumer обработал событие, но упал до отправки ACK, broker не знает, что обработка завершена, и доставит событие повторно. Consumer обработает его второй раз.

Сценарии сбоев

Разберём, что происходит при падении consumer-а на разных этапах:

Сценарий 1: Consumer упал ДО обработки
─────────────────────────────────────────
Broker ──▶ Consumer (crash!)
Broker: "ACK не получен, доставлю повторно"
Broker ──▶ Consumer (перезапущен) ──▶ обрабатывает ──▶ ACK
Результат: событие обработано 1 раз ✅

Сценарий 2: Consumer упал ПОСЛЕ обработки, но ДО ACK
─────────────────────────────────────────
Broker ──▶ Consumer ──▶ обработал ──▶ (crash before ACK!)
Broker: "ACK не получен, доставлю повторно"
Broker ──▶ Consumer (перезапущен) ──▶ обработал СНОВА ──▶ ACK
Результат: событие обработано 2 раза ⚠️ ДУБЛИКАТ

Сценарий 2 - это именно тот случай, когда нужна идемпотентность consumer-а.

Kafka: auto-commit vs manual commit

В Kafka consumer читает сообщения из partition и коммитит offset (позицию). Offset указывает, до какого сообщения consumer дочитал.

// Auto-commit: offset коммитится автоматически через заданный интервал
// Проблема: если consumer упал между auto-commit и обработкой - сообщение потеряно
reader := kafka.NewReader(kafka.ReaderConfig{
    Brokers:        []string{"localhost:9092"},
    Topic:          "progress.lesson.completed",
    GroupID:        "email-sender",
    CommitInterval: 1 * time.Second, // коммит каждую секунду - рискованно
})

// Manual commit: offset коммитится только после успешной обработки
func consumeManual(ctx context.Context, reader *kafka.Reader) error {
    for {
        // FetchMessage НЕ коммитит offset
        msg, err := reader.FetchMessage(ctx)
        if err != nil {
            return fmt.Errorf("fetch: %w", err)
        }

        // Обработка сообщения
        if err := processMessage(ctx, msg); err != nil {
            slog.Error("process failed, will retry on next fetch",
                slog.String("err", err.Error()),
                slog.Int64("offset", msg.Offset),
            )
            continue // не коммитим - сообщение будет доставлено повторно
        }

        // Коммитим offset только после успешной обработки
        if err := reader.CommitMessages(ctx, msg); err != nil {
            return fmt.Errorf("commit: %w", err)
        }
    }
}
<?php
// src/Infrastructure/Messaging/KafkaManualCommitConsumer.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

use Psr\Log\LoggerInterface;
use RdKafka\Conf;
use RdKafka\KafkaConsumer;
use RdKafka\Message;

// Manual commit: offset коммитится только после успешной обработки.
// Auto-commit (enable.auto.commit=true) рискован: если consumer упал
// между авто-коммитом и обработкой - сообщение потеряно.
final class KafkaManualCommitConsumer
{
    private readonly KafkaConsumer $consumer;

    public function __construct(
        string $brokers,
        string $groupId,
        private readonly MessageProcessor $processor,
        private readonly LoggerInterface $logger,
    ) {
        $conf = new Conf();
        $conf->set('group.id', $groupId);
        $conf->set('bootstrap.servers', $brokers);
        $conf->set('enable.auto.commit', 'false'); // ОТКЛЮЧАЕМ авто-коммит
        $conf->set('auto.offset.reset', 'earliest');

        $this->consumer = new KafkaConsumer($conf);
    }

    public function run(string $topic): void
    {
        $this->consumer->subscribe([$topic]);

        while (true) {
            $msg = $this->consumer->consume(timeoutMs: 10_000);
            if ($msg === null || $msg->err !== RD_KAFKA_RESP_ERR_NO_ERROR) {
                continue;
            }

            try {
                $this->processor->process($msg);
            } catch (\Throwable $e) {
                $this->logger->error('process failed, will retry on next fetch', [
                    'err' => $e->getMessage(),
                    'offset' => $msg->offset,
                ]);
                continue; // не коммитим - сообщение будет доставлено повторно
            }

            // Коммитим offset только после успешной обработки (синхронный commit)
            $this->consumer->commit($msg);
        }
    }
}

В Symfony Messenger Kafka transport то же поведение настраивается опциями commit_async и enable.auto.commit в DSN транспорта - ручной цикл consume() обычно прячется внутри MessengerBundle worker-а.

Manual commit даёт at-least-once: если consumer упал после обработки, но до коммита, сообщение будет доставлено повторно. Auto-commit даёт at-most-once: если consumer упал после коммита, но до обработки, сообщение потеряно.

RabbitMQ: ack, nack, reject

В RabbitMQ подтверждение работает на уровне отдельных сообщений:

ack - "обработал успешно, удали из очереди"
nack - "не смог обработать, верни в очередь (requeue)" или "отправь в DLQ"
reject - "отклоняю это сообщение" (аналог nack без requeue)
// RabbitMQ consumer с manual ack
func consumeRabbit(ch *amqp.Channel, queueName string) error {
    msgs, err := ch.Consume(
        queueName,
        "",    // consumer tag
        false, // auto-ack = false (manual ack!)
        false, // exclusive
        false, // no-local
        false, // no-wait
        nil,
    )
    if err != nil {
        return fmt.Errorf("consume: %w", err)
    }

    for msg := range msgs {
        err := processMessage(msg.Body)
        if err != nil {
            slog.Error("processing failed",
                slog.String("err", err.Error()),
            )
            // Возвращаем в очередь для повторной обработки
            _ = msg.Nack(false, true) // multiple=false, requeue=true
            continue
        }

        // Подтверждаем успешную обработку
        _ = msg.Ack(false) // multiple=false
    }
    return nil
}
<?php
// src/Infrastructure/Messaging/RabbitConsumer.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

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

// RabbitMQ consumer с manual ack
final class RabbitConsumer
{
    public function __construct(
        private readonly AMQPChannel $channel,
        private readonly MessageProcessor $processor,
        private readonly LoggerInterface $logger,
    ) {}

    public function consume(string $queueName): void
    {
        $callback = function (AMQPMessage $msg): void {
            try {
                $this->processor->process($msg->getBody());
            } catch (\Throwable $e) {
                $this->logger->error('processing failed', ['err' => $e->getMessage()]);
                // Возвращаем в очередь для повторной обработки
                $msg->nack(requeue: true); // multiple=false (по умолчанию)

                return;
            }

            // Подтверждаем успешную обработку
            $msg->ack(); // multiple=false
        };

        $this->channel->basic_consume(
            queue: $queueName,
            consumer_tag: '',
            no_local: false,
            no_ack: false, // manual ack!
            exclusive: false,
            nowait: false,
            callback: $callback,
        );

        while ($this->channel->is_consuming()) {
            $this->channel->wait();
        }
    }
}

Exactly-once: миф или реальность?

Строгий exactly-once в распределённой системе невозможен (это следствие теоремы о двух генералах). То, что Kafka называет «exactly-once semantics» - это at-least-once delivery + idempotent producer + transactional consumer. По сути: сообщение может быть доставлено несколько раз, но эффект от обработки будет как от одной.

Практический подход - at-least-once + идемпотентный consumer:

// IdempotentConsumer проверяет, обработано ли уже событие
type IdempotentConsumer struct {
    db      *sql.DB
    handler func(ctx context.Context, event Event) error
    logger  *slog.Logger
}

func (c *IdempotentConsumer) Handle(ctx context.Context, event Event) error {
    // Проверяем: обрабатывали ли мы уже это событие?
    var exists bool
    err := c.db.QueryRowContext(ctx,
        "SELECT EXISTS(SELECT 1 FROM processed_events WHERE event_id = $1)",
        event.ID,
    ).Scan(&exists)
    if err != nil {
        return fmt.Errorf("check processed: %w", err)
    }

    if exists {
        c.logger.Info("event already processed, skipping",
            slog.String("event_id", event.ID),
        )
        return nil // идемпотентность: повторная обработка - просто skip
    }

    // Обрабатываем событие и записываем ID в одной транзакции
    tx, err := c.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    if err := c.handler(ctx, event); err != nil {
        return fmt.Errorf("handle event: %w", err)
    }

    _, err = tx.ExecContext(ctx,
        "INSERT INTO processed_events (event_id, processed_at) VALUES ($1, NOW())",
        event.ID,
    )
    if err != nil {
        return fmt.Errorf("save processed event: %w", err)
    }

    return tx.Commit()
}
<?php
// src/Infrastructure/Messaging/IdempotentConsumer.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

use App\Domain\Event\Event;
use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;

// IdempotentConsumer проверяет, обработано ли уже событие
final class IdempotentConsumer
{
    /**
     * @param  callable(Event): void  $handler
     */
    public function __construct(
        private readonly Connection $db,
        private readonly LoggerInterface $logger,
        private $handler,
    ) {}

    public function handle(Event $event): void
    {
        // Проверяем: обрабатывали ли мы уже это событие?
        $exists = (bool) $this->db->fetchOne(
            'SELECT EXISTS(SELECT 1 FROM processed_events WHERE event_id = :id)',
            ['id' => $event->id],
        );

        if ($exists) {
            $this->logger->info('event already processed, skipping', [
                'event_id' => $event->id,
            ]);

            return; // идемпотентность: повторная обработка - просто skip
        }

        // Обрабатываем событие и записываем ID в одной транзакции
        $this->db->beginTransaction();
        try {
            ($this->handler)($event);

            $this->db->executeStatement(
                'INSERT INTO processed_events (event_id, processed_at) VALUES (:id, NOW())',
                ['id' => $event->id],
            );

            $this->db->commit();
        } catch (\Throwable $e) {
            $this->db->rollBack();
            throw new \RuntimeException('handle event: ' . $e->getMessage(), previous: $e);
        }
    }
}

Таблица processed_events - это deduplication store. Она хранит ID обработанных событий. Если событие пришло повторно, consumer видит, что ID уже есть, и пропускает обработку. Бизнес-логика и запись ID выполняются в одной транзакции - атомарно.

Это золотой стандарт для 95% систем. At-most-once теряет события. Exactly-once дорогой и сложный. At-least-once + идемпотентный consumer даёт надёжность без чрезмерной сложности. Правило простое: каждый subscriber должен корректно обработать повторное получение одного и того же события. Если consumer падает между auto-commit и обработкой сообщения - сообщение потеряно. В продакшне всегда используй manual commit: `FetchMessage` → process → `CommitMessages`. Это единственный способ гарантировать at-least-once.

Сводная таблица

ГарантияПотеря событийДубликатыСложностьКогда использовать
At-most-onceДаНетНизкаяМетрики, аналитика, логи
At-least-onceНетДаСредняяEmail, бонусы (+ идемпотентность)
Exactly-onceНетНетВысокаяФинансовые операции, billing

Типичные ошибки

  • Auto-commit + долгая обработка - Kafka закоммитил offset до того, как ты успел сделать работу; consumer крашится → событие потеряно. Manual commit после успешной обработки (либо после записи в БД с idempotency).

  • Manual ack/commit перед обработкой - «сначала ack, потом работа» = at-most-once в одежде at-least-once. Acknowledgement - после успешного действия. Если упал перед ack - broker повторит.

  • «Exactly-once на уровне брокера, мне не нужна идемпотентность» - у Kafka Exactly-Once-Semantics работает только внутри Kafka-pipeline (Kafka→Kafka). Как только пишешь во внешнюю БД или дёргаешь HTTP - у тебя at-least-once, нужен idempotency-store.

  • Дедупликация по event.Type + aggregate_id - у двух событий одного типа на один aggregate один и тот же ключ → второе событие считается дублем и теряется. Дедуп только по event_id (уникальный UUID).

  • processed_events без cleanup - таблица растёт безгранично, через год запросы тормозят. TTL/партиционирование/архив (для большинства задач достаточно 7-30 дней - повторы за гранью ретеншна брокера невозможны).

  • Бесконечный retry без exponential backoff - consumer падает, broker сразу redeliver, consumer снова падает - retry storm загружает 100% CPU. Exponential backoff (100ms, 500ms, 2s, 10s...) + max retries + DLQ.

  • Записал в БД, и в defer отправляешь событие - defer сработает даже если транзакция откатилась. Событие уйдёт без реальных данных в БД. Публикация - после tx.Commit() или через Outbox (урок 7).

  • Python - Kafka с aiokafka: idempotent producer, exactly-once - практическая реализация гарантий доставки в aiokafka: транзакции, consumer commit

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

  • Реализуй publisher с retry и exponential backoff (3 попытки, начальная задержка 100ms)
  • Напиши consumer с manual commit: FetchMessage → process → CommitMessages
  • Создай таблицу processed_events и оберни subscriber в IdempotentConsumer
  • Для каждого subscriber в своём проекте определи: какая гарантия доставки нужна и почему
  • Смоделируй сценарий: consumer упал после обработки, но до ACK - что произойдёт?

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