Modello di convoglio sequenziale

Raggruppare i messaggi correlati tramite una chiave di categoria ed elaborare ogni gruppo in sequenza, un messaggio alla volta, durante l'elaborazione di gruppi diversi in parallelo.

Questo modello risolve la tensione tra il mantenimento della correttezza FIFO (First-In, First-Out) all'interno di ogni gruppo logico e la scalabilità orizzontale dell'elaborazione simultanea tra gruppi. La progettazione garantisce che i vincoli di ordinamento non diventino un collo di bottiglia a livello di sistema.

Contesto e problema

Le applicazioni spesso devono elaborare i messaggi correlati nell'ordine in cui arrivano, continuando al contempo a scalare orizzontalmente per gestire un carico maggiore. In un'architettura distribuita, questo requisito è difficile da soddisfare perché i processi worker prelevano in modo indipendente i messaggi da una coda condivisa. Quando più lavoratori competono per i messaggi, come nel modello Consumer concorrenti, l'ordinamento si suddivide.

Si consideri un sistema di rilevamento degli ordini che riceve un flusso di operazioni, ad esempio la creazione di un ordine, l'aggiunta di una transazione, la modifica di una transazione precedente e l'eliminazione di un ordine. Le operazioni di ogni ordine devono essere elaborate nell'ordine FIFO, perché l'applicazione di tali operazioni fuori sequenza potrebbe danneggiare lo stato dell'ordine. Tuttavia, la coda in ingresso intercala le operazioni di molti ordini. Un singolo consumer che impone un ordinamento globale diventa un collo di bottiglia, e consumer multipli possono elaborare fuori sequenza le operazioni relative allo stesso ordine.

Gli approcci più semplici a questo problema falliscono ciascuno in un modo diverso:

  • Singolo consumatore. Un singolo consumer mantiene l'ordine dei messaggi perché elabora un messaggio alla volta, ma non può essere ridimensionato per gestire una maggiore velocità effettiva.

  • Più consumer concorrenti. Più consumer aumentano il throughput prelevando i messaggi in parallelo, ma perdono le garanzie di ordinamento all’interno di ciascun gruppo. Due worker possono prelevare messaggi consecutivi per lo stesso ordine ed elaborarli simultaneamente o fuori sequenza, corrompendo lo stato dell'ordine.

Soluzione

Il modello Convoglio sequenziale suddivide i messaggi correlati in categorie ed elabora ciascuna categoria in sequenza, un messaggio alla volta, mentre le categorie vengono elaborate in parallelo.

Il modello funziona assegnando a ogni messaggio una chiave di categoria che identifica il gruppo a cui appartiene. Un broker di messaggi usa questa chiave per partizionare i messaggi in gruppi logici. All'interno di ogni gruppo, il broker applica l'ordinamento FIFO in modo che un consumer che blocca un gruppo riceva i messaggi rigorosamente nella sequenza in cui sono stati accodati. I diversi gruppi possono essere elaborati contemporaneamente da diversi consumer, quindi il sistema viene ridimensionato orizzontalmente tra i gruppi senza sacrificare l'ordinamento all'interno di un singolo gruppo.

In Azure, bus di servizio di Azure sessioni di messaggio forniscono un'implementazione predefinita di questo modello.

Il diagramma seguente illustra il modello di convoglio sequenziale generale.

Diagramma dello schema di convoglio sequenziale. Mostra un produttore, una coda centrale e tre consumatori.

Nella coda, i messaggi per categorie diverse potrebbero essere intervallati, come illustrato nel diagramma seguente.

Diagramma che illustra quattro categorie di messaggi intercalati in una singola coda. Ogni categoria occupa una propria corsia orizzontale.

Questo modello offre diversi vantaggi principali:

  • Elaborazione ordinata per gruppo. I messaggi all'interno di ogni categoria vengono elaborati rigorosamente in sequenza, il che evita condizioni di competizione, modifiche dello stato fuori ordine e la necessità di ricorrere a espedienti di riordinamento.

  • Scalabilità orizzontale tra gruppi. Ogni categoria è un'unità indipendente di concorrenza. L'aggiunta di consumatori aumenta il throughput in modo proporzionale al numero di categorie attive, senza compromettere le garanzie di ordinamento.

  • Disaccoppiamento produttore-consumatore. I produttori accoderanno i messaggi senza conoscere quale consumer li elabora o quando. I consumer sono scalabili e sostituibili in modo indipendente.

Problemi e considerazioni

