HardПрактика31 min

Очереди сообщений: сравнение и выбор

RabbitMQ, Kafka, SQS, NATS, Redis Streams -- семантики доставки, DLQ, retry, ordering. Примеры на Symfony Messenger и Go

Зачем нужны очереди сообщений

Очередь сообщений (Message Queue) -- промежуточный слой между отправителем и получателем. Вместо прямого вызова сервис кладёт сообщение в очередь, а потребитель забирает и обрабатывает его асинхронно.

Без очереди:                        С очередью:

┌──────────┐   sync    ┌──────────┐  ┌──────────┐  async  ┌───────┐  async  ┌──────────┐
│ Order    │──────────>│ Payment  │  │ Order    │───────>│ Queue │───────>│ Payment  │
│ Service  │<──────────│ Service  │  │ Service  │        │       │        │ Service  │
└──────────┘  timeout? └──────────┘  └──────────┘        └───────┘        └──────────┘
     │                                    │                                     │
     │ блокирует поток                    │ мгновенный ответ                    │ обработка
     │ если Payment упал --               │ клиенту (202 Accepted)              │ в своём темпе
     │ Order тоже падает                  │                                     │

Четыре ключевые задачи

Задача Описание Пример
Decoupling Сервисы не знают друг о друге Order не импортирует Payment
Async processing Тяжёлые задачи вынесены из запроса Отправка email, генерация PDF
Load leveling Очередь буферизирует всплески Чёрная пятница: 10x трафик
Retry / Reliability Повторная обработка при сбоях Платёж не прошёл -- retry через 30s

Два подхода: Message Queue vs Event Stream

Системы обмена сообщениями делятся на два фундаментально разных подхода:

Message Queue (Pull & Delete)           Event Stream (Append-only Log)
─────────────────────────────           ────────────────────────────────
┌──────────┐    ┌───────────┐           ┌──────────┐    ┌──────────────────────┐
│ Producer │───>│  Queue    │           │ Producer │───>│ Topic / Partition    │
└──────────┘    │           │           └──────────┘    │ [0][1][2][3][4][5].. │
                │ [msg][msg]│                           │   ^         ^        │
                │           │                           │   │         │        │
                └─────┬─────┘                           └───┼─────────┼────────┘
                      │                                     │         │
                 ┌────▼────┐                           ┌────┴──┐ ┌────┴──┐
                 │Consumer │                           │Group A│ │Group B│
                 │ ACK -->│                           │offset3│ │offset1│
                 │ delete  │                           └───────┘ └───────┘
                 └─────────┘
                                                    Сообщения не удаляются!
 Сообщение удаляется после ACK.                     Каждая группа читает
 Один потребитель получает                          с собственного offset.
 каждое сообщение.                                  Можно перечитать заново.
Характеристика Message Queue Event Stream
Модель Competing consumers Consumer groups + offsets
После обработки Сообщение удаляется Сообщение остаётся (retention)
Replay Невозможен Читай с любого offset
Ordering Per-queue (один поток) Per-partition
Типичные системы RabbitMQ, SQS, NATS Core Kafka, Redpanda, Redis Streams

Сравнение систем

RabbitMQ (AMQP)

RabbitMQ -- классический message broker с мощной маршрутизацией через exchanges и bindings.

                        ┌──────────────────────────────────────┐
                        │            RabbitMQ Broker            │
                        │                                      │
┌──────────┐            │  ┌──────────┐    ┌────────────────┐  │  ┌──────────┐
│ Producer │──publish──>│  │ Exchange │───>│ Queue: orders  │──┼─>│Consumer 1│
│          │            │  │ (topic)  │    └────────────────┘  │  └──────────┘
└──────────┘            │  │          │    ┌────────────────┐  │  ┌──────────┐
                        │  │          │───>│ Queue: emails  │──┼─>│Consumer 2│
                        │  └──────────┘    └────────────────┘  │  └──────────┘
                        │                  ┌────────────────┐  │
                        │                  │ Queue: DLQ     │  │  (poison msgs)
                        │                  └────────────────┘  │
                        └──────────────────────────────────────┘

Ключевые концепции:

  • Exchange -- маршрутизатор: принимает сообщение и решает, в какую очередь направить
  • Binding -- правило маршрутизации (routing key pattern)
  • Типы exchange: direct, topic, fanout, headers
  • ACK/NACK -- подтверждение обработки или отказ
  • Dead Letter Queue (DLQ) -- очередь для необработанных сообщений
  • Publisher Confirms -- брокер подтверждает, что сообщение записано на диск

Apache Kafka

