Зачем нужны очереди сообщений
Очередь сообщений (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',
);
}
}
В 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 -- абстракция над транспортами (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(),
) {}
}
# Просмотр сообщений в 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()
}
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()
}
Поскольку 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),
// );
Гарантирует, что запись в БД и отправка сообщения либо оба произойдут, либо оба нет:
┌─────────────────────────────────────────────────────┐
│ 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;
}
}
Нужна очередь сообщений
│
┌───────────┴───────────┐
│ │
Нужен 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