Ultima parte din seria despre Azure Service Bus. Avem fundamentele, producători și consumatori și gestionarea eșecurilor. A rămas cea mai subtilă problemă: consistența dintre baza de date și mesajele publicate.
Problema dual-write
Codul din partea 1 arăta așa:
public async Task<IActionResult> PlaceOrder(OrderRequest request)
{
var order = await _orderRepository.CreateAsync(request); // scriere 1: DB
await _publisher.PublishOrderPlacedAsync(new OrderPlacedEvent(order)); // scriere 2: Service Bus
return Accepted(order);
}
Două scrieri în două sisteme diferite, fără tranzacție comună. Ce poate merge prost:
- DB reușește, publish eșuează — comanda există, dar niciun serviciu nu află de ea. Stocul nu se rezervă, email-ul nu pleacă. Comandă-fantomă.
- Publish reușește, dar procesul moare înainte de commit (în ordinea inversă) — serviciile reacționează la o comandă care nu există în DB.
Nu poți rezolva cu try/catch: între cele două scrieri procesul poate muri oricând (deploy, OOM, scale-down — exact evenimentele frecvente din Container Apps). Nici tranzacțiile distribuite (2PC) nu sunt o opțiune realistă: Cosmos DB și Service Bus nu participă într-o tranzacție comună.
Ideea Outbox: o singură scriere atomică
Soluția: nu scrii în două sisteme. Scrii o singură dată, în baza de date, atât starea cât și evenimentul de publicat — atomic. Un proces separat (dispatcher) citește evenimentele nepublicate și le trimite în Service Bus.
1. [Tranzactie DB] salveaza Order + OutboxEvent (atomic)
2. [Dispatcher] citeste OutboxEvent-uri nepublicate
3. [Dispatcher] publica in Service Bus
4. [Dispatcher] marcheaza evenimentul ca publicat
Dacă procesul moare între pașii 1 și 3, evenimentul rămâne în outbox și dispatcher-ul îl publică la următoarea rulare. Garanția: dacă starea s-a salvat, evenimentul va fi publicat (eventual). Semantica e at-least-once — revenim la asta.
Implementare cu Cosmos DB
În Cosmos DB, atomicitatea există la nivel de partition key: un transactional batch scrie mai multe documente atomic dacă au aceeași cheie de partiție. Stocăm evenimentul outbox în același container cu comanda, pe aceeași partiție:
public class OutboxEvent
{
public string Id { get; set; } = Guid.NewGuid().ToString();
public string Type { get; set; } = "OutboxEvent"; // discriminator
public string OrderId { get; set; } = default!; // partition key comun
public string EventType { get; set; } = default!; // "OrderPlaced"
public string Payload { get; set; } = default!; // JSON-ul evenimentului
public bool Published { get; set; }
public DateTime CreatedAt { get; set; } = DateTime.UtcNow;
public DateTime? PublishedAt { get; set; }
}
public async Task CreateOrderAsync(Order order)
{
var outboxEvent = new OutboxEvent
{
OrderId = order.Id,
EventType = "OrderPlaced",
Payload = JsonConvert.SerializeObject(new OrderPlacedEvent(order))
};
// Scriere ATOMICA: ambele documente sau niciunul
var batch = _container.CreateTransactionalBatch(new PartitionKey(order.Id))
.CreateItem(order)
.CreateItem(outboxEvent);
var response = await batch.ExecuteAsync();
if (!response.IsSuccessStatusCode)
throw new InvalidOperationException(
$"Salvarea comenzii a esuat: {response.StatusCode}");
}
Aici e toată magia: CreateTransactionalBatch garantează că Order și OutboxEvent se salvează împreună sau deloc. Problema dual-write a dispărut — a rămas o singură scriere.
Dispatcher-ul
Un BackgroundService care rulează periodic, citește evenimentele nepublicate și le trimite în Service Bus:
public class OutboxDispatcherService : BackgroundService
{
private readonly Container _container;
private readonly ServiceBusSender _sender;
private readonly ILogger<OutboxDispatcherService> _logger;
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
try
{
await DispatchPendingAsync(stoppingToken);
}
catch (Exception ex)
{
_logger.LogError(ex, "Eroare la dispatch outbox");
}
await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
}
}
private async Task DispatchPendingAsync(CancellationToken ct)
{
// Cross-partition query -- sortam in memorie (lectie Cosmos DB!)
var query = new QueryDefinition(
"SELECT * FROM c WHERE c.Type = 'OutboxEvent' AND c.Published = false");
var pending = new List<OutboxEvent>();
using var iterator = _container.GetItemQueryIterator<OutboxEvent>(query);
while (iterator.HasMoreResults)
pending.AddRange(await iterator.ReadNextAsync(ct));
foreach (var evt in pending.OrderBy(e => e.CreatedAt))
{
var message = new ServiceBusMessage(evt.Payload)
{
ContentType = "application/json",
MessageId = evt.Id, // deduplicare pe id-ul evenimentului
Subject = evt.EventType,
CorrelationId = evt.OrderId
};
await _sender.SendMessageAsync(message, ct);
// Patch cu casing-ul EXACT al documentului stocat
await _container.PatchItemAsync<OutboxEvent>(
evt.Id,
new PartitionKey(evt.OrderId),
new[]
{
PatchOperation.Set("/Published", true),
PatchOperation.Set("/PublishedAt", DateTime.UtcNow)
},
cancellationToken: ct);
}
}
}
Două detalii din lecțiile Cosmos DB ale seriei: query-ul cross-partition nu folosește ORDER BY (sortăm în memorie), iar operațiile patch folosesc exact casing-ul câmpurilor stocate (/Published, nu /published) — altfel primești 500.
At-least-once și consumatori idempotenți
Priviți atent dispatcher-ul: dacă procesul moare după SendMessageAsync dar înainte de patch, evenimentul rămâne nepublicat în outbox și va fi trimis din nou. Semantica e at-least-once: garantăm publicarea, dar posibil de mai multe ori.
Două apărări, în straturi:
1. Duplicate detection în Service Bus
# Fereastra de deduplicare pe MessageId (max 7 zile pe Standard)
az servicebus queue create \
--name order-processing \
--namespace-name my-servicebus \
--resource-group my-rg \
--enable-duplicate-detection true \
--duplicate-detection-history-time-window PT30M
Pentru că MessageId e id-ul evenimentului outbox (determinist), Service Bus aruncă automat dublurile trimise în fereastra configurată.
2. Idempotență în consumator
Duplicate detection nu acoperă totul (fereastră limitată, relivrări după abandon). Ultima linie de apărare e consumatorul idempotent — procesarea aceluiași eveniment de două ori are efectul unei singure procesări:
public async Task HandleAsync(OrderPlacedEvent evt, CancellationToken ct)
{
// Verificare: am procesat deja acest eveniment?
var alreadyProcessed = await _processedEventsRepository
.ExistsAsync(evt.EventId, ct);
if (alreadyProcessed)
{
_logger.LogInformation("Eveniment {EventId} deja procesat -- skip", evt.EventId);
return;
}
await _inventoryService.ReserveStockAsync(evt.OrderId, evt.Items, ct);
// Marcheaza ca procesat -- ideal in aceeasi tranzactie cu efectul
await _processedEventsRepository.MarkProcessedAsync(evt.EventId, ct);
}
Ideal, verificarea și marcarea stau în aceeași tranzacție cu efectul de business (același transactional batch în Cosmos DB) — altfel ai recreat problema dual-write la nivel de consumator.
Checklist Outbox + încheierea seriei
- Stare + eveniment salvate atomic — transactional batch pe aceeași partition key
- Dispatcher separat de request — BackgroundService cu polling, tolerant la eșecuri
- MessageId determinist — id-ul evenimentului outbox, pentru deduplicare
- Duplicate detection activat — prima plasă pentru at-least-once
- Consumatori idempotenți — a doua plasă, obligatorie
- Evenimente vechi curățate — TTL pe documentele publicate (Cosmos DB face curățenia gratuit)
Cu asta, seria Service Bus e completă: de la decizia sincron/asincron, prin implementarea producătorilor și consumatorilor, prin gestionarea eșecurilor cu DLQ și retry, până la consistența garantată cu Outbox. Împreună cu mini-seria Container Apps (unde KEDA scalează exact acești consumatori), ai o arhitectură asincronă completă și rezilientă.
Dacă ai întrebări sau vrei să discuți cum aplici Outbox pattern în proiectul tău, scrie-mi la contact@ludoprogramming.com.