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
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:
- Definició explícita de l'esquema (
StructType) — en streaming l'esquema sempre ha de ser explícit, mai inferit automàticament. - Un
readStreamsobre la carpeta DBFS, tractant cada fitxer CSV nou com una arribada contínua de dades. - 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. - 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.pyigenera_dataset.py. - Captura de Kafka UI mostrant el topic
sensor_dataamb les 3 particions actives. - Captura del consumidor detectant almenys una anomalia en temps real.
- Notebook Databricks exportat (
.dbco.ipynb) amb l'agregació per finestres funcionant. - Respostes a les 3 preguntes de reflexió i a l'exercici complementari.
Rúbrica: Rúbrica PR507504