Распределённый лог событий с высочайшей пропускной способностью. Сообщения хранятся на диске и не удаляются после чтения.

  • Topics + Partitions -- горизонтальное масштабирование
  • Consumer Groups -- каждая группа получает все сообщения, внутри группы -- параллелизм
  • Offsets -- потребитель управляет своей позицией
  • Log Compaction -- хранение последнего состояния по ключу
  • Kafka Transactions -- exactly-once семантика (EOS)

AWS SQS / SNS

Полностью управляемые сервисы Amazon:

  • SQS Standard -- at-least-once, best-effort ordering, почти безлимитный throughput
  • SQS FIFO -- exactly-once processing, строгий порядок, до 3000 msg/s (с batching)
  • SNS -- pub/sub: один SNS topic рассылает в несколько SQS очередей, Lambda, HTTP
  • Dead-Letter Queue -- встроенная через RedrivePolicy
  • Нет сервера для обслуживания -- pay-per-request

NATS

Легковесный, высокопроизводительный брокер. Написан на Go.

  • NATS Core -- at-most-once, pub/sub, request/reply
  • JetStream -- persistence, at-least-once, consumer groups, replay
  • Простота -- один бинарник, минимум конфигурации
  • Скорость -- десятки миллионов msg/s на кластере
  • Когда подходит -- IoT, edge computing, микросервисы с допустимой потерей

Redis Streams

Встроенная в Redis структура данных для потоковой обработки.

  • XADD / XREAD -- запись и чтение
  • Consumer Groups -- как в Kafka, с pending entries list (PEL)
  • XACK -- подтверждение обработки
  • XTRIM / MAXLEN -- ограничение размера потока
  • Когда подходит -- Redis уже в стеке, умеренная нагрузка, допустима потеря при OOM

Сравнительная таблица

Критерий RabbitMQ Kafka SQS/SNS NATS Redis Streams
Throughput ~50K msg/s ~1M+ msg/s ~∞ (managed) ~10M+ msg/s ~100K msg/s
Ordering Per-queue Per-partition FIFO: per-group Per-subject Per-stream
Persistence Disk (опционально) Disk (всегда) Managed JetStream: disk AOF/RDB
Протокол AMQP 0-9-1 Binary (TCP) HTTPS / SDK NATS protocol RESP (Redis)
Маршрутизация Гибкая (exchanges) Только topics SNS filters Subject-based Нет
Replay Нет Да (offsets) Нет JetStream: да Да (по ID)
Операционная сложность Средняя Высокая Нулевая Низкая Низкая
Exactly-once Нет (at-least) Да (транзакции) FIFO: да Нет Нет
Лицензия MPL 2.0 Apache 2.0 Проприетарная Apache 2.0 BSD

Семантики доставки

Одна из важнейших характеристик любой системы обмена сообщениями -- гарантии доставки.

At-most-once (Fire and Forget)

Producer ──send──> Broker ──deliver──> Consumer
                                         │
                                    обработал?
                                    не важно!
  • Сообщение отправлено без подтверждения со стороны потребителя
  • При сбое consumer -- сообщение потеряно навсегда
  • Максимальная скорость, минимальная задержка
  • Применение: метрики, логи, аналитика в реальном времени

At-least-once (ACK after Processing)

Producer ──send──> Broker ──deliver──> Consumer
                                         │
                                    обработал ✓
                                         │
                              Consumer ──ACK──> Broker
                                                  │
                                             удалить msg
  • Consumer отправляет ACK только после успешной обработки
  • При сбое до ACK -- брокер доставит повторно (duplicate!)
  • Consumer ОБЯЗАН быть идемпотентным
  • Применение: платежи, заказы, любая бизнес-логика

Exactly-once (Transactional)

┌─────────────────────────────────────────────────┐
│                   Стратегия 1:                   │
│          Transactional Outbox + Dedup            │
│                                                  │
│  ┌────────┐  TX  ┌──────────────────┐            │
│  │ App DB │<────>│ outbox table     │            │
│  │ orders │      │ (msg_id, payload)│            │
│  └────────┘      └───────┬──────────┘            │
│                          │ relay                  │
│                    ┌─────▼─────┐                  │
│                    │  Broker   │                  │
│                    └─────┬─────┘                  │
│                    ┌─────▼──────────┐             │
│                    │ Consumer       │             │
│                    │ dedup by msg_id│             │
│                    └────────────────┘             │
│                                                  │
│            Стратегия 2: Kafka Transactions        │
│  Producer.initTransactions()                     │
│  Producer.beginTransaction()                     │
│  Producer.send(record)                           │
│  Producer.sendOffsetsToTransaction(offsets)       │
│  Producer.commitTransaction()                    │
└─────────────────────────────────────────────────┘
  • Сообщение обработано ровно один раз, без дубликатов
  • Достигается через: Transactional Outbox + идемпотентный consumer с дедупликацией
  • Или через Kafka Transactions (producer + consumer в одной транзакции)
  • Самая сложная семантика, наибольший overhead
  • Применение: финансовые транзакции, критичные данные

