RO EN

Azure Service Bus (2) — producători și consumatori în .NET

Azure Service Bus (2) — producători și consumatori în .NET ✨ Imagine generată cu AI
Doru Bulubașa
22 iulie 2026
48 vizualizări

A doua parte din seria despre Azure Service Bus. În prima parte am ales între Queue și Topic. Acum scriem producătorii și consumatorii în .NET.


Setup: pachet și autentificare

dotnet add package Azure.Messaging.ServiceBus
dotnet add package Azure.Identity

Consistent cu restul seriei: fără connection string. ServiceBusClient se autentifică prin Managed Identity:

# Rol pentru trimitere
az role assignment create \
  --assignee $PRINCIPAL_ID \
  --role "Azure Service Bus Data Sender" \
  --scope $SERVICEBUS_ID

# Rol pentru consum
az role assignment create \
  --assignee $PRINCIPAL_ID \
  --role "Azure Service Bus Data Receiver" \
  --scope $SERVICEBUS_ID
// Program.cs -- un singur ServiceBusClient, singleton
builder.Services.AddSingleton(_ => new ServiceBusClient(
    "my-servicebus.servicebus.windows.net",  // fully qualified namespace
    new DefaultAzureCredential()));

ServiceBusClient e thread-safe și gestionat intern cu connection pooling — îl înregistrezi o singură dată ca singleton și creezi sender-e/processor-e din el.


Producătorul

Publicare simplă

public class OrderEventPublisher
{
    private readonly ServiceBusSender _sender;

    public OrderEventPublisher(ServiceBusClient client)
    {
        // Sender per queue/topic -- si el e thread-safe, il pastrezi
        _sender = client.CreateSender("order-events");
    }

    public async Task PublishOrderPlacedAsync(OrderPlacedEvent evt)
    {
        var message = new ServiceBusMessage(JsonConvert.SerializeObject(evt))
        {
            ContentType = "application/json",
            MessageId = evt.OrderId.ToString(),   // pentru deduplicare
            Subject = "OrderPlaced",              // tipul evenimentului
            CorrelationId = evt.CorrelationId     // pentru tracing end-to-end
        };

        // Proprietati custom -- pe ele se aplica filtrele de subscription
        message.ApplicationProperties["TotalAmount"] = evt.TotalAmount;
        message.ApplicationProperties["Region"] = evt.Region;

        await _sender.SendMessageAsync(message);
    }
}

Detalii care contează:

  • MessageId — setat determinist (ex. OrderId), permite deduplicarea automată dacă activezi duplicate detection pe queue/topic.
  • Subject — tipul evenimentului; consumatorii pot rula handler-e diferite pe baza lui.
  • ApplicationProperties — pe acestea se aplică filtrele SQL de subscription din partea 1 (TotalAmount > 1000).

Batching pentru volum

Trimiterea mesaj-cu-mesaj la volume mari e ineficientă. ServiceBusMessageBatch împachetează mai multe mesaje într-o singură operație, respectând limita de dimensiune:

public async Task PublishBatchAsync(IEnumerable<OrderPlacedEvent> events)
{
    using ServiceBusMessageBatch batch = await _sender.CreateMessageBatchAsync();

    foreach (var evt in events)
    {
        var message = new ServiceBusMessage(JsonConvert.SerializeObject(evt));

        if (!batch.TryAddMessage(message))
        {
            // Batch-ul e plin -- trimite si incepe unul nou
            await _sender.SendMessagesAsync(batch);
            batch.Dispose();
            batch = await _sender.CreateMessageBatchAsync();

            if (!batch.TryAddMessage(message))
                throw new InvalidOperationException("Mesaj prea mare pentru un batch gol");
        }
    }

    if (batch.Count > 0)
        await _sender.SendMessagesAsync(batch);
}

Mesaje programate

// Livreaza mesajul peste 30 de minute (ex. reminder, retry intarziat)
message.ScheduledEnqueueTime = DateTimeOffset.UtcNow.AddMinutes(30);
await _sender.SendMessageAsync(message);

