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

Project: end-to-end real-time pipeline

Build a complete pipeline from Kafka to ClickHouse to live dashboards.

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

What you will learn

  • 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

Project: end-to-end real-time pipeline

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.

  1. Avvia Kafka e ClickHouse e verifica che entrambi rispondano prima di produrre dati.
  2. Pubblica eventi ecommerce con producer Python a ritmo costante e chiavi stabili.
  3. Crea tabelle di ingestion e viste aggregate per minuto con metriche di ordini e ricavi.
  4. Costruisci tre pannelli Grafana per ordini, ricavi e funnel con baseline esplicite.
  5. Configura un alert di anomalia con soglia, owner e 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 with owner e runbook.

Delivery

  • Running Python producer (≥10 msg/sec)
  • ClickHouse with MV populated in real time
  • Grafana dashboard with 3 working panels
  • Alert configured and tested

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

  1. Quale ritmo e chiavi usi nel producer per test realistici?
  2. Quale vista aggregata alimenta ordini, ricavi e funnel?
  3. Quali tre pannelli mostrano se la pipeline è sana?
  4. Quando il tuo alert di anomalia deve scattare e chi avvisa?
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