Go to main content
Kafka Producer and Consumer - official lesson image on GinnyTech

Producer, Consumer, and Serialization

Implement robust Kafka producers and consumers with optimal serialization patterns for analytics.

AD
Created byAndrii Dyshkantiuk
Lesson 113 / 236Level: AdvancedDuration: 22 minPrerequisites: 1

What you will learn

  • Configurare acks, idempotenza e commit degli offset in base alla garanzia di delivery promessa
  • Scegliere la chiave di partizionamento sulla distribuzione reale del traffico
  • Versionare lo schema dei messaggi per garantire compatibilità tra team

Producer, Consumer, and Serialization

Questa lezione appartiene al binario ml-tabellare e ti porta dal lato pratico del modulo: producer e consumer non sono client da configurare, ma contratti da progettare, e ogni parametro dichiara una garanzia esplicita su perdita, ordine, duplicati e formato.

L’idea in una frase

Producer e consumer robusti sono contratti eseguibili che dichiarano garanzie esplicite su perdita accettabile, ordine, duplicati e formato dei messaggi.

La procedura in sei passi

  1. Dichiara la garanzia di delivery: perdita accettabile, ordine necessario, duplicati tollerati.
  2. Imposta acks di conseguenza: all per dati critici, 1 per analytics tolleranti.
  3. Attiva l’idempotenza del producer e rendi il consumer tollerante ai duplicati.
  4. Scegli la chiave guardando la distribuzione reale del traffico, non le ipotesi.
  5. Fissa il formato come interfaccia pubblica con schema versionato e check di compatibilità.
  6. Configura il commit degli offset dopo il processamento, con gestione esplicita dei crash.

Cosa promette davvero un producer

Un producer sembra semplice finché un retry duplica eventi, un consumer resta indietro o una serializzazione rompe un servizio a valle. Il lavoro professionale è progettare chiavi, acks, batch, commit degli offset e formato del messaggio come parti dello stesso contratto. Producer e consumer robusti non sono client configurati bene: sono contratti eseguibili tra sistemi. Ogni scelta risponde a una domanda: quale garanzia prometti su perdita accettabile, duplicati tollerati, ordine necessario, compatibilità del formato e recupero dopo un errore.

Inviare un messaggio sembra banale: istanzi un client e invochi un metodo di produzione. La differenza tra implementazione ingenua e professionale si misura in durabilità, throughput e latenza. Un producer non spedisce ogni messaggio subito sulla rete. Usa un accumulatore di record, un buffer dove i messaggi sono raggruppati in batch per partizione, e un thread di invio in background che spedisce i batch ai broker. La competenza sta nel configurare questo meccanismo sul caso d’uso.

La configurazione che decide la durabilità

The most critical configuration is acks, che definisce la garanzia sulla scrittura. Qui dichiari in modo esplicito quanto rischi in cambio di latenza più bassa.

Valore di acksGaranziaWhen to Use
0Nessuna conferma, throughput massimoTelemetria sacrificabile, log non critici
1Conferma dal solo leaderAnalytics dove un evento perso ogni tanto non sposta le metriche
allConferma da tutte le ISR, nessuna perditaTransazioni e dati che non possono andare persi

La scelta non è solo tecnica. È una decisione di prodotto travestita da parametro: stai fissando una soglia di rischio accettabile.

Chiavi, ordine e partizionamento

La chiave decide la partizione, e la partizione decide l’ordine. Kafka garantisce l’ordine solo dentro la singola partizione, quindi se gli eventi di un ordine devono arrivare in sequenza devi usare order_id come chiave. Il rovescio è che una chiave a bassa cardinalità, o una chiave calda, concentra il traffico su poche partizioni e sbilancia i consumer. La scelta della chiave è quindi un compromesso tra ordine e parallelismo, e va fatta sui pattern di accesso reali, non su un’ipotesi astratta.

Consumer e gestione degli offset