Какая система что поддерживает

Семантика RabbitMQ Kafka SQS NATS Core NATS JetStream Redis Streams
At-most-once ✓ (no-ack) ✓ — ✓ (default) ✓ ✓
At-least-once ✓ (ack) ✓ ✓ (default) — ✓ (default) ✓ (XACK)
Exactly-once — (app-level) ✓ (native) ✓ (FIFO dedup) — — (app-level) — (app-level)

Dead Letter Queue (DLQ)

DLQ -- специальная очередь для сообщений, которые не удалось обработать после нескольких попыток. Poison messages -- сообщения, которые никогда не будут обработаны успешно (невалидные данные, баг в обработчике).

Основная очередь                    DLQ
┌────────────────┐                  ┌────────────────┐
│ [msg1] [msg2]  │   3 retries     │ [msg_poison]   │
│ [msg3]         │──────fail──────>│ [msg_corrupt]  │
│                │                  │                │
└───────┬────────┘                  └───────┬────────┘
        │                                   │
   ┌────▼────┐                         ┌────▼────────┐
   │Consumer │                         │ Admin Panel  │
   │ ACK ✓   │                         │ Inspect      │
   └─────────┘                         │ Fix & Replay │
                                       └──────────────┘

Настройка DLQ в RabbitMQ

DLQ в RabbitMQ реализуется через x-dead-letter-exchange и x-dead-letter-routing-key:

<?php

declare(strict_types=1);

use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Exchange\AMQPExchangeType;

/**
 * RabbitMQ DLQ configuration example
 */
final class RabbitMqDlqSetup
{
    public function configure(AMQPStreamConnection $connection): void
    {
        $channel = $connection->channel();

        // Declare DLQ exchange and queue
        $channel->exchange_declare(
            exchange: 'dlx.exchange',
            type: AMQPExchangeType::DIRECT,
            durable: true,
        );

        $channel->queue_declare(
            queue: 'orders.dlq',
            durable: true,
        );

        $channel->queue_bind(
            queue: 'orders.dlq',
            exchange: 'dlx.exchange',
            routing_key: 'orders',
        );

        // Declare main queue with DLQ settings
        $channel->queue_declare(
            queue: 'orders',
            durable: true,
            arguments: [
                'x-dead-letter-exchange' => ['S', 'dlx.exchange'],
                'x-dead-letter-routing-key' => ['S', 'orders'],
                'x-message-ttl' => ['I', 60000],        // 60s TTL
                'x-max-delivery-count' => ['I', 3],      // max 3 attempts (RabbitMQ 3.10+)
            ],
        );

        $channel->exchange_declare(
            exchange: 'orders.exchange',
            type: AMQPExchangeType::DIRECT,
            durable: true,
        );

        $channel->queue_bind(
            queue: 'orders',
            exchange: 'orders.exchange',
            routing_key: 'orders.created',
        );
    }
}
### Настройка DLQ в AWS SQS

В SQS DLQ настраивается через RedrivePolicy:

{
    "QueueName": "orders-queue",
    "Attributes": {
        "RedrivePolicy": "{\"deadLetterTargetArn\":\"arn:aws:sqs:eu-west-1:123456:orders-dlq\",\"maxReceiveCount\":\"3\"}"
    }
}
  • maxReceiveCount -- после стольких неудачных receive сообщение уходит в DLQ
  • deadLetterTargetArn -- ARN целевой DLQ очереди

Мониторинг и Reprocessing DLQ

Мониторинг:
┌─────────┐     ┌──────────┐     ┌───────────────┐
│   DLQ   │────>│CloudWatch│────>│ Alert (Slack)  │
│ depth>0 │     │ Alarm    │     │ "3 msgs in DLQ"│
└─────────┘     └──────────┘     └───────────────┘

Reprocessing:
┌─────────┐  read   ┌───────────┐  fix & publish  ┌──────────────┐
│   DLQ   │───────>│ Admin CLI  │────────────────>│ Main Queue   │
│         │        │ inspect    │                  │ (повторно)   │
└─────────┘        └───────────┘                  └──────────────┘

