Salta el contingut

Pràctica PR507504: Streaming amb Kafka i Databricks

Informació de la pràctica

Camp Detall
Codi PR507504
Mòdul 5075 — Big Data Aplicat
RA RA2 — Ecosistema i processament distribuït
Durada estimada 6-8 hores (2 sessions)
Modalitat Individual
Lliurament Codi font + notebook Databricks exportat + informe amb captures
Qualificació Sobre 10 punts, amb rúbrica

Objectius

  • Desplegar un clúster Kafka local amb Docker Compose i crear topics amb múltiples particions.
  • Implementar un productor Python que simuli sensors IoT enviant lectures contínues.
  • Implementar un consumidor Python que detecti anomalies en temps real a partir de llindars.
  • Processar el mateix flux de dades amb Spark Structured Streaming a Databricks Community Edition, aplicant finestres temporals i agregacions.

Preparació de l'entorn

# docker-compose.yml
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.6.1
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  kafka:
    image: confluentinc/cp-kafka:7.6.1
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "29092:29092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'

  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    depends_on:
      - kafka
    ports:
      - "8080:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092
docker compose up -d
pip install kafka-python==2.0.2

Crea un compte gratuït a Databricks Community Edition i un clúster amb Databricks Runtime 14.3 LTS.

Per què un fitxer i no Kafka directament a Databricks?

Databricks Community Edition executa els clústers en un entorn aïllat sense accés a la xarxa local, així que la integració directa Kafka→Databricks requereix Databricks Standard o Confluent Cloud. El mecanisme de Structured Streaming és exactament el mateix independentment de la font (Kafka, fitxers, Delta Lake), així que la Fase 3 simula l'arribada contínua de dades amb fitxers CSV.

Fase 1: Topics i cua manual (45 min)

docker exec -it kafka kafka-topics --create \
  --bootstrap-server localhost:29092 \
  --topic sensor_data --partitions 3 --replication-factor 1

docker exec -it kafka kafka-topics --describe \
  --bootstrap-server localhost:29092 --topic sensor_data

Prova, en dos terminals, un kafka-console-producer i un kafka-console-consumer --from-beginning sobre el topic, enviant manualment un parell de lectures JSON ({"sensor_id": "S001", "temperatura": 23.5}).

Pregunta de reflexió 1

Per què creem 3 particions? Quina relació té el nombre de particions amb la paral·lelització del processament?

Fase 2: Productor Python — sensors IoT (60 min)

Crea producer.py: envia, cada segon, una lectura per a cada un de 5 sensors simulats (S001-S005), cadascun amb una temperatura base diferent i un 5% de probabilitat d'anomalia (pic de +10 a +20 °C). Usa KafkaProducer amb value_serializer JSON i el sensor_id com a clau del missatge, perquè totes les lectures d'un mateix sensor vagin sempre a la mateixa partició i es puguin consultar en ordre.

from kafka import KafkaProducer
import json, time, random
from datetime import datetime

producer = KafkaProducer(
    bootstrap_servers=['localhost:9092'],
    value_serializer=lambda v: json.dumps(v).encode('utf-8'),
    key_serializer=lambda k: k.encode('utf-8'))

TEMPERATURA_BASE = {'S001': 22.0, 'S002': 35.0, 'S003': 18.5, 'S004': 25.0, 'S005': 40.0}

def genera_lectura(sensor_id):
    temp = TEMPERATURA_BASE[sensor_id] + random.uniform(-2.0, 2.0)
    if random.random() < 0.05:
        temp += random.uniform(10.0, 20.0)
    return {"sensor_id": sensor_id, "temperatura": round(temp, 2),
            "timestamp": datetime.now().isoformat(), "ubicacio": f"Zona_{sensor_id[-1]}"}

while True:
    for sensor_id in TEMPERATURA_BASE:
        event = genera_lectura(sensor_id)
        future = producer.send("sensor_data", key=sensor_id, value=event)
        meta = future.get(timeout=10)
        print(f"Particio {meta.partition} Offset {meta.offset} -> {event}")
    time.sleep(1)

Pregunta de reflexió 2

Per quin motiu s'usa el sensor_id com a clau del missatge? Quina garantia dona respecte a l'ordenació dels events?

Fase 3: Consumidor Python — detecció d'anomalies (60 min)

Crea consumer.py: un KafkaConsumer (group_id='sensor_monitor_group') que, per a cada lectura rebuda, actualitzi una mitjana mòbil de les 10 últimes lectures del sensor i verifiqui si la temperatura supera un llindar {min, max} propi de cada sensor, imprimint una alerta (fred/calor) quan calgui.

Pregunta de reflexió 3

Quina diferència hi ha entre auto_offset_reset='latest' i auto_offset_reset='earliest'? En quins casos d'ús usaries cadascun?

Exercici complementari

Executa dos terminals amb consumer.py, cadascun amb un group_id diferent, i després amb el mateix group_id. Observa quins missatges rep cadascun i explica la diferència (vegeu consumer groups).

Fase 4: Generació del dataset per a Databricks (30 min)

Crea genera_dataset.py: genera 60 fitxers CSV (un per "minut" simulat, amb 5 lectures per sensor cada un), reproduint el mateix esquema que el productor Kafka, dins d'una carpeta sensor_stream_data/. Comprimeix la carpeta i puja-la a Databricks (Data → Add Data → Upload File → DBFS, ruta /FileStore/sensor_stream_data/).

Fase 5: Structured Streaming a Databricks (90 min)

Crea un notebook streaming_sensors amb:

  1. Definició explícita de l'esquema (StructType) — en streaming l'esquema sempre ha de ser explícit, mai inferit automàticament.
  2. Un readStream sobre la carpeta DBFS, tractant cada fitxer CSV nou com una arribada contínua de dades.
  3. Una agregació amb finestra temporal (window) que calculi, per a cada sensor i finestra de 1 minut, la temperatura mitjana, màxima i mínima.
  4. Marcatge de les files amb temperatura fora del llindar del sensor (when/otherwise), reaprofitant la mateixa lògica del consumidor Python.
df_stream = spark.readStream.schema(schema).csv("/FileStore/sensor_stream_data/")

df_agregat = df_stream \
    .withColumn("event_time", to_timestamp(col("timestamp"))) \
    .groupBy(window(col("event_time"), "1 minute"), col("sensor_id")) \
    .agg(avg("temperatura").alias("temp_mitjana"),
         max("temperatura").alias("temp_maxima"),
         min("temperatura").alias("temp_minima"),
         count("*").alias("num_lectures"))

query = df_agregat.writeStream.outputMode("complete").format("memory") \
    .queryName("sensors_agregats").start()

Consulta Structured Streaming i Windowing i Joins per a aprofundir en finestres temporals amb Spark.

Entrega

  • Codi de producer.py, consumer.py i genera_dataset.py.
  • Captura de Kafka UI mostrant el topic sensor_data amb les 3 particions actives.
  • Captura del consumidor detectant almenys una anomalia en temps real.
  • Notebook Databricks exportat (.dbc o .ipynb) amb l'agregació per finestres funcionant.
  • Respostes a les 3 preguntes de reflexió i a l'exercici complementari.

Rúbrica: Rúbrica PR507504