RO EN

Cosmos DB patterns (2) — Change Feed ca trigger pentru evenimente

Cosmos DB patterns (2) — Change Feed ca trigger pentru evenimente ✨ Imagine generată cu AI
Doru Bulubașa
28 iulie 2026
42 vizualizări

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.