Правила мониторинга DLQ:

  • Алерт при depth > 0 -- любое сообщение в DLQ требует внимания
  • Дашборд -- график заполнения DLQ по времени
  • Reprocessing -- ручной или автоматический: прочитать, исправить, отправить обратно
  • Retention -- DLQ тоже имеет TTL (обычно 14 дней), не забыть обработать

Message Ordering

Глобальный vs Per-Partition Ordering

Глобальный порядок (1 partition / 1 queue):

  [msg1] → [msg2] → [msg3] → [msg4] → [msg5]
  ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
  Все сообщения строго упорядочены.
  Один потребитель. Нет параллелизма.

Per-partition ordering (3 partitions):

  Partition 0: [user_A:msg1] → [user_A:msg4]    ← порядок ✓
  Partition 1: [user_B:msg2] → [user_B:msg5]    ← порядок ✓
  Partition 2: [user_C:msg3] → [user_C:msg6]    ← порядок ✓

  3 потребителя параллельно!
  Порядок внутри каждого ключа гарантирован.

Ordering Keys

Система Механизм Пример
Kafka Partition Key producer.send(topic, key=userId, value)
SQS FIFO MessageGroupId SendMessage(MessageGroupId=orderId)
NATS JetStream Subject prefix orders.user123.created
Redis Streams Один stream = порядок Разделение по разным stream-ам

Когда ordering не нужен

  • Независимые задачи -- отправка email, ресайз картинок
  • Идемпотентные операции -- обновление кеша, пересчёт статистики
  • Fan-out уведомления -- push на все устройства пользователя

Правило: Не требуйте ordering без необходимости. Ordering = ограничение параллелизма.

Retry Strategies

Immediate Retry vs Exponential Backoff

Immediate retry:
  attempt 1 ──fail──> retry ──fail──> retry ──fail──> DLQ
             0s              0s              0s

Exponential backoff:
  attempt 1 ──fail──> wait 1s ──fail──> wait 4s ──fail──> wait 16s ──> DLQ
                      (1s)               (4s)               (16s)

Exponential backoff + jitter:
  attempt 1 ──fail──> wait 1.2s ──fail──> wait 3.7s ──fail──> wait 17.1s ──> DLQ
                      (random)            (random)             (random)

Delayed Retry в RabbitMQ

RabbitMQ не имеет встроенных delayed retry. Реализация через DLX + TTL:

┌──────────┐ nack  ┌──────────────┐ TTL=5s ┌──────────────┐ expire ┌──────────┐
│  Main    │──────>│ Retry Queue  │───────>│ DLX routes   │──────>│  Main    │
│  Queue   │       │ (x-msg-ttl)  │        │ back to main │       │  Queue   │
└──────────┘       └──────────────┘        └──────────────┘       └──────────┘
                   TTL: 5s, 30s, 120s (3 retry queues for different delays)

Circuit Breaker + Retry Budget

<?php

declare(strict_types=1);

/**
 * Circuit breaker prevents cascading failures when downstream is unhealthy.
 */
final class CircuitBreaker
{
    private int $failures = 0;
    private float $lastFailure = 0;

    public function __construct(
        private readonly int $threshold = 5,
        private readonly int $resetTimeoutSeconds = 60,
    ) {}

    public function isOpen(): bool
    {
        if ($this->failures < $this->threshold) {
            return false;
        }

        // After timeout -- allow one probe request (half-open)
        if ((microtime(true) - $this->lastFailure) > $this->resetTimeoutSeconds) {
            return false;
        }

        return true;
    }

    public function recordSuccess(): void
    {
        $this->failures = 0;
    }

    public function recordFailure(): void
    {
        $this->failures++;
        $this->lastFailure = microtime(true);
    }
}

/**
 * Retry budget limits total retry percentage to prevent thundering herd.
 * Rule: retries should not exceed 10% of successful requests.
 */
final class RetryBudget
{
    private int $totalRequests = 0;
    private int $retryCount = 0;

    public function __construct(
        private readonly float $maxRetryRatio = 0.1, // 10%
        private readonly int $minRetriesPerSecond = 3,
    ) {}

    public function canRetry(): bool
    {
        if ($this->totalRequests === 0) {
            return true;
        }

        // Always allow minimum retries
        if ($this->retryCount < $this->minRetriesPerSecond) {
            return true;
        }

        return ($this->retryCount / $this->totalRequests) < $this->maxRetryRatio;
    }

    public function recordRequest(): void
    {
        $this->totalRequests++;
    }

    public function recordRetry(): void
    {
        $this->retryCount++;
    }
}
## Symfony Messenger: полный пример

