Гарантии доставки: at-most-once / at-least-once / exactly-once
Гарантии доставки: at-most-once / at-least-once / exactly-once
Когда событие путешествует от publisher к subscriber через сеть, что-то может пойти не так. Сеть может оборваться, consumer может упасть в середине обработки, broker может потерять сообщение. Гарантии доставки - это ответ на вопрос: что произойдёт при сбое?
Этот вопрос критически важен для продакшна. Представь: subscriber начисляет бонусные баллы за пройденный урок. Если событие потеряется - пользователь не получит бонус. Если событие обработается дважды - получит двойной. Выбор гарантии доставки определяет, какой из этих сценариев возможен.
Три уровня гарантий
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 выполняются в одной транзакции - атомарно.
Сводная таблица
| Гарантия | Потеря событий | Дубликаты | Сложность | Когда использовать |
|---|---|---|---|---|
| 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 - что произойдёт?