Vai al contenuto principale
Kafka Streams: processare eventi con Java - immagine ufficiale della lezione su GinnyTech, creata da AD

Kafka Streams: processare eventi con Java

Introduzione a Kafka Streams per trasformazioni stateful su flussi di eventi senza cluster esterno.

AD
Creato daAndrii Dyshkantiuk
Lezione 116 / 236Livello: AvanzatoDurata: 22 minPrerequisiti: 1

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

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

  1. Disegna la topology come flusso da topic sorgente a topic di destinazione.
  2. Dimensiona stato locale, retention dei changelog e gestione dei late events.
  3. Configura aggregazioni con finestre e join con tabelle compatte.
  4. Attiva la garanzia exactly-once solo dove i duplicati hanno conseguenze reali.
  5. 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.

La scelta tra Kafka Streams e Flink dipende da quanto è pesante l’elaborazione e da quanta infrastruttura vuoi gestire.

Kafka StreamsFlink
DeployLibreria embeddedCluster separato
ComplessitàSemplice, Java/Kafka nativoComplesso, richiede infrastruttura
Use case idealeTrasformazioni leggere, arricchimentoAggregazioni pesanti, ML su stream, event time complesso
Throughput massimoMilioni/msg secDecine 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 osservataLettura prudenteAzione consigliata
Il throughput miglioraPotrebbe essere un effetto reale o una variazione normale di caricoCercare un confronto e un segmento
Una istanza accumula più stato delle altreLa distribuzione delle chiavi è sbilanciataRivedere la chiave o il numero di partizioni
Il costo cresce insieme allo statoL’impatto va letto sul margine, non sul totaleStimare 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

  1. Dove vive lo stato della tua topology e come si ricostruisce dopo un restart?
  2. Quanta retention tieni sui changelog e perché basta per i tuoi late events?
  3. Quando attivi la garanzia exactly-once e quale costo accetti in cambio?
  4. Perché Kafka Streams basta per il tuo caso invece di un cluster Flink separato?
Serve una mano concreta?

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

Prenota una call