Symfony Messenger -- абстракция над транспортами (AMQP, Doctrine, Redis, Amazon SQS). Один интерфейс -- разные брокеры.

Конфигурация

# config/packages/messenger.yaml
framework:
    messenger:
        # Serialization
        serializer:
            default_serializer: messenger.transport.symfony_serializer

        # Transports
        transports:
            async:
                dsn: '%env(MESSENGER_TRANSPORT_DSN)%'   # amqp://guest:guest@rabbitmq:5672/%2f/messages
                options:
                    exchange:
                        name: app.exchange
                        type: topic
                    queues:
                        orders:
                            binding_keys: ['order.#']
                        notifications:
                            binding_keys: ['notification.#']
                retry_strategy:
                    max_retries: 3
                    delay: 1000          # 1 second
                    multiplier: 4        # 1s → 4s → 16s
                    max_delay: 60000     # max 60 seconds

            async_priority_high:
                dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
                options:
                    queues:
                        high_priority: ~
                retry_strategy:
                    max_retries: 5
                    delay: 500

            failed:
                dsn: 'doctrine://default?queue_name=failed'

        # Failure transport (DLQ)
        failure_transport: failed

        # Routing: Message class → Transport
        routing:
            'App\Message\OrderCreated': async
            'App\Message\SendNotification': async
            'App\Message\ProcessPayment': async_priority_high

Message + Handler

<?php

declare(strict_types=1);

namespace App\Message;

/**
 * Dispatched when a new order is placed.
 * Immutable message — all data set at construction.
 */
final readonly class OrderCreated
{
    public function __construct(
        public string $orderId,
        public string $userId,
        public float $totalAmount,
        public string $currency = 'USD',
        public \DateTimeImmutable $occurredAt = new \DateTimeImmutable(),
    ) {}
}
### Управление Failure Transport (DLQ)
# Просмотр сообщений в DLQ
php bin/console messenger:failed:show

# Детали конкретного сообщения
php bin/console messenger:failed:show 42

# Повторная обработка одного сообщения
php bin/console messenger:failed:retry 42

# Повторная обработка всех
php bin/console messenger:failed:retry

# Удалить из DLQ
php bin/console messenger:failed:remove 42

Go: RabbitMQ Consumer с ручным ACK

package consumer

import (
	"context"
	"encoding/json"
	"fmt"
	"log/slog"
	"time"

	amqp "github.com/rabbitmq/amqp091-go"
)

// OrderCreated represents an incoming order event.
type OrderCreated struct {
	OrderID     string    `json:"orderId"`
	UserID      string    `json:"userId"`
	TotalAmount float64   `json:"totalAmount"`
	Currency    string    `json:"currency"`
	OccurredAt  time.Time `json:"occurredAt"`
}

// RabbitConsumer processes messages from RabbitMQ with manual acknowledgment.
type RabbitConsumer struct {
	conn    *amqp.Connection
	channel *amqp.Channel
	logger  *slog.Logger
}

func NewRabbitConsumer(url string, logger *slog.Logger) (*RabbitConsumer, error) {
	conn, err := amqp.Dial(url)
	if err != nil {
		return nil, fmt.Errorf("dial rabbitmq: %w", err)
	}

	ch, err := conn.Channel()
	if err != nil {
		return nil, fmt.Errorf("open channel: %w", err)
	}

	// Prefetch: process one message at a time
	if err := ch.Qos(1, 0, false); err != nil {
		return nil, fmt.Errorf("set qos: %w", err)
	}

	return &RabbitConsumer{conn: conn, channel: ch, logger: logger}, nil
}

// Consume starts consuming messages from the given queue.
// Blocks until context is cancelled.
func (c *RabbitConsumer) Consume(ctx context.Context, queue string, handler func(OrderCreated) error) error {
	msgs, err := c.channel.Consume(
		queue,
		"",    // consumer tag (auto-generated)
		false, // auto-ack: false = manual ACK
		false, // exclusive
		false, // no-local
		false, // no-wait
		nil,
	)
	if err != nil {
		return fmt.Errorf("consume: %w", err)
	}

	c.logger.Info("consumer started", "queue", queue)

	for {
		select {
		case <-ctx.Done():
			c.logger.Info("consumer stopping", "reason", ctx.Err())
			return ctx.Err()

		case msg, ok := <-msgs:
			if !ok {
				return fmt.Errorf("channel closed")
			}

			var order OrderCreated
			if err := json.Unmarshal(msg.Body, &order); err != nil {
				c.logger.Error("invalid message payload", "error", err)
				// Reject without requeue — goes to DLQ
				_ = msg.Nack(false, false)
				continue
			}

			if err := handler(order); err != nil {
				c.logger.Error("handler failed",
					"orderId", order.OrderID,
					"error", err,
				)
				// Nack with requeue: broker will redeliver
				_ = msg.Nack(false, true)
				continue
			}

			// Successful processing — acknowledge
			if err := msg.Ack(false); err != nil {
				c.logger.Error("ack failed", "error", err)
			}

			c.logger.Info("order processed",
				"orderId", order.OrderID,
				"amount", order.TotalAmount,
			)
		}
	}
}

