Sobre el proyecto
RealTime-ECommerce-Analytics (RTEA) es un pipeline de streaming de alto rendimiento escrito en Python que captura datos de clickstream de comercio electrónico desde una API REST y los enruta a través de Confluent Cloud Kafka hacia ClickHouse para consultas analíticas en tiempo real.
## Arquitectura
El pipeline consta de cuatro capas:
**Fuente de Datos** — Un generador de clickstream simulado mediante API alojado en Railway.app que produce aproximadamente 10.000 mensajes por segundo en formato JSON.
**Broker de Mensajes** — Confluent Cloud Kafka con un tema `clickstream-events` distribuido en 6 particiones, autenticación SASL_SSL, compresión Snappy y retención de 7 días.
**Capa de Procesamiento** — Servicio productor que utiliza `kafka-python` con configuraciones optimizadas (batches de 16 KB, compresión gzip, buffer de 32 MB, limitado a 10K msg/seg) y servicio consumidor con procesamiento por lotes (1000 eventos por lote, auto-commit de offsets, transformación de campos).
**Almacenamiento Analítico** — ClickHouse Cloud con tabla de motor MergeTree ordenada por `(timestamp, user_id)`, compresión LZ4 y TTL de 30 días.
## Flujo de Datos
1. El productor consulta la API externa cada segundo, recolecta mensajes en una cola y los publica en batches en el tema de Kafka.
2. Múltiples instancias consumidoras se suscriben a las particiones asignadas, deserializan JSON, transforman timestamps y campos, y escriben batches en ClickHouse.
3. Los usuarios consultan la tabla resultante `ecommerce.clickstream_events` mediante SQL.
## Configuración Clave
### Productor
- Tamaño del batch: 16 KB
- Compresión: gzip
- Memoria del buffer: 32 MB
- Reintentos: 3
### Consumidor
- Tamaño del batch: 1000 eventos
- Timeout del batch: 5 segundos
- Auto-commit de offsets habilitado
### Tabla ClickHouse
```sql
CREATE TABLE clickstream_events (
user_id String,
session_id String,
timestamp UInt64,
event_type String,
product_id String,
product_category String,
price Float64,
quantity Int32,
source String,
received_time DateTime DEFAULT now()
) ENGINE = MergeTree()
ORDER BY (timestamp, user_id)
```
## Rendimiento Reportado
- Throughput sostenido: más de 650 eventos/segundo
- Total de eventos procesados: 510.414+
- Latencia end-to-end inferior a un segundo
## Inicio Rápido
1. Instalar dependencias con `uv sync`
2. Copiar `.env.example` a `.env` y completar las credenciales de Confluent Cloud y ClickHouse
3. Ejecutar el productor: `uv run python services/producer/src/simple_producer.py`
4. Ejecutar el consumidor: `uv run python services/consumer/src/simple_consumer.py`
5. Verificar datos: `uv run python check_clickhouse.py`
## Solución de Problemas
El repositorio incluye scripts de diagnóstico (`debug_kafka.py`, `test_clickhouse_insert.py`) y aborda modos de fallo comunes como timeouts de conexión de API, lag del consumidor de Kafka, errores de deserialización de mensajes e incompatibilidades de esquema en ClickHouse.
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.