HardТеория2 min

Azure Service Bus

Queues vs Topics, Dead Letter Queue, порядок сообщений, дедупликация

Service Bus -- корпоративный брокер сообщений для надежной асинхронной коммуникации между сервисами. Поддерживает расширенные возможности: транзакции, сессии, Dead Letter Queue, дедупликацию.

Queues vs Topics

Queues (очереди)

Модель point-to-point: один отправитель, один получатель. Сообщение обрабатывается одним получателем.

using Azure.Messaging.ServiceBus;

// Отправка сообщения
var client = new ServiceBusClient(connectionString);
var sender = client.CreateSender("order-queue");

var order = new OrderCreatedEvent { OrderId = "123", Total = 150m };
var message = new ServiceBusMessage(JsonSerializer.SerializeToUtf8Bytes(order))
{
    ContentType = "application/json",
    Subject = "OrderCreated",
    MessageId = order.OrderId, // Для дедупликации
    ScheduledEnqueueTime = DateTimeOffset.UtcNow.AddMinutes(5) // Отложенная доставка
};
await sender.SendMessageAsync(message);

Topics & Subscriptions (топики и подписки)

Модель publish-subscribe: один отправитель, несколько получателей. Каждая подписка получает копию сообщения.

// Отправка в топик
var sender = client.CreateSender("order-events");
var message = new ServiceBusMessage(body)
{
    Subject = "OrderCreated",
    ApplicationProperties =
    {
        { "Region", "EU" },
        { "Priority", "High" }
    }
};
await sender.SendMessageAsync(message);

Подписки могут иметь фильтры:

# Создание подписки с SQL-фильтром
az servicebus topic subscription rule create \
  --namespace-name sb-orderapp \
  --topic-name order-events \
  --subscription-name payment-processor \
  --name HighPriorityOnly \
  --filter-sql-expression "Priority = 'High'"

Получение сообщений

var processor = client.CreateProcessor("order-queue", new ServiceBusProcessorOptions
{
    MaxConcurrentCalls = 5,
    AutoCompleteMessages = false // Manual complete для надежности
});

processor.ProcessMessageAsync += async args =>
{
    var order = JsonSerializer.Deserialize<OrderCreatedEvent>(args.Message.Body);

    try
    {
        await ProcessOrderAsync(order);
        await args.CompleteMessageAsync(args.Message); // Подтверждение обработки
    }
    catch (TransientException)
    {
        // Не вызываем Complete -- сообщение вернется в очередь
        await args.AbandonMessageAsync(args.Message);
    }
    catch (PermanentException ex)
    {
        // Перемещаем в Dead Letter Queue
        await args.DeadLetterMessageAsync(args.Message,
            deadLetterReason: "ProcessingFailed",
            deadLetterErrorDescription: ex.Message);
    }
};

processor.ProcessErrorAsync += async args =>
{
    Console.WriteLine($"Error: {args.Exception.Message}");
};

await processor.StartProcessingAsync();

Dead Letter Queue (DLQ)

Сообщения попадают в DLQ когда:

  • Превышено maxDeliveryCount (по умолчанию 10 попыток)
  • Явный вызов DeadLetterMessageAsync
  • Истек TTL сообщения
  • Ошибка обработки подписчиком
// Обработка Dead Letter Queue
[Function("ProcessDeadLetters")]
public async Task RunDlq(
    [ServiceBusTrigger("order-queue/$deadletterqueue", Connection = "ServiceBus")]
    ServiceBusReceivedMessage message)
{
    _logger.LogWarning(
        "Dead letter: Reason={Reason}, Description={Description}",
        message.DeadLetterReason,
        message.DeadLetterErrorDescription);

    // Сохранить для ручного анализа
    await _failedMessageRepo.SaveAsync(new FailedMessage
    {
        MessageId = message.MessageId,
        Body = message.Body.ToString(),
        Reason = message.DeadLetterReason,
        ReceivedAt = DateTime.UtcNow
    });
}

Расширенные возможности

Дедупликация

Предотвращает обработку одного сообщения дважды. Включается на очереди/топике. Ключ -- свойство MessageId.

Message Sessions

Гарантирует упорядоченную обработку сообщений с одним SessionId. Полезно для обработки событий одного заказа в правильном порядке.

var message = new ServiceBusMessage(body)
{
    SessionId = orderId.ToString() // Все события одного заказа в одной сессии
};

Scheduled Messages

Отложенная доставка -- сообщение появится в очереди в указанное время.

Transactions

Атомарные операции с несколькими сообщениями:

using var scope = new TransactionScope(TransactionScopeAsyncFlowOption.Enabled);
await sender.SendMessageAsync(message1);
await sender.SendMessageAsync(message2);
scope.Complete();

Service Bus vs Queue Storage

Характеристика Service Bus Queue Storage
Размер сообщения До 256 KB (Standard) / 100 MB (Premium) До 64 KB
Порядок FIFO гарантирован (Sessions) Примерный FIFO
Дедупликация Встроенная Нет
Dead Letter Queue Встроенная Нет
Topics/Subscriptions Да Нет
Transactions Да Нет
Стоимость Выше Ниже

Проверь себя

В чем разница между Queue и Topic в Service Bus?

Когда сообщение автоматически попадает в Dead Letter Queue?

Как гарантировать порядок обработки сообщений одного заказа в Service Bus?

Какой метод вызывается для подтверждения успешной обработки сообщения?