Quando si decide come implementare questo modello, tenere presente quanto segue:

  • Categoria e unità di scala. Determinare quale proprietà dei messaggi in arrivo usare per aumentare il numero di istanze. La chiave di categoria definisce l'unità di parallelismo: ogni valore di chiave distinta diventa un gruppo elaborabile in modo indipendente. Nello scenario di rilevamento degli ordini, questa proprietà è l'ID dell'ordine. La scelta di una chiave troppo grossolana (ad esempio, un singolo ID cliente per tutti gli ordini) limita il parallelismo, mentre la scelta di una chiave troppo fine non offre vantaggi significativi per l'ordinamento.

  • Limiti di portata. Valuta la velocità di trasmissione dei messaggi prevista. Poiché questo modello applica l'elaborazione sequenziale all'interno di ogni categoria, la velocità effettiva per categoria è limitata dal tempo per elaborare un singolo messaggio. Ottimizzare il tempo di elaborazione per messaggio, ad esempio, usando operazioni di I/O asincrone o invio in batch di scritture downstream, perché tale tempo determina direttamente la velocità effettiva massima per ogni categoria. Se il requisito complessivo di capacità effettiva è molto elevato, riconsiderate se sia necessario un ordinamento FIFO rigoroso durante l'intero ciclo di vita dei messaggi. Le alternative includono l'applicazione di un messaggio di inizio e di fine a una sequenza o l'ordinamento dei messaggi in base al timestamp all'interno di una finestra batch e l'invio del batch per l'elaborazione parallela.

  • Funzionalità del servizio. Verificare se il message broker scelto supporta l'elaborazione di un messaggio alla volta all'interno di una coda o di una sua categoria. Non tutti i servizi di messaggistica forniscono garanzie di blocco a livello di sessione o FIFO all'interno di una partizione. Se il broker non supporta in modo nativo questa funzionalità, il consumer deve implementare la propria logica di coordinamento, che aggiunge complessità e rischi di elaborazione duplicata, messaggi persi o esecuzione non ordinata. Il supporto delle sessioni potrebbe anche limitare la scelta del livello di messaggistica o dello SKU, che influisce sui costi.

  • Evolvibilità. Pianificare la modalità di aggiunta di nuove categorie di messaggi al sistema. Il modello deve supportare la crescita della cardinalità delle categorie senza richiedere modifiche strutturali ai consumatori. Si supponga, ad esempio, che il sistema libro mastro descritto in precedenza sia specifico di un cliente. Se devi acquisire un nuovo cliente, dovresti poter aggiungere un insieme di processori del ledger che distribuiscono il lavoro in base all'ID cliente senza riprogettare la topologia della coda.

  • Recapito non ordinato dei messaggi. I messaggi possono arrivare fuori ordine a causa della latenza di rete variabile tra il produttore e il broker, prima che l'ordinamento della sessione del broker entri in vigore. È consigliabile usare i numeri di sequenza per verificare l'ordinamento all'interno di ogni categoria. È anche possibile includere un flag di fine sequenza nell'ultimo messaggio di una transazione in modo che i consumer possano rilevare quando una sequenza è stata completata.

  • Gestione dei messaggi poison. Un messaggio che ha ripetutamente esito negativo durante l'elaborazione all'interno di una sessione blocca tutti i messaggi successivi in tale sessione perché il modello applica un ordinamento sequenziale rigoroso. Progettare una strategia per rilevare i messaggi non elaborabili, ad esempio tenere traccia dei conteggi dei tentativi di recapito e spostarli in una coda di messaggi non recapitabili dopo una soglia di ripetizione dei tentativi definita in modo che i messaggi rimanenti nella sessione possano continuare l'elaborazione.

  • Disponibilità del broker. Il broker di messaggi è una dipendenza condivisa per tutte le categorie. La disponibilità e la durabilità influiscono direttamente sulle garanzie di affidabilità del modello. Valutare le funzionalità di resilienza a livello di broker, ad esempio le zone di disponibilità e il ripristino di emergenza geografico in base ai requisiti di disponibilità e al budget del carico di lavoro, perché le configurazioni a durabilità più elevata aumentano in genere i costi.

  • Correttezza della chiave producer. Il modello presuppone che i producer impostino correttamente la chiave di categoria (ID sessione) in ogni messaggio. Se un producer imposta una chiave non corretta, accidentalmente o a causa di un bug, il messaggio instrada alla sessione errata e danneggia lo stato del gruppo. Verificare che i produttori assegnino le chiavi di categoria in modo coerente e valutare l'aggiunta di una logica di validazione della chiave nel consumer se le conseguenze di un messaggio instradato erroneamente sono gravi.

  • Complessità operativa. Il monitoraggio dell'elaborazione basata su sessione comporta un sovraccarico operativo superiore a quello dell'utilizzo standard della coda. Gli operatori necessitano di visibilità sui backlog delle sessioni (il numero di sessioni attive e la profondità dei messaggi in attesa in ogni sessione) per identificare le categorie che si trovano in ritardo. Le sessioni non recapitabili richiedono un flusso di lavoro di monitoraggio e correzione separato per analizzare i messaggi non riusciti, risolvere la causa radice e riprodurre i messaggi corretti nella sessione.

  • Contenzione e latenza del blocco sessione. Il blocco della sessione introduce un sovraccarico di latenza perché ogni consumer deve ottenere un blocco esclusivo su una sessione prima di elaborare i messaggi. Quando un consumer detiene un blocco di sessione, nessun altro consumer può elaborare i messaggi di tale sessione, anche se il consumer è lento o temporaneamente in stallo. Se la durata del blocco è troppo breve, la scadenza del blocco può causare la rielaborazione del messaggio. Se la durata del blocco è troppo lunga, un consumatore in stallo ritarda il ripristino. Ottimizzare la durata del blocco della sessione in base al tempo di elaborazione dei messaggi previsto e implementare il rinnovo del blocco per le operazioni con esecuzione più lunga.

  • Scalabilità e costi per il mercato consumer. Il parallelismo tra sessioni si traduce in istanze consumer concorrenti. In un modello serverless come Funzioni di Azure, ogni sessione attiva corrisponde a un'esecuzione concorrente, mentre in un modello dedicato corrisponde a un'istanza o a un thread. Il numero di sessioni attive influisce quindi direttamente sul costo di calcolo. Pianificate i limiti di scalabilità dei consumer e i controlli della concorrenza per bilanciare la capacità effettiva e i costi.

