Go to main content
Project: end-to-end Kafka pipeline - official lesson image on GinnyTech, created by AD

Project: end-to-end Kafka pipeline

Build a complete pipeline with Kafka, producer, consumer, and Kafka Streams.

AD
Created byAndrii Dyshkantiuk
Lesson 119 / 236Level: AdvancedDuration: 28 minPrerequisites: 1

What you will learn

  • Costruire una pipeline end-to-end con producer, join Kafka Streams e monitoraggio del lag
  • Scegliere chiave, partizioni e finestra di join per preservare l'ordine degli eventi
  • Documentare owner, retention e runbook per difendere il rilascio in review

Project: end-to-end Kafka pipeline

Questo laboratorio chiude il binario ml-tabellare del modulo: metti insieme producer, topic, join e monitoraggio in una pipeline completa, e impari a difenderla in review come se fosse già in produzione.

L’idea in una frase

Questo progetto è una pipeline completa che unisce producer, topic, join e monitoraggio in un rilascio difendibile in review.

La procedura in cinque passi

  1. Crea i topic orders, deliveries ed enriched_orders con partizioni e replica dichiarate.
  2. Avvia il producer Python con chiave order_id e verifica l’ordinamento per ordine.
  3. Attiva il join Kafka Streams con finestra di 30 minuti verso enriched_orders.
  4. Misura il consumer lag e tienilo sotto la soglia prima del rilascio.
  5. Documenta owner, retention, contratti di schema e runbook di incidente.

Come leggere il progetto

Questo progetto unisce sorgenti, topic, schema, consumer, sink e monitoraggio in un’unica pipeline funzionante. La sfida non è far passare un messaggio in una demo, cosa che riesce a chiunque, ma tenere governati replay, duplicati, evoluzione dello schema e lag quando la pipeline diventa una dipendenza reale del business. È un lab: il punto non è accumulare definizioni ma arrivare a una pipeline difendibile in review.

Affrontalo come review architetturale, non come copia e incolla. Ogni topic deve avere un owner, una retention dichiarata, uno schema, un insieme atteso di consumer e un criterio di qualità. Se una parte della pipeline non si spiega in termini di responsabilità e modalità di guasto, non è pronta per la produzione, per quanto giri bene in locale.

La prima domanda non è “quale metrica calcolo” ma “quale decisione dovrà migliorare grazie a questa pipeline”. Una dashboard o una query hanno valore solo se riducono l’incertezza di una scelta concreta. Se non cambiano alcuna decisione, sono documentazione o teatro tecnico.

Architettura del caso

Il caso ricostruisce il flusso di una piattaforma di food delivery: gli ordini arrivano da un microservizio, sono arricchiti con lo stato della consegna tramite Kafka Streams e finiscono in un topic da cui leggono i consumer analitici. La decisione finale, mandare in produzione o no, dipende da test di replay, contratti di schema, gestione dei duplicati e runbook per gli incidenti.

Prima di procedere fissa l’unità di lavoro (topic, evento, schema, producer, consumer o stream processor), il segnale osservato (latenza, throughput, lag, compatibilità schema, perdita dati), la baseline di lettura e la decisione attesa. Il rischio costante è scambiare un numero disponibile per una prova sufficiente: una pipeline che gira non è ancora una pipeline affidabile.

Fase 1: setup e topic (20 min)

docker-compose up -d kafka zookeeper schema-registry
kafka-topics --create --topic orders --partitions 8 --replication-factor 1
kafka-topics --create --topic deliveries --partitions 8
kafka-topics --create --topic enriched_orders --partitions 8

Il numero di partizioni va deciso ora, perché aumentarlo dopo ridistribuisce le chiavi e può rompere le garanzie di ordine su cui contano i consumer.

Fase 2: producer Python (20 min)

# Simulate orders from the orders microservice
for i in range(1000):
    order = {"order_id": i, "restaurant_id": random.randint(1,50),
             "amount": round(random.uniform(10,100),2),
             "timestamp": time.time()}
    producer.produce('orders', key=str(order['order_id']),
                     value=json.dumps(order))

La chiave dell’ordine determina la partizione, quindi tutti gli eventi di uno stesso ordine restano ordinati tra loro, ed è ciò che serve per arricchirli correttamente più avanti.

Phase 3: Kafka Streams (30 min)

// Enrich orders with delivery status
KStream<String, Order> orders = builder.stream("orders");
KStream<String, Delivery> deliveries = builder.stream("deliveries");
orders.join(deliveries, (order, delivery) ->
    new EnrichedOrder(order, delivery.getStatus()),
    JoinWindows.of(Duration.ofMinutes(30)))
.to("enriched_orders");

La finestra di join di 30 minuti è la decisione più delicata: troppo stretta e perdi gli abbinamenti delle consegne lente, troppo larga e tieni stato in memoria più del necessario.

Fase 4: consumer e monitoring (15 min)

Verifica il consumer lag con il comando di descrizione dei consumer group. Un lag che cresce in modo lineare dice che il consumer non recupererà da solo, e va affrontato prima del rilascio.

Delivery

  • Order producer working
  • Kafka Streams join active
  • Consumer lag <1000
  • Topic enriched_orders populated

Errori tipici nel progetto

L’errore più comune è trattare la pipeline come una definizione: imparare i nomi dei componenti, ricordare due configurazioni, applicare un template. Il lavoro reale è diverso, perché bisogna capire quale problema risolve ogni pezzo, quali assunzioni contiene e cosa succede quando quelle assunzioni saltano. Una pipeline reale non vive isolata: sta dentro un sistema di decisioni, dati disponibili, vincoli tecnici, incentivi organizzativi e qualità dell’esecuzione.

Un secondo errore è presentare un numero senza dire quale decisione cambia, quale baseline lo rende interpretabile e quale rischio resta aperto. In quel caso il dato sembra preciso ma non guida l’azione. La domanda di controllo è semplice: se questo risultato fosse instabile, quale scelta sbaglieresti? Se non sai rispondere, manca ancora il collegamento tra analisi e azione.

L’esempio che fa da riferimento

Kafka nasce in LinkedIn e viene pubblicato open source nel 2011 per collegare produttori e consumatori attraverso topic condivisi invece di integrazioni punto a punto. Nel 2014 gli autori fondano Confluent e portano in produzione gli stessi pezzi di questo progetto: producer con chiavi ordinate, join per l’arricchimento e monitoraggio del lag. La finestra di join, il versionamento degli schemi e le soglie di lag restano le tre decisioni che separano una demo da una pipeline difendibile in review. Questo progetto le mette in fila nello stesso ordine in cui si rompono in produzione.

Verdetto: la pipeline è pronta per la produzione solo quando ogni topic ha owner, retention dichiarata, schema e criterio di qualità, e la finestra di join, il versionamento degli schemi e le soglie di lag reggono replay, duplicati e incidenti: una pipeline che gira non è ancora una pipeline affidabile.

Domande per verificare la lezione

  1. Quale decisione di rilascio dipende dai test di replay della tua pipeline?
  2. Come gestisci duplicati ed evoluzione dello schema senza fermare i consumer?
  3. Quale soglia di lag blocca il rilascio e perché è quella giusta?
  4. Chi possiede ogni topic e quale runbook segue in caso di incidente?
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