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.