Quando usare questo modello

Usare questo modello quando:

  • I messaggi arrivano in ordine e devono essere elaborati nello stesso ordine.
  • I messaggi possono essere classificati in modo che ogni categoria diventi un'unità di scala indipendente per il sistema.

Questo modello potrebbe non essere adatto quando:

  • Si prevedono scenari con velocità effettiva estremamente elevata (milioni di messaggi al minuto), perché il requisito FIFO limita il ridimensionamento che il sistema può raggiungere.

  • L'ordinamento dei messaggi non è obbligatorio. Quando i messaggi possono essere elaborati in modo indipendente in un ordine qualsiasi, il pattern Competing Consumers offre una scalabilità orizzontale più semplice senza il sovraccarico di coordinamento dovuto al blocco della sessione.

Progettazione del carico di lavoro

Valutare come usare il convoglio sequenziale nella progettazione di un carico di lavoro per soddisfare gli obiettivi e i principi trattati nei pilastri di Azure Well-Architected Framework. La tabella seguente fornisce indicazioni su come questo modello supporta gli obiettivi di ogni pilastro.

Pilastro Come questo modello supporta gli obiettivi di pilastro
decisioni di progettazione dell'affidabilità consentono al carico di lavoro di diventare resiliente a un malfunzionamento e assicurano che ripristini a uno stato completamente funzionante dopo che si verifica un guasto. Questo modello usa l'ordinamento FIFO basato sulle sessioni per eliminare le race condition, la logica di gestione dei messaggi soggetta a contese e altre soluzioni temporanee adottate per gestire messaggi ordinati in modo errato che possono causare malfunzionamenti.

- Flussi critici RE:02
- RE:07 Processi in background

Se questo modello introduce compromessi all'interno di un pilastro, considerarli contro gli obiettivi degli altri pilastri.

Esempio

In Azure è possibile implementare questo modello usando bus di servizio sessioni di messaggio. Per i consumer, è possibile utilizzare App per la logica di Azure con il connettore bus di servizio peek-lock oppure Funzioni di Azure con il trigger di bus di servizio.

Quando un producer imposta la SessionId proprietà su un messaggio, bus di servizio raggruppa tutti i messaggi che condividono lo stesso ID sessione in una singola sessione logica. Un utente accetta una sessione e ottiene un blocco esclusivo su di essa. Questo blocco garantisce che un solo consumer elabori i messaggi di tale sessione in un dato momento e che i messaggi arrivino in ordine FIFO. Altri consumer possono accettare ed elaborare contemporaneamente sessioni diverse, offrendo velocità effettiva parallela tra gruppi.

