Proyecto Final End-to-End
Integrando todo lo aprendido. Vamos a diseñar y construir una arquitectura de datos capaz de soportar millones de eventos por día usando los componentes estándar de la industria.
El Problema de Negocio
"Globo Rides" es una aplicación de viajes compartidos ficticia. Necesitan saber dos cosas:
- En tiempo real: Dónde hay picos de demanda de viajes (surge pricing) para enviar alertas a los conductores en menos de 30 segundos.
- Histórico (Batch): Reportes financieros diarios y semanales perfectos para los inversores.
Para lograrlo, implementaremos una Arquitectura Lambda, que maneja ambos casos de uso simultáneamente.
Diseño de Arquitectura
⬇ (Debezium CDC)
[Apache Kafka]
↙ ↘
[Spark Streaming] [S3 / Data Lake (Raw)]
(Speed Layer) (Batch Layer)
⬇ ⬇
[Redis / DynamoDB] [Snowflake / dbt]
↘ ↙
[Power BI / Dashboard App]
Todos estos componentes serán orquestados por Apache Airflow y empaquetados en un archivo docker-compose.yml gigante para desarrollo local.
Capa Speed (Tiempo Real)
Extraemos datos de la base de datos transaccional usando Change Data Capture (CDC).
- Debezium: Lee el log de transacciones (WAL) de PostgreSQL y envía un mensaje a Kafka cada vez que se inserta un viaje.
- Kafka: Recibe el evento en el topic
rides.events. - Spark Streaming: Un script en PySpark lee el topic, agrupa los viajes por zona cada 30 segundos, y si detecta un pico (>100 viajes/minuto en una zona), escribe la alerta en Redis.
# PySpark Streaming pseudocódigo
df = spark.readStream.format("kafka").option("subscribe", "rides.events").load()
parsed_df = parse_json_from_kafka(df)
# Ventana de 1 minuto deslizable cada 30 segundos
demand_df = parsed_df.groupBy(
window(col("timestamp"), "1 minute", "30 seconds"),
col("zone_id")
).count()
surge_alerts = demand_df.filter(col("count") > 100)
surge_alerts.writeStream.foreachBatch(save_to_redis).start()
Capa Batch (Histórico Correcto)
Mientras Spark calcula el tiempo real, otro proceso lee de Kafka y guarda todos los eventos crudos en Amazon S3 como archivos Parquet. Aquí entra Airflow.
- Un DAG en Airflow se ejecuta todos los días a las 2:00 AM.
- El DAG invoca el operador de Snowflake para ejecutar
COPY INTOdesde S3 a Snowflake (Capa Bronze). - El DAG ejecuta
dbt build. - dbt limpia los datos (Capa Silver), y luego calcula las facturaciones agregadas, impuestos, y comisiones de conductores (Capa Gold).
- Un test de dbt asegura que los ingresos nunca sean negativos.
Capa Serving
Los analistas y directivos no consumen ni Kafka ni S3. Usan la capa de servicio.
- La App de conductores de la empresa lee directamente de Redis para mostrar el mapa de calor de "Surge Pricing" en milisegundos.
- El equipo de Finanzas usa Power BI conectado a los modelos Gold de Snowflake usando DirectQuery para sus reportes mensuales.
Muestra que no eres solo un "SQL monkey". Entiendes infraestructura, orquestación, tradeoffs entre latencia (Speed) y precisión (Batch), y dominas las herramientas estándar de 2026. Sube esto a tu GitHub con un buen README y destacarás.