Apache Kafka Parte 2: event-driven, Kafka Streams e Connect

Riassunto

Seconda parte della guida ad Apache Kafka: come scalare i consumer group, trasformare i dati in volo con Kafka Streams, integrare Elasticsearch, Neo4j e InfluxDB con Kafka Connect e scegliere tra architettura Lambda e Kappa. Il filo conduttore è una piattaforma di social network analytics in tempo reale dal corso del Politecnico di Torino.

Avatar Alessandro Fiori

PhD · Politecnico di Torino

In sintesi. Nella Parte 1 hai visto topic, partizioni, offset e consumer group. Qui passiamo alla produzione: come i consumer group scalano in orizzontale, come Kafka Streams trasforma e aggrega i dati mentre scorrono, come Kafka Connect li porta verso Elasticsearch, Neo4j e InfluxDB senza scrivere consumer a mano, e quando scegliere un’architettura Lambda o Kappa. Il filo narrativo è un progetto reale del corso di Advanced Data Modeling del Politecnico di Torino: una piattaforma di social network analytics alimentata in tempo reale.

Una startup torinese lancia una piattaforma tematica sul patrimonio culturale e gastronomico piemontese. Gli utenti pubblicano post su luoghi, ricette ed eventi, si seguono a vicenda, mettono like e commentano. Il team di prodotto vuole sapere tre cose, e le vuole adesso, non domani mattina: quali hashtag stanno diventando virali, chi sono gli utenti più influenti, quali persone suggerire a ciascun iscritto nella sezione “potresti conoscere”.

Questo è lo scenario del progetto P4 del corso — Social Network Analytics Platform — e mette insieme cinque tecnologie: MongoDB per profili e post, Neo4j per il grafo sociale, Elasticsearch per la ricerca full-text e i trend, Grafana per i KPI, e Kafka come sistema nervoso che collega tutto. Nella Parte 1 abbiamo usato Kafka per trasportare letture di sensori. Qui lo usiamo per qualcosa di più ambizioso: far reagire più sistemi allo stesso evento, nello stesso momento, senza che nessuno di loro sappia dell’esistenza degli altri.

Architettura event-driven Kafka: eventi social nei topic, Kafka Streams, sink verso MongoDB, Neo4j, Elasticsearch, Grafana
L’architettura del progetto P4: un solo flusso di eventi, quattro sistemi che lo consumano in modo indipendente

Dall’architettura request-response all’event-driven

Nella maggior parte delle applicazioni web tradizionali, quando un utente mette like a un post succede questo: il backend riceve la richiesta, scrive nel database, poi chiama in sequenza il servizio notifiche, aggiorna il contatore, magari invia qualcosa al sistema di analytics. Ogni chiamata è un punto di rottura. Se il servizio notifiche è lento, l’utente aspetta. Se il sistema di analytics è giù, il like rischia di andare perso o di generare un errore visibile.

In un’architettura event-driven il backend fa una cosa sola: pubblica un evento su Kafka — {"event_type": "like", "user_id": "u42", "target_id": "p1337", "timestamp": "..."} — e risponde subito. Tutto il resto avviene a valle, in modo asincrono. Chi è interessato a quell’evento si iscrive al topic e fa il proprio lavoro con i propri tempi.

L’evento come fatto immutabile

Il cambio di prospettiva è sottile ma importante. Un evento non è un comando (“aggiorna il contatore”) ma un fatto accaduto (“l’utente u42 ha messo like al post p1337 alle 18:04”). I fatti non si modificano e non si cancellano: si accumulano. Per questo il log persistente di Kafka è il posto naturale dove conservarli, e per questo ogni consumer può ricostruire il proprio stato rileggendo la storia dall’inizio.

