
Progetto: pipeline real-time end-to-end
Costruire una pipeline completa da Kafka a ClickHouse a dashboard live.
Cosa imparerai
- Costruire una pipeline end-to-end da Kafka a ClickHouse a dashboard Grafana
- Creare viste aggregate per minuto con metriche di ordini e ricavi
- Configurare un alert di anomalia con soglia, owner e prova di notifica
Collegamenti
Progetto: pipeline real-time end-to-end
Chiudiamo il modulo con un lavoro completo, sempre sul binario ml-tabellare: qui non si spiega, si costruisce. Dovrai collegare producer Kafka, tabelle ClickHouse, viste aggregate, dashboard Grafana e alert in un’unica pipeline che regga duplicati, ritardi e picchi.
L’obiettivo del progetto
Il progetto collega producer Kafka, tabelle ClickHouse, viste aggregate, dashboard Grafana e alert in un’unica pipeline che regge duplicati, ritardi e picchi. Ogni pezzo ha un motivo e un criterio di fallimento esplicito.
La sequenza di montaggio
Segui questi cinque passi nell’ordine indicato.
- Avvia
Kafkae ClickHouse e verifica che entrambi rispondano prima di produrre dati. - Pubblica eventi ecommerce con producer
Pythona ritmo costante e chiavi stabili. - Crea tabelle di ingestion e viste aggregate per minuto con metriche di ordini e ricavi.
- Costruisci tre pannelli
Grafanaper ordini, ricavi e funnel con baseline esplicite. - Configura un alert di anomalia con soglia,
ownere prova di notifica prima della consegna.
Fase 1: setup Kafka e ClickHouse
docker-compose up -d kafka clickhouse
# Verifica: docker-compose logs kafka | grep "started"
Fase 2: producer Python
from kafka import KafkaProducer
producer = KafkaProducer(bootstrap_servers='localhost:9092',
value_serializer=lambda v: json.dumps(v).encode('utf-8'))
while True:
event = {"user_id": random.randint(1,1000), "event": random.choice(
["page_view","add_cart","purchase"]), "amount": round(random.uniform(10,200),2),
"timestamp": time.time()}
producer.send('ecommerce_events', event)
time.sleep(0.1) # 10 msg/sec
Fase 3: tabelle ClickHouse
CREATE TABLE kafka_events (...) ENGINE = Kafka SETTINGS ...;
CREATE MATERIALIZED VIEW mv_events_per_min
ENGINE = SummingMergeTree() ORDER BY (minute, event)
AS SELECT toStartOfMinute(toDateTime(timestamp)) AS minute,
event, count() AS cnt, sum(amount) AS revenue
FROM kafka_events GROUP BY minute, event;
Fase 4: dashboard con Grafana
Crea tre pannelli. Il primo mostra gli ordini al minuto, come time series su mv_events_per_min filtrato per event='purchase'. Il secondo mostra il revenue al minuto, come time series su somma dei ricavi. Il terzo costruisce il funnel da pagina vista ad aggiunta carrello fino ad acquisto con tassi di conversione tra stadi.
Fase 5: alert di anomalia
Configura un alert in Grafana: se il tasso di purchase scende sotto la baseline del cinquanta per cento, invia una notifica Slack con owner e runbook.
Consegna
- Producer Python in esecuzione (≥10 msg/sec)
- ClickHouse con MV popolata in tempo reale
- Dashboard Grafana con 3 pannelli funzionanti
- Alert configurato e testato
Leggi il caso come design review. Ogni componente deve avere motivo, owner e criterio di fallimento. La domanda non è se farlo in real-time ma quali decisioni migliorano abbastanza da giustificare complessità, monitoraggio e costo operativo.
Gli errori da non ripetere
Il primo errore è lavorare su medie globali che nascondono segmenti opposti. Il secondo è ignorare qualità del dato con duplicati, timezone incoerenti e definizioni instabili. Il terzo è scambiare correlazione per causalità su feature e conversioni. Ogni analisi porta definizione esplicita, confronto per segmento e verifica contro periodo precedente o controllo.
Il progetto che ha ispirato tutto: Kafka a LinkedIn
LinkedIn crea Kafka nel 2011 per unificare log, eventi e pipeline in un unico flusso durevole e replayabile. Il progetto separa producer, broker e consumatori così ogni team rilegge gli stessi eventi senza export notturni. Da quel disegno discende la pipeline della lezione: producer Python, tabelle di ingestion, viste aggregate e dashboard live sullo stesso stream. Il criterio resta quello: un solo flusso ordinato alimenta operative, analitica e alert senza copie incoerenti.
Verdetto: il real-time vince solo quando le decisioni migliorano abbastanza da giustificare complessità, monitoraggio e costo operativo: ogni componente deve avere motivo, owner e criterio di fallimento, e l’alert di anomalia deve scattare su una soglia scritta con baseline esplicita.
Domande per ripassare
- Quale ritmo e chiavi usi nel producer per test realistici?
- Quale vista aggregata alimenta ordini, ricavi e funnel?
- Quali tre pannelli mostrano se la pipeline è sana?
- Quando il tuo alert di anomalia deve scattare e chi avvisa?
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.