func (c *RabbitConsumer) Close() error {
	if err := c.channel.Close(); err != nil {
		return fmt.Errorf("close channel: %w", err)
	}
	return c.conn.Close()
}
## Kafka Consumer с Consumer Group
package consumer

import (
	"context"
	"encoding/json"
	"fmt"
	"log/slog"
	"time"

	"github.com/segmentio/kafka-go"
)

// KafkaConsumer reads messages from Kafka topic using consumer group.
type KafkaConsumer struct {
	reader *kafka.Reader
	logger *slog.Logger
}

// NewKafkaConsumer creates a consumer with consumer group for parallel processing.
func NewKafkaConsumer(brokers []string, topic, groupID string, logger *slog.Logger) *KafkaConsumer {
	reader := kafka.NewReader(kafka.ReaderConfig{
		Brokers:        brokers,
		Topic:          topic,
		GroupID:        groupID,
		MinBytes:       1,              // 1 byte
		MaxBytes:       10 * 1024 * 1024, // 10 MB
		MaxWait:        500 * time.Millisecond,
		CommitInterval: time.Second,     // Auto-commit interval
		StartOffset:    kafka.FirstOffset,
	})

	return &KafkaConsumer{reader: reader, logger: logger}
}

// Event is a generic event read from Kafka.
type Event struct {
	Key       string
	Value     json.RawMessage
	Topic     string
	Partition int
	Offset    int64
	Timestamp time.Time
}

// Consume reads messages until context is cancelled.
// Handler receives deserialized events. Offset is committed
// only after successful processing (at-least-once).
func (c *KafkaConsumer) Consume(ctx context.Context, handler func(ctx context.Context, event Event) error) error {
	c.logger.Info("kafka consumer started",
		"topic", c.reader.Config().Topic,
		"group", c.reader.Config().GroupID,
	)

	for {
		// FetchMessage does NOT auto-commit offset
		msg, err := c.reader.FetchMessage(ctx)
		if err != nil {
			if ctx.Err() != nil {
				return ctx.Err()
			}
			c.logger.Error("fetch failed", "error", err)
			continue
		}

		event := Event{
			Key:       string(msg.Key),
			Value:     msg.Value,
			Topic:     msg.Topic,
			Partition: msg.Partition,
			Offset:    msg.Offset,
			Timestamp: msg.Time,
		}

		c.logger.Info("message received",
			"partition", msg.Partition,
			"offset", msg.Offset,
			"key", string(msg.Key),
		)

		// Process message
		if err := handler(ctx, event); err != nil {
			c.logger.Error("handler error",
				"partition", msg.Partition,
				"offset", msg.Offset,
				"error", err,
			)
			// Do not commit — message will be redelivered on restart
			continue
		}

		// Commit offset only after successful processing
		if err := c.reader.CommitMessages(ctx, msg); err != nil {
			c.logger.Error("commit failed", "error", err)
		}
	}
}

func (c *KafkaConsumer) Close() error {
	return c.reader.Close()
}

// Producer demonstrates sending messages with partition key.
type Producer struct {
	writer *kafka.Writer
}

func NewProducer(brokers []string, topic string) *Producer {
	return &Producer{
		writer: &kafka.Writer{
			Addr:     kafka.TCP(brokers...),
			Topic:    topic,
			Balancer: &kafka.Hash{}, // Same key → same partition
		},
	}
}

// Send publishes a message with partition key for ordering guarantee.
func (p *Producer) Send(ctx context.Context, key string, payload any) error {
	data, err := json.Marshal(payload)
	if err != nil {
		return fmt.Errorf("marshal: %w", err)
	}

	return p.writer.WriteMessages(ctx, kafka.Message{
		Key:   []byte(key),   // Partition key: all events for this key go to same partition
		Value: data,
		Time:  time.Now(),
	})
}

func (p *Producer) Close() error {
	return p.writer.Close()
}
## Idempotent Consumer Pattern

Поскольку at-least-once доставка может привести к дубликатам, consumer обязан быть идемпотентным:

<?php

declare(strict_types=1);

namespace App\Service;