Consumatorul: ServiceBusProcessor

Pentru consum, folosește ServiceBusProcessor — gestionează automat lock renewal, concurrency și reconectarea. Îl integrezi natural într-un BackgroundService:

public class OrderProcessorService : BackgroundService
{
    private readonly ServiceBusProcessor _processor;
    private readonly IServiceScopeFactory _scopeFactory;
    private readonly ILogger<OrderProcessorService> _logger;

    public OrderProcessorService(
        ServiceBusClient client,
        IServiceScopeFactory scopeFactory,
        ILogger<OrderProcessorService> logger)
    {
        _scopeFactory = scopeFactory;
        _logger = logger;

        _processor = client.CreateProcessor(
            topicName: "order-events",
            subscriptionName: "inventory",
            new ServiceBusProcessorOptions
            {
                MaxConcurrentCalls = 5,       // 5 mesaje procesate in paralel
                PrefetchCount = 20,            // pre-incarca 20 pentru throughput
                AutoCompleteMessages = false,  // completam explicit, dupa succes
                MaxAutoLockRenewalDuration = TimeSpan.FromMinutes(5)
            });
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _processor.ProcessMessageAsync += OnMessageAsync;
        _processor.ProcessErrorAsync += OnErrorAsync;

        await _processor.StartProcessingAsync(stoppingToken);

        // Graceful shutdown -- oprire curata la SIGTERM
        stoppingToken.Register(() =>
            _processor.StopProcessingAsync().GetAwaiter().GetResult());
    }

    private async Task OnMessageAsync(ProcessMessageEventArgs args)
    {
        // Scope nou per mesaj -- repository-urile scoped functioneaza corect
        using var scope = _scopeFactory.CreateScope();
        var handler = scope.ServiceProvider.GetRequiredService<IOrderHandler>();

        var evt = JsonConvert.DeserializeObject<OrderPlacedEvent>(
            args.Message.Body.ToString())!;

        await handler.HandleAsync(evt, args.CancellationToken);

        // Complete DOAR dupa succes -- altfel mesajul revine in coada
        await args.CompleteMessageAsync(args.Message);
    }

    private Task OnErrorAsync(ProcessErrorEventArgs args)
    {
        _logger.LogError(args.Exception,
            "Eroare Service Bus: {ErrorSource}, entitate {EntityPath}",
            args.ErrorSource, args.EntityPath);
        return Task.CompletedTask;
    }
}

Setările care contează

  • AutoCompleteMessages = false — cel mai important. Completezi mesajul explicit după procesarea reușită. Dacă handler-ul aruncă excepție, mesajul revine în coadă și e relivrat — baza retry-ului natural.
  • MaxConcurrentCalls — paralelismul per instanță. 5 înseamnă 5 mesaje simultan; corelează cu resursele containerului și cu capacitatea downstream-ului (nu bombarda baza de date).
  • PrefetchCount — câte mesaje pre-încarcă clientul. Crește throughput-ul, dar mesajele prefetch-uite au lock-ul pornit — un prefetch prea mare + procesare lentă = lock-uri expirate. Regulă practică: 2-4x MaxConcurrentCalls.
  • MaxAutoLockRenewalDuration — processor-ul reînnoiește automat lock-ul pentru procesări lungi. Setează peste durata maximă reală de procesare.

Legătura cu seria: acest BackgroundService respectă CancellationToken (graceful shutdown din primul articol al seriei), iar în Container Apps, KEDA scalează exact acest consumator pe lungimea cozii (partea 2 din mini-seria ACA).


Ce urmează

Ce se întâmplă când procesarea eșuează repetat? În partea a treia: delivery count, dead-letter queue, retry policies cu backoff exponențial și cum monitorizezi și reprocesezi mesajele moarte.

Întrebări? Scrie-mi la contact@ludoprogramming.com.