Tre vantaggi concreti

  • Disaccoppiamento. Il servizio che produce i like non conosce Neo4j né Elasticsearch. Puoi aggiungere un quinto consumer domani senza toccare una riga del produttore.
  • Replay. Hai sbagliato la logica di calcolo dei trend? Correggi il codice, riporti l’offset del consumer indietro di sette giorni e ricalcoli tutto.
  • Fan-out. Lo stesso evento alimenta il grafo sociale, l’indice di ricerca e la dashboard. Un solo dato di origine, molte viste specializzate: è l’essenza della polyglot persistence, di cui parleremo in dettaglio nelle prossime settimane.

Consumer group e scalabilità orizzontale in produzione

Nella Parte 1 hai visto che un consumer group divide le partizioni di un topic tra i suoi membri. In produzione questo meccanismo diventa lo strumento principale per scalare. Se il topic social-events ha 6 partizioni e il consumer che scrive su Neo4j non tiene il passo, avvii una seconda e poi una terza istanza con lo stesso group.id: Kafka redistribuisce le partizioni e il carico si divide. Oltre le 6 istanze, però, non guadagni nulla — le istanze in più restano inattive. Il numero di partizioni è il tetto massimo del parallelismo, ed è per questo che va deciso con attenzione alla creazione del topic.

Il rebalancing e il nuovo protocollo

Ogni volta che un consumer entra o esce dal gruppo, Kafka esegue un rebalance: riassegna le partizioni. Con il protocollo classico, durante il rebalance tutto il gruppo si fermava per qualche secondo — un problema serio con decine di consumer. Kafka 4.0 ha reso disponibile il nuovo protocollo di gruppo lato broker (KIP-848), che rende i rebalance incrementali: solo le partizioni che devono davvero cambiare proprietario vengono spostate. Con Kafka 4.2 (febbraio 2026) lo stesso approccio è arrivato in versione GA, con un set di funzionalità ancora limitato, anche per le applicazioni Kafka Streams (KIP-1071). Se parti oggi con un cluster nuovo, vale la pena abilitarlo.

Semantiche di consegna: quante volte arriva un messaggio?

Kafka offre tre garanzie, e scegliere quella giusta dipende da quando fai il commit dell’offset:

  • At-most-once: fai commit prima di elaborare. Se il consumer va in crash a metà, il messaggio è perso. Accettabile per metriche approssimate.
  • At-least-once: fai commit dopo aver elaborato. Se il consumer va in crash dopo aver scritto ma prima del commit, rielaborerà lo stesso messaggio. È il default più diffuso.
  • Exactly-once: possibile con producer idempotenti e transazioni Kafka, a patto che anche la destinazione sia Kafka (è il caso di Kafka Streams).

Nel mondo reale, la strategia più robusta è at-least-once con scritture idempotenti sulla destinazione: se lo stesso evento arriva due volte, il risultato non cambia. In Neo4j questo si ottiene con MERGE al posto di CREATE, esattamente come suggerito agli studenti del corso:

from kafka import KafkaConsumer
from neo4j import GraphDatabase
import json

consumer = KafkaConsumer(
    "social-events",
    bootstrap_servers="localhost:9092",
    group_id="neo4j-writer",
    enable_auto_commit=False,
    value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)
driver = GraphDatabase.driver("neo4j+s://xxxx.databases.neo4j.io", auth=("neo4j", "***"))

CYPHER = {
    "follow": "MERGE (a:User {id:$uid}) MERGE (b:User {id:$tid}) MERGE (a)-[:FOLLOWS]->(b)",
    "like":   "MERGE (a:User {id:$uid}) MERGE (p:Post {id:$tid}) MERGE (a)-[:LIKES]->(p)",
}

for msg in consumer:
    ev = msg.value
    query = CYPHER.get(ev["event_type"])
    if query:
        with driver.session() as s:
            s.run(query, uid=ev["user_id"], tid=ev["target_id"])
    consumer.commit()  # commit solo dopo la scrittura: at-least-once

Grazie a MERGE, un replay del topic non crea relazioni duplicate: il grafo resta coerente anche se rileggi una settimana di eventi.

Kafka Streams: trasformare i dati mentre scorrono