use Doctrine\DBAL\Connection;

/**
 * Idempotent message processor using deduplication table.
 * Prevents double-processing on redelivery.
 */
final readonly class IdempotentProcessor
{
    public function __construct(
        private Connection $db,
    ) {}

    /**
     * Process message only if not already processed.
     * Uses DB transaction: dedup check + business logic are atomic.
     *
     * @param string $messageId Unique message identifier
     * @param callable $handler Business logic to execute
     */
    public function processOnce(string $messageId, callable $handler): void
    {
        $this->db->transactional(function () use ($messageId, $handler): void {
            // Check if already processed (SELECT FOR UPDATE prevents race)
            $exists = $this->db->fetchOne(
                'SELECT 1 FROM processed_messages WHERE message_id = ? FOR UPDATE',
                [$messageId],
            );

            if ($exists !== false) {
                return; // Already processed — skip
            }

            // Execute business logic
            $handler();

            // Mark as processed
            $this->db->executeStatement(
                'INSERT INTO processed_messages (message_id, processed_at) VALUES (?, NOW())',
                [$messageId],
            );
        });
    }
}

// Usage in handler:
// $this->idempotentProcessor->processOnce(
//     $message->orderId,
//     fn() => $this->processOrder($message),
// );
## Transactional Outbox Pattern

Гарантирует, что запись в БД и отправка сообщения либо оба произойдут, либо оба нет:

┌─────────────────────────────────────────────────────┐
│                  Database (1 TX)                     │
│                                                     │
│  ┌─────────────┐      ┌───────────────────────┐     │
│  │ orders      │      │ outbox                │     │
│  │ ──────────  │      │ ───────────────────── │     │
│  │ id: ord_123 │      │ id: msg_456           │     │
│  │ status: new │      │ aggregate_type: order │     │
│  │ amount: 99  │      │ event_type: created   │     │
│  └─────────────┘      │ payload: {...}        │     │
│                       │ published: false      │     │
│                       └───────────────────────┘     │
└─────────────────────────────────────────────────────┘
          │
          │  Relay process (CDC / Polling)
          ▼
    ┌───────────┐
    │  Broker   │ ← published: false → true
    │ (RabbitMQ)│
    └───────────┘
<?php

declare(strict_types=1);

namespace App\Service;

use App\Entity\Order;
use App\Entity\OutboxMessage;
use Doctrine\ORM\EntityManagerInterface;

/**
 * Transactional outbox: business write + message in single DB transaction.
 * A relay process (cron/daemon) publishes outbox messages to the broker.
 */
final readonly class OrderService
{
    public function __construct(
        private EntityManagerInterface $em,
    ) {}

    public function createOrder(string $userId, float $amount, string $currency): Order
    {
        return $this->em->wrapInTransaction(function () use ($userId, $amount, $currency): Order {
            // Business logic: create order
            $order = new Order(
                userId: $userId,
                amount: $amount,
                currency: $currency,
            );
            $this->em->persist($order);

            // Outbox: store message in same transaction
            $outbox = new OutboxMessage(
                aggregateType: 'order',
                aggregateId: $order->getId(),
                eventType: 'order.created',
                payload: json_encode([
                    'orderId' => $order->getId(),
                    'userId' => $userId,
                    'amount' => $amount,
                    'currency' => $currency,
                ], JSON_THROW_ON_ERROR),
            );
            $this->em->persist($outbox);

            return $order;
        });
    }
}

/**
 * Relay: reads unpublished outbox messages and sends to broker.
 * Runs as a daemon or cron job.
 */