Nell'esempio di monitoraggio degli ordini, il sistema elabora ogni messaggio del registro nell'ordine in cui viene ricevuto e invia ogni transazione a un'altra coda, in cui la categoria è impostata sull'ID dell'ordine. Una transazione non coinvolge mai più ordini in questo scenario, quindi i consumer elaborano ciascuna categoria in parallelo, ma seguendo l'ordine FIFO all'interno della categoria.

Il processore del registro distribuisce i messaggi scomponendo in singoli elementi il contenuto di ciascun messaggio nella prima coda:

Diagramma dell'architettura di esempio Sequential Convoy. Mostra un produttore, una coda del libro mastro, un processore del libro mastro, una coda delle transazioni e tre processori degli ordini.

Il processore libro mastro esegue tre passaggi:

  1. Scorre il libro mastro una transazione alla volta.
  2. Imposta l'ID sessione del messaggio in modo che corrisponda all'ID dell'ordine.
  3. Invia ogni transazione del libro mastro a una coda secondaria con l'ID della sessione impostato sull'ID dell'ordine.

I consumatori ascoltano la coda secondaria ed elaborano tutti i messaggi con ID dell’ordine corrispondenti in ordine FIFO. I consumatori usano la modalità peek-lock.

La coda del libro mastro è un punto di transizione da seriale a parallelo: tutte le transazioni vi passano in modo sequenziale prima di essere smistate verso un'elaborazione parallela basata sulle sessioni. Questa fase di serializzazione è il principale collo di bottiglia per la scalabilità perché limita il throughput dell'intera pipeline a valle. Tuttavia, dopo che il processore del registro distribuisce i messaggi alla coda secondaria, i consumer possono essere scalati in modo indipendente tra le sessioni, uno per ciascun ID ordine.

Tecnologie di supporto

  • Sessioni di messaggi di bus di servizio: raggruppa i messaggi in base all'ID di sessione e garantisce l'elaborazione FIFO all'interno di ogni sessione. Le sessioni di messaggio sono il meccanismo principale di Azure per l'implementazione del modello di convoglio sequenziale.

  • trigger Funzioni di Azure bus di servizio: supporta trigger basati su sessione che consentono alle istanze della funzione di elaborare i messaggi da una singola sessione alla volta.

  • Connettore bus di servizio di Logic Apps: fornisce un connettore bus di servizio con supporto peek-lock per l'utilizzo di code con supporto per le sessioni nell'elaborazione basata su workflow.

Contributori

Microsoft gestisce questo articolo. I seguenti collaboratori hanno scritto questo articolo.

Autore principale:

  • Taka Venkata Cheruvu | Senior Cloud Solution Architect e infrastruttura di intelligenza artificiale

Per visualizzare i profili LinkedIn non pubblici, accedere a LinkedIn.

  • Modello consumer concorrenti: più consumer eseguono il pull di messaggi da una coda condivisa in parallelo, aumentando la velocità effettiva, ma rimuovendo le garanzie di ordinamento per messaggio. Il pattern Sequential Convoy affronta il problema di ordinamento introdotto da Competing Consumers. Risolve questo divario partizionando i messaggi in sessioni con chiave di categoria ed elaborando ogni sessione in sequenza.

  • Queue-Based modello di livellamento del carico: un buffer di accodamento funziona tra produttori e consumer per assorbire picchi e carico irregolare senza problemi. Il pattern Convoy sequenziale si basa su questo buffering aggiungendo il partizionamento basato su sessioni, in modo che la coda bilanci il carico tra le categorie e preservi l'ordinamento FIFO all'interno di ogni categoria.

  • Modello di coda con priorità: i messaggi vengono instradati verso code separate oppure viene assegnata loro una priorità all'interno di una coda, in modo che le attività con priorità più alta vengano elaborate prima di quelle con priorità più bassa. Quando è necessario preservare anche l'ordinamento all'interno di un livello di priorità, il pattern del convoglio sequenziale può essere combinato con una coda con priorità per garantire l'elaborazione FIFO all'interno di ogni sessione identificata dalla priorità.

  • Messaggio Peek-Lock (lettura non distruttiva): questa operazione recupera e blocca atomicamente un messaggio da una coda o da una sottoscrizione per l'elaborazione.

  • Recapito in ordine dei messaggi correlati in Logic Apps usando le sessioni di bus di servizio: questo post del blog descrive il supporto di Logic Apps per il pattern Sequential Convoy.