Dal lato consumer, il punto delicato è il commit degli offset, cioè il modo in cui Kafka registra cosa hai già letto. Se committi prima di processare, un crash fa perdere il messaggio, con delivery al massimo una volta. Se committi dopo, un crash a metà lavoro fa rileggere lo stesso messaggio, con delivery almeno una volta, e l’applicazione deve essere idempotente per non duplicare gli effetti. Anche la scelta tra commit automatico e manuale è una dichiarazione di garanzia, non una preferenza stilistica.

Esempio: progettare il contratto degli eventi

Il contratto si progetta su un flusso concreto, per esempio transazioni: la domanda non è qual è la configurazione corretta in assoluto, ma quale scelta diventa meno rischiosa se progettata bene. La tabella mostra come leggere le situazioni tipiche.

SituationCautious interpretationDecision
Un retry sembra duplicare eventiIl producer non è idempotenteAttivare l’idempotenza e rendere il consumer tollerante ai duplicati
Una partizione riceve molto più trafficoLa chiave scelta crea hotspotRivedere la chiave o aumentare la cardinalità
Un servizio a valle si rompe dopo un deployIl formato del messaggio è cambiato in modo incompatibileVersionare lo schema e garantire la compatibilità
Il consumer accumula lagCapacità di consumo insufficienteAggiungere istanze fino al numero di partizioni

Serializzazione e compatibilità del formato

Il formato del messaggio è parte del contratto quanto chiave e acks. JSON è leggibile e flessibile ma verboso e senza garanzie di schema, mentre Avro o Protobuf con schema registry rendono esplicita la struttura e permettono di evolvere il formato senza rompere i consumer. La regola è semplice: appena più di un team legge un topic, il formato smette di essere un dettaglio implementativo e diventa un’interfaccia pubblica, con tutte le responsabilità di compatibilità che ne derivano.

Errori tipici da evitare

L’errore più frequente è trattare producer e consumer come configurazioni da copiare invece che come contratti da progettare. Succede quando scegli acks senza chiederti quanta perdita è accettabile, o una chiave senza guardare la distribuzione reale del traffico. Il sintomo è costante: tutto funziona in sviluppo e si rompe sotto carico.

Tre controlli minimi riducono il rischio. Primo, dichiarare la garanzia di delivery promessa. Secondo, verificare la distribuzione delle chiavi su un campione di traffico reale prima della produzione. Terzo, versionare lo schema dei messaggi e testare la compatibilità prima di ogni cambio di formato.

L’esempio che fa da riferimento

Kafka nasce in LinkedIn per collegare decine di sistemi produttori e consumatori senza una rete ingestibile di pipeline punto a punto, e viene pubblicato open source nel 2011. Il disegno di chiavi che decidono la partizione, conferme di scrittura configurabili e offset gestiti dai consumer viene da quell’esperienza operativa. Nel 2014 gli stessi ingegneri fondano Confluent e portano lo stesso modello sul mercato, con producer idempotenti e serializzazione governata. Da quel percorso discende la regola della lezione: ogni parametro del client è una garanzia dichiarata, non un dettaglio.

Verdetto: ogni parametro del client è una garanzia dichiarata, non un dettaglio: acks=all per i dati critici e 1 per l’analytics tollerante, chiave scelta sulla distribuzione reale del traffico, e schema versionato appena più di un team legge il topic.

Domande per verificare la lezione

  1. Quale livello di acks hai scelto e quanta perdita di dati dichiari accettabile?
  2. Quale chiave di partizionamento usi e come hai verificato la distribuzione del traffico?
  3. Quando committi gli offset e come gestisci i duplicati dopo un crash?
  4. Quale formato di serializzazione usi e come garantisci la compatibilità tra team?
Serve una mano concreta?

Bloccato su questo argomento o vuoi applicarlo al tuo caso? Prenota una call di 15 minuti con un analista esperto.

Book a call