Finora i consumer si limitano a spostare dati da Kafka a un database. Ma spesso vuoi elaborare il flusso prima che arrivi a destinazione: filtrare gli eventi di test, arricchire un like con la categoria del post, contare gli hashtag degli ultimi dieci minuti. Kafka Streams è la libreria Java (fa parte del progetto Apache Kafka) pensata esattamente per questo. Non è un cluster separato da gestire come Spark o Flink: è una normale applicazione che legge da topic, elabora e scrive su altri topic. Per scalarla, avvii più istanze — il meccanismo è lo stesso dei consumer group.

Operazioni stateless

Sono le trasformazioni che guardano un evento alla volta, senza memoria del passato: filter (scarta gli eventi dei bot), map e mapValues (normalizza il formato del timestamp), flatMap (da un post con cinque hashtag genera cinque eventi separati), split/branch (instrada like, follow e commenti su topic distinti). Sono veloci, semplici da scalare e non richiedono storage locale.

Operazioni stateful: aggregazioni, finestre e join

Le cose si fanno interessanti quando serve memoria. Contare gli hashtag richiede di ricordare i conteggi precedenti; unire un evento di like con i metadati del post richiede di avere quei metadati a portata di mano. Kafka Streams gestisce questo stato in uno state store locale (RocksDB di default) e ne salva ogni modifica su un topic interno di changelog. Se un’istanza muore, un’altra ricostruisce lo stato rileggendo il changelog e riprende da dove si era interrotta.

Ecco la topologia che calcola gli hashtag di tendenza per finestre di 10 minuti:

StreamsBuilder builder = new StreamsBuilder();

builder.stream("posts", Consumed.with(Serdes.String(), postSerde))
    .flatMapValues(post -> post.getHashtags())          // 1 post -> N hashtag
    .selectKey((k, tag) -> tag.toLowerCase())
    .groupByKey(Grouped.with(Serdes.String(), Serdes.String()))
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(10)))
    .count(Materialized.as("hashtag-counts"))            // state store
    .toStream()
    .map((wKey, count) -> KeyValue.pair(wKey.key(), count))
    .to("trending-hashtags", Produced.with(Serdes.String(), Serdes.Long()));

Il risultato è un topic trending-hashtags sempre aggiornato, pronto per essere indicizzato in Elasticsearch o visualizzato in Grafana.

KStream e KTable: due modi di leggere lo stesso log

Kafka Streams introduce una distinzione concettuale che vale la pena capire bene. Un KStream interpreta il topic come una sequenza di eventi indipendenti: tre like dello stesso utente sono tre fatti distinti. Una KTable interpreta il topic come un changelog: per ogni chiave conta solo l’ultimo valore, come una riga di una tabella che viene aggiornata. Il profilo di un utente è una KTable (vuoi l’ultima versione), i suoi like sono un KStream. Il join tra i due — “arricchisci ogni like con la città dell’utente” — è una delle operazioni più comuni in produzione.

E se non scrivi in Java?

Kafka Streams è una libreria JVM. Se il tuo team lavora in Python hai alternative valide: Quix Streams e faust-streaming offrono API simili con finestre e stato; Flink SQL (disponibile anche su Confluent Cloud) ti permette di esprimere la stessa aggregazione in SQL, senza scrivere codice applicativo. Per un progetto didattico da 20 ore, un consumer Python che aggrega in memoria su finestre brevi è spesso sufficiente; per la produzione, meglio uno strumento che gestisca stato e fault tolerance per te.

Kafka Connect: integrazioni senza scrivere consumer

Il consumer Neo4j che abbiamo scritto sopra funziona, ma moltiplicalo per ogni destinazione e ogni sorgente e ti ritrovi con decine di piccoli script da mantenere, monitorare e aggiornare. Kafka Connect risolve questo problema con un framework standard: invece di scrivere codice, configuri un connettore con un file JSON. Connect si occupa di offset, retry, parallelismo e gestione degli errori.

Kafka Connect con connettori source da PostgreSQL via Debezium e connettori sink verso Elasticsearch, Neo4j e InfluxDB
Kafka Connect: i connettori source portano dati dentro Kafka, i sink li consegnano ai database di destinazione

