
Kafka Streams: processare eventi con Java
Introduzione a Kafka Streams per trasformazioni stateful su flussi di eventi senza cluster esterno.
Cosa imparerai
- Disegnare una topology Kafka Streams con stato locale ricostruibile dai changelog
- Configurare aggregazioni con finestre e join con tabelle compatte
- Attivare la semantica exactly-once solo dove i duplicati hanno conseguenze reali
Collegamenti
Kafka Streams: processare eventi con Java
Questa lezione del binario ml-tabellare ti porta dalla trasmissione degli eventi alla loro elaborazione: con Kafka Streams la logica di trasformazione vive dentro la tua applicazione, e il prezzo da pagare è la gestione consapevole dello stato.
L’idea in una frase
Kafka Streams è una libreria che trasforma eventi con stato locale ricostruibile, senza cluster separato da mantenere.
La procedura in cinque passi
- Disegna la topology come flusso da topic sorgente a topic di destinazione.
- Dimensiona stato locale, retention dei changelog e gestione dei late events.
- Configura aggregazioni con finestre e join con tabelle compatte.
- Attiva la garanzia exactly-once solo dove i duplicati hanno conseguenze reali.
- Pianifica rebalance e restart verificando la ricostruzione dello stato dai changelog.
Perché lo stato cambia tutto
Prima o poi un’applicazione deve arricchire eventi, calcolare aggregati e reagire a sequenze di comportamento senza uscire da Kafka. Kafka Streams porta questa logica dentro un’applicazione deployabile, ma con un costo: introduce stato locale, changelog, repartition e failure mode propri. Leggilo quindi non come libreria di trasformazioni, ma come progetto di un servizio stateful. La domanda giusta non è quale metrica calcoli, ma quale decisione cambia se l’applicazione resta corretta mentre scala o riparte.
Una topology si valuta seguendo lo stato: dove nasce, dove viene salvato, come si ricostruisce e cosa succede durante un rebalance. È corretta quando conserva semantica, ordine e recuperabilità mentre cambia numero di istanze o riparte dopo un crash. Tre vincoli decidono il disegno: quanto stato locale serve, quanta retention tenere sui changelog e come gestire i late events. Se ignori uno di questi tre punti, l’applicazione funziona in demo e si rompe in produzione.
Il modello di Kafka Streams
Uno stream processing in Kafka Streams segue sempre questo pattern:
Source Topic → Stream Processor → Sink Topic
Per esempio, per filtrare le transazioni sospette bastano poche righe:
KStream<String, Transaction> transactions = builder.stream("transactions");
KStream<String, Transaction> suspicious = transactions
.filter((key, txn) -> txn.getAmount() > 10000 && txn.getCountry().equals("NG"));
suspicious.to("suspicious_transactions");
Leggi, filtri, scrivi. Sotto questa semplicità, Kafka Streams gestisce stato, ribilanciamento tra istanze e tolleranza ai guasti. È questa automazione che rende il sistema potente e, insieme, difficile da debuggare quando qualcosa va storto.
Operazioni stateful: aggregazioni e finestre
Le operazioni interessanti sono quelle che mantengono stato. Per contare gli eventi per tipo ogni 5 minuti:
KTable<Windowed<String>, Long> counts = transactions
.groupBy((key, txn) -> txn.getType())
.windowedBy(TimeWindows.of(Duration.ofMinutes(5)))
.count();
Per arricchire ogni evento con dati esterni si usa il join tra stream e tabella:
KStream<String, Transaction> enriched = transactions
.join(userTable, (txn, user) ->
new EnrichedTransaction(txn, user.getName()));
Il join unisce ogni transazione con i dati utente più recenti, dove la tabella è un log compatto di Kafka. È il pattern che arricchisce eventi in tempo reale senza interrogare un database esterno per ogni messaggio.
Exactly-once semantics
Kafka Streams implementa la semantica exactly-once nativamente dal 2017. Usa le transazioni Kafka per scrivere in modo atomico su più topic, così ogni evento è processato esattamente una volta anche dopo crash e riavvio. Si attiva con la garanzia di processamento exactly-once. La garanzia non è gratis (costa latenza e throughput), quindi va richiesta solo quando la duplicazione avrebbe conseguenze reali, come in un conteggio finanziario.
Kafka Streams o Flink
La scelta tra Kafka Streams e Flink dipende da quanto è pesante l’elaborazione e da quanta infrastruttura vuoi gestire.
| Kafka Streams | Flink | |
|---|---|---|
| Deploy | Libreria embedded | Cluster separato |
| Complessità | Semplice, Java/Kafka nativo | Complesso, richiede infrastruttura |
| Use case ideale | Trasformazioni leggere, arricchimento | Aggregazioni pesanti, ML su stream, event time complesso |
| Throughput massimo | Milioni/msg sec | Decine di milioni/msg sec |
In pratica Kafka Streams vince quando vuoi trasformazioni leggere e arricchimento senza un cluster da mantenere, mentre Flink vale la complessità in più quando hai aggregazioni pesanti, machine learning sullo stream o una gestione del tempo evento complessa.
Verdetto: scegli Kafka Streams per trasformazioni leggere e arricchimento senza cluster da mantenere, e Flink solo quando aggregazioni pesanti, ML sullo stream o event time complesso giustificano l’infrastruttura; attiva la garanzia exactly-once solo dove i duplicati hanno conseguenze reali.
Esempio: scegliere se usare Kafka Streams
Il caso tipico è un team che vuole calcolare sessioni utente e alert comportamentali direttamente dallo stream. Prima di scegliere Kafka Streams deve valutare dimensione dello stato, retention dei changelog, gestione dei late events e impatto dei rebalancing, perché l’applicazione diventa parte della piattaforma dati. La tabella aiuta a leggere i segnali tipici.
| Evidenza osservata | Lettura prudente | Azione consigliata |
|---|---|---|
| Il throughput migliora | Potrebbe essere un effetto reale o una variazione normale di carico | Cercare un confronto e un segmento |
| Una istanza accumula più stato delle altre | La distribuzione delle chiavi è sbilanciata | Rivedere la chiave o il numero di partizioni |
| Il costo cresce insieme allo stato | L’impatto va letto sul margine, non sul totale | Stimare il trade-off di retention dei changelog |
Errori tipici da evitare
L’errore più frequente è confondere una pipeline dimostrativa con un servizio stateful pronto per la produzione. Succede quando trascuri la retention dei changelog, ignori i late events o non provi mai un restart con ricostruzione dello stato. Il sintomo è costante: la topology gira in demo e perde correttezza al primo rebalance.
Sul lato dati le trappole sono tre. La prima è sottodimensionare lo stato locale rispetto alle chiavi reali, con istanze sbilanciate. La seconda è tenere finestre troppo strette o troppo larghe rispetto ai ritardi effettivi degli eventi. La terza è attivare exactly-once ovunque, pagando latenza anche dove i duplicati sarebbero innocui. Tre controlli minimi riducono il rischio: retention dei changelog dichiarata, prova di restart superata e garanzia di processamento scelta per motivi espliciti.
L’esempio che fa da riferimento
Kafka Streams implementa la semantica exactly-once dal 2017 usando le transazioni di Kafka per scritture atomiche su più topic. Ben Stopford ha sistematizzato nel 2018, in Designing Event-Driven Systems, i pattern di stato, changelog e join che rendono queste applicazioni recuperabili dopo crash e rebalance. Da allora la regola operativa è stabile: lo stato locale si ricostruisce dai changelog e la correttezza si gioca su retention e gestione dei late events. La lezione resta quella del disegno stateful: la semplicità del deploy non elimina il costo dello stato.
Domande per verificare la lezione
- Dove vive lo stato della tua topology e come si ricostruisce dopo un restart?
- Quanta retention tieni sui changelog e perché basta per i tuoi late events?
- Quando attivi la garanzia exactly-once e quale costo accetti in cambio?
- Perché Kafka Streams basta per il tuo caso invece di un cluster Flink separato?
Bloccato su questo argomento o vuoi applicarlo al tuo caso? Prenota una call di 15 minuti con un analista esperto.
Percorso collegato
Lezioni da leggere insieme
Questi collegamenti portano la lezione dentro il resto del corso: basi da riprendere, passaggi successivi e connessioni tematiche tra moduli.