final readonly class OutboxRelay
{
    public function __construct(
        private EntityManagerInterface $em,
        private \Symfony\Component\Messenger\MessageBusInterface $bus,
    ) {}

    public function publishPending(int $batchSize = 100): int
    {
        $messages = $this->em->getRepository(OutboxMessage::class)
            ->findBy(['published' => false], ['createdAt' => 'ASC'], $batchSize);

        $count = 0;
        foreach ($messages as $outbox) {
            $this->bus->dispatch(new \App\Message\OutboxEvent(
                eventType: $outbox->getEventType(),
                payload: $outbox->getPayload(),
            ));

            $outbox->markPublished();
            $count++;
        }

        $this->em->flush();

        return $count;
    }
}
## Как выбрать: Decision Tree
                              Нужна очередь сообщений
                                       │
                           ┌───────────┴───────────┐
                           │                       │
                    Нужен replay?             Нет replay
                    Event sourcing?
                           │                       │
                      ┌────▼────┐          ┌───────▼────────┐
                      │  KAFKA  │          │  Managed в     │
                      │         │          │  облаке?       │
                      └─────────┘          │                │
                                      ┌────┴────┐    ┌─────▼──────┐
                                      │   Да    │    │    Нет     │
                                      │         │    │            │
                                  ┌───▼───┐     │ ┌──▼───────────┐
                                  │AWS SQS│     │ │ Сложная      │
                                  │  SNS  │     │ │ маршрутизация│
                                  └───────┘     │ │ (routing)?   │
                                                │ └──┬───────┬───┘
                                                │    │       │
                                                │   Да      Нет
                                                │    │       │
                                           ┌────▼────▼┐  ┌──▼──────────────┐
                                           │ RabbitMQ │  │  Допустима      │
                                           │          │  │  потеря msg?    │
                                           └──────────┘  └──┬──────────┬──┘
                                                            │          │
                                                           Да         Нет
                                                            │          │
                                                      ┌─────▼───┐ ┌───▼──────────┐
                                                      │  NATS   │ │ Redis уже    │
                                                      │  Core   │ │ в стеке?     │
                                                      └─────────┘ └──┬────────┬──┘
                                                                     │        │
                                                                    Да       Нет
                                                                     │        │
                                                               ┌─────▼────┐ ┌─▼────────┐
                                                               │  Redis   │ │ RabbitMQ │
                                                               │ Streams  │ │ или NATS │
                                                               └──────────┘ │ JetStream│
                                                                            └──────────┘

Краткая таблица выбора

Сценарий Рекомендация Причина
Event sourcing, replay, высокий throughput Kafka Immutable log, consumer groups, offset management
Сложная маршрутизация, priority queues RabbitMQ Exchanges, bindings, message TTL, priority
Облако AWS, zero ops SQS + SNS Managed, pay-per-use, интеграция с Lambda
Максимальная скорость, потери ОК NATS Core Минимальная задержка, простота
Persistence + скорость NATS JetStream Лёгкий Kafka-like функционал
Redis уже в проекте, простые задачи Redis Streams Нет нового сервиса, consumer groups

Message Flow: полная диаграмма

┌──────────┐        ┌──────────────────────────────────────────────────┐
│ Producer │        │                  Broker                          │
│          │        │                                                  │
│  1. serialize     │   2. route                3. store               │
│  message  │──────>│  ┌──────────┐  binding  ┌─────────────┐         │
│           │       │  │ Exchange │──────────>│ Queue /     │         │
│           │       │  │ / Topic  │           │ Partition   │         │
│           │       │  └──────────┘           │ [msg][msg]  │         │
│           │       │                         └──────┬──────┘         │
│           │       │                                │                │
│           │       │  6. DLQ                  4. deliver              │
│           │       │  ┌──────────┐  max_retry ┌─────▼──────┐         │
│           │       │  │   DLQ    │<───────────│  Consumer  │         │
│           │       │  │          │   failed   │            │         │
│           │       │  └──────────┘            │ 5. ACK/NACK│         │
│           │       │                          └────────────┘         │
└──────────┘        └──────────────────────────────────────────────────┘

Жизненный цикл сообщения:
1. Producer сериализует и отправляет (publish)
2. Broker маршрутизирует (exchange → binding → queue)
3. Broker сохраняет на диск (persistence)
4. Broker доставляет consumer-у (deliver / fetch)
5. Consumer обрабатывает и подтверждает (ACK) или отклоняет (NACK)
6. После N неудач — сообщение уходит в DLQ для ручного разбора

Итоги

  • Message Queue -- классические очереди с удалением после ACK (RabbitMQ, SQS)
  • Event Stream -- append-only лог с replay (Kafka, Redis Streams, NATS JetStream)
  • Семантики доставки: at-most-once, at-least-once, exactly-once — выбирайте осознанно
  • At-least-once + идемпотентный consumer -- самый практичный вариант для 90% задач
  • DLQ -- обязателен в production, мониторьте и обрабатывайте
  • Transactional Outbox -- гарантия консистентности между БД и брокером
  • Ordering -- требуйте только там, где бизнес-логика требует, иначе теряете параллелизм
  • Retry: exponential backoff + jitter + circuit breaker + retry budget

Проверь себя

5 из 8

Какую систему стоит выбрать, если нужна сложная маршрутизация сообщений (routing по паттернам, fanout, headers)?

Какой механизм в RabbitMQ используется для реализации delayed retry (повторная обработка через N секунд)?

Почему retry budget (бюджет повторов) важен в распределённых системах?

В Kafka, что гарантирует ordering (порядок) сообщений?

Для чего используется Dead Letter Queue (DLQ)?