A doua parte din seria despre pattern-uri avansate în Cosmos DB. Cu partition key-ul stabilit, trecem la una dintre cele mai puternice funcționalități: Change Feed.
Ce este Change Feed
Change Feed e fluxul ordonat al modificărilor dintr-un container: fiecare insert și fiecare update apare în feed, în ordinea în care s-a întâmplat (per partition key). Practic, containerul tău devine gratuit o sursă de evenimente — fără să scrii cod de publicare, fără dual-write.
Ce NU conține feed-ul: ștergerile (folosește soft delete cu TTL dacă ai nevoie să reacționezi la ele) și versiunile intermediare (dacă un document e modificat de 3 ori rapid, poți vedea doar starea finală).
De ce contează pentru arhitectură: multe scenarii care altfel cer mesaje explicite se rezolvă elegant citind feed-ul:
- Indexare — documentul salvat trebuie indexat pentru căutare (embeddings, full-text). Feed-ul e trigger-ul natural.
- Invalidare de cache — configurația unui tenant s-a schimbat, cache-ul trebuie invalidat pe toate instanțele.
- Proiecții / materialized views — menții un container de citire denormalizat, actualizat din feed-ul containerului de scriere.
- Arhivare și analytics — copiezi modificările spre storage ieftin sau spre un sistem analitic.
Change Feed Processor în .NET
SDK-ul oferă Change Feed Processor — gestionează automat checkpointing-ul (unde ai rămas), distribuirea partițiilor între instanțe și reluarea după restart. Are nevoie de un lease container: un container mic în care își ține starea.
public class MessageIndexerService : BackgroundService
{
private readonly CosmosClient _client;
private readonly ILogger<MessageIndexerService> _logger;
private ChangeFeedProcessor? _processor;
public MessageIndexerService(CosmosClient client, ILogger<MessageIndexerService> logger)
{
_client = client;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
var database = _client.GetDatabase("ChatDb");
var monitored = database.GetContainer("messages"); // containerul urmarit
var leases = database.GetContainer("leases"); // starea processor-ului
_processor = monitored
.GetChangeFeedProcessorBuilder<ChatMessage>(
processorName: "message-indexer", // identitate stabila!
onChangesDelegate: HandleChangesAsync)
.WithInstanceName(Environment.MachineName) // unic per instanta
.WithLeaseContainer(leases)
.WithStartTime(DateTime.MinValue.ToUniversalTime()) // de la inceput
.Build();
await _processor.StartAsync();
stoppingToken.Register(() =>
_processor.StopAsync().GetAwaiter().GetResult());
}
private async Task HandleChangesAsync(
ChangeFeedProcessorContext context,
IReadOnlyCollection<ChatMessage> changes,
CancellationToken ct)
{
foreach (var message in changes)
{
await _searchIndexer.IndexAsync(message, ct);
}
_logger.LogInformation(
"Procesate {Count} modificari, RU: {RU}",
changes.Count, context.Headers.RequestCharge);
}
}
Detaliile care contează:
- processorName stabil — identifică logic acest consumator; îl schimbi = feed-ul e reluat de la StartTime. Două procesoare cu nume diferite pe același container primesc fiecare toate modificările (ca două subscriptions).
- WithInstanceName unic — pe el se face distribuirea lease-urilor între instanțe.
- Lease container cu partition key /id — convenția standard; îl creezi o dată cu throughput mic (400 RU sau shared).
Scalarea pe mai multe instanțe
Rulezi 3 instanțe ale aceluiași BackgroundService (același processorName, instanceName diferit)? Processor-ul împarte automat partițiile între ele prin lease-uri. O instanță moare? Lease-urile ei expiră și sunt preluate de celelalte. Zero cod suplimentar.
Corelația cu seria Container Apps: consumatorul de Change Feed e exact tipul de workload care rulează natural în ACA — dar atenție la scale-to-zero: cu 0 replici, nimeni nu citește feed-ul. Pentru procesare continuă, min-replicas 1; pentru procesare periodică tolerantă la întârziere, scale-to-zero e acceptabil (feed-ul te așteaptă).
Gestionarea erorilor: feed-ul nu are DLQ
Diferență crucială față de Service Bus: dacă delegate-ul aruncă excepție, batch-ul e relivrat integral, la nesfârșit. Un document otrăvit blochează partiția lui. Nu există dead-letter automat — ți-l construiești:
private async Task HandleChangesAsync(
ChangeFeedProcessorContext context,
IReadOnlyCollection<ChatMessage> changes,
CancellationToken ct)
{
foreach (var message in changes)
{
try
{
await _searchIndexer.IndexAsync(message, ct);
}
catch (Exception ex)
{
// NU re-arunca -- ar bloca partitia pe acest batch
_logger.LogError(ex, "Indexare esuata pentru {Id}", message.Id);
// Dead-letter manual: salveaza esecul pentru reprocesare
await _failedItemsContainer.CreateItemAsync(new FailedChange
{
SourceId = message.Id,
PartitionKey = message.SessionId,
Error = ex.Message,
Payload = JsonConvert.SerializeObject(message),
FailedAt = DateTime.UtcNow
}, cancellationToken: ct);
}
}
}
Regula: excepțiile nu părăsesc delegate-ul. Eșecurile individuale merg într-un container de failed items (cu alertă pe el, ca DLQ-ul din Service Bus), iar batch-ul avansează.
Change Feed vs. mesaje explicite
| Criteriu | Change Feed | Service Bus + Outbox |
|---|---|---|
| Efort de publicare | Zero — automat la scriere | Outbox pattern necesar |
| Ce transportă | Starea documentului | Eveniment de business modelat explicit |
| Ștergeri | Nu (doar cu soft delete) | Da, ca eveniment |
| Consumatori externi sistemului | Nepotrivit (acces direct la DB) | Natural |
| Retry / DLQ | Manual | Built-in |
Regula practică: Change Feed pentru reacții interne la schimbarea datelor (indexare, cache, proiecții); Service Bus pentru evenimente de business consumate de alte servicii.
Ce urmează
În partea a treia: multi-region writes — când ai nevoie de scrieri în mai multe regiuni, ce se întâmplă la conflicte și cum configurezi conflict resolution policies.
Întrebări? Scrie-mi la contact@ludoprogramming.com.