Source e sink connector

I connettori sono di due tipi. I source leggono da un sistema esterno e scrivono su Kafka: il più importante è Debezium, che legge il transaction log di PostgreSQL, MySQL o MongoDB e pubblica ogni INSERT, UPDATE e DELETE come evento. È la tecnica chiamata Change Data Capture (CDC), lo standard di fatto per sincronizzare database diversi senza scritture doppie nel codice. I sink fanno il percorso inverso: leggono da un topic e scrivono su una destinazione.

Ecco un sink verso Elasticsearch per indicizzare i post del progetto P4:

{
  "name": "posts-to-elasticsearch",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "topics": "posts",
    "connection.url": "http://elasticsearch:9200",
    "key.ignore": "false",
    "schema.ignore": "true",
    "tasks.max": "3",
    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "dlq-posts",
    "errors.deadletterqueue.context.headers.enable": "true"
  }
}

Nota le ultime tre righe: con errors.tolerance e una dead letter queue, un documento malformato non blocca l’intera pipeline ma finisce in un topic a parte, dove puoi analizzarlo con calma. È uno dei dettagli che separano un prototipo da un sistema in produzione.

Neo4j e InfluxDB come destinazioni

Per il grafo esiste il Neo4j Connector for Kafka, che permette di definire una query Cypher per ogni topic: in pratica la stessa logica MERGE del nostro consumer Python, ma dichiarata in configurazione. Per le serie temporali hai due strade: un connettore sink per InfluxDB oppure Telegraf, l’agente di InfluxData, che ha un plugin di input kafka_consumer nativo. Nel progetto P1 Smart City è la soluzione più semplice per portare le letture dei sensori in InfluxDB con scritture in batch.

Il caso P4: social network analytics in tempo reale

Mettiamo insieme i pezzi. Il simulatore del progetto genera eventi ogni 0,5–2 secondi con un payload JSON uniforme. Ogni tipo di evento ha una o più destinazioni, ciascuna scelta per il tipo di domanda a cui deve rispondere:

EventoTopicChi lo consumaDomanda a cui risponde
Nuovo postpostsMongoDB + Elasticsearch (Connect)“Mostrami i post su Barolo vicino ad Alba”
Followsocial-eventsNeo4j (consumer idempotente)“Chi potrei seguire?” (friend-of-friend a 2 hop)
Like / commentosocial-eventsNeo4j + Kafka Streams“Quali post crescono più in fretta nelle prime 2 ore?”
Hashtag aggregatitrending-hashtagsElasticsearch / Grafana“Cosa è di tendenza nell’ultima ora?”
KPI di engagementengagement-metricsInfluxDB → Grafana“Post al giorno, like ratio, crescita follower”

Un dettaglio che spesso sorprende: Grafana non ha un datasource nativo per Kafka. Per mostrare i KPI in tempo reale devi prima scriverli in un database che Grafana sa interrogare — InfluxDB o PostgreSQL. Non è un limite, è coerente con la filosofia di Kafka: il log trasporta e conserva, i database specializzati servono le query. Se vuoi ripassare come costruire i pannelli, trovi tutto nella guida a Grafana per dati tecnici e IoT, mentre la modellazione del grafo sociale è spiegata nell’articolo su Neo4j.

La query che dà valore al grafo

Una volta che i follow arrivano in Neo4j in tempo reale, la raccomandazione “potresti conoscere” diventa una query Cypher di poche righe:

MATCH (me:User {id: $uid})-[:FOLLOWS]->(friend)-[:FOLLOWS]->(suggestion)
WHERE NOT (me)-[:FOLLOWS]->(suggestion) AND suggestion <> me
RETURN suggestion.id, count(friend) AS amici_in_comune
ORDER BY amici_in_comune DESC
LIMIT 10

Il valore di Kafka qui non è la query in sé, ma il fatto che il risultato riflette i follow di pochi secondi fa, non quelli di ieri notte.

Lambda vs Kappa: dove si colloca Kafka nell’architettura

Quando i dati arrivano in streaming, prima o poi si pone una domanda: come gestisco anche le analisi sullo storico completo? Le due risposte classiche hanno un nome.

Confronto tra Lambda Architecture con batch layer e speed layer e Kappa Architecture con solo stream processing
Lambda mantiene due pipeline parallele; Kappa ne usa una sola e ricalcola rileggendo il log

Lambda Architecture

Proposta da Nathan Marz nel 2011, combina tre livelli: un batch layer che ricalcola periodicamente tutto lo storico (la “verità” definitiva), uno speed layer che elabora in tempo reale gli eventi più recenti per compensare la latenza del batch, e un serving layer che unisce i due risultati. Il vantaggio è l’affidabilità: il batch corregge sempre eventuali errori dello streaming. Il difetto è che scrivi e mantieni la stessa logica due volte, con due tecnologie diverse.

Kappa Architecture

Jay Kreps — uno dei creatori di Kafka in LinkedIn — propose nel 2014 di eliminare il batch layer. Tutto passa dallo streaming; se serve ricalcolare lo storico, si rilegge il log dall’inizio con una nuova versione del codice. Una sola codebase, operazioni più semplici. Il prezzo: serve un’infrastruttura di streaming matura e una retention sufficientemente lunga su Kafka.

CriterioLambdaKappa
Rielaborazione dello storicoNaturale (job batch)Richiede retention lunga su Kafka
Complessità del codiceDue pipelineUna pipeline
Garanzia di correttezzaIl batch è la veritàSi corregge rielaborando
Competenze del teamBatch + streamingSolo streaming
Quando sceglierlaAudit e compliance stringentiAlta velocità, team agili

La raccomandazione che do agli studenti è pragmatica: parti con Kappa se il team è a suo agio con lo streaming; passa a Lambda solo se un vincolo normativo richiede un ricalcolo batch certificato. Nel progetto P4, Kappa è la scelta naturale: tutto nasce come evento e il log di Kafka è già la fonte di verità.

Gli errori più comuni (e come evitarli)

Dopo aver seguito decine di gruppi di progetto, questi sono gli inciampi che vedo più spesso:

  • Troppe poche partizioni. Un topic con 1 partizione non scala oltre un consumer. Per la didattica va bene, in produzione no. Le partizioni si possono aumentare ma mai ridurre.
  • CREATE invece di MERGE. Al primo replay il grafo si riempie di duplicati. Rendi idempotente ogni scrittura.
  • Chiave del messaggio sbagliata. L’ordine è garantito solo all’interno di una partizione. Se gli eventi di uno stesso utente devono restare in ordine, usa user_id come chiave.
  • Nessuna dead letter queue. Un solo JSON malformato può bloccare un connettore per ore.
  • Dimenticare il consumer lag. Se i consumer restano indietro, i dati “in tempo reale” diventano vecchi di minuti. Monitora il lag e imposta un alert: è la metrica più importante di un cluster Kafka.

Da dove partire

Per sperimentare senza installare nulla puoi usare il piano gratuito di Confluent Cloud, come abbiamo fatto nella Parte 1; se preferisci lavorare in locale, un docker-compose con l’immagine single-node di Kafka in modalità KRaft (senza ZooKeeper, rimosso definitivamente da Kafka 4.0) ti dà un ambiente completo in pochi minuti. Il percorso consigliato: prima un producer e un consumer Python, poi un connettore sink verso Elasticsearch, infine una piccola topologia Kafka Streams con una finestra temporale.

Con questa seconda parte chiudiamo il capitolo Kafka. Nel prossimo articolo della serie Stack Digitale 2026 scendiamo di un livello e guardiamo dove finiscono le letture dei sensori: InfluxDB, il database progettato per le serie temporali. Se stai valutando un’architettura event-driven per la tua azienda o il tuo ente e vuoi capire se Kafka è davvero necessario o se basta qualcosa di più semplice, richiedi un audit dati e AI: partiamo dai tuoi flussi reali, non dalla tecnologia di moda.