Sobre o projeto
RealTime-ECommerce-Analytics (RTEA) é um pipeline de streaming de alta performance escrito em Python que captura dados de clickstream de comércio eletrônico de uma API REST e os encaminha pelo Confluent Cloud Kafka para o ClickHouse, permitindo consultas analíticas em tempo real.
## Arquitetura
O pipeline é composto por quatro camadas:
**Fonte de Dados** — Um gerador simulado de clickstream hospedado na Railway.app, produzindo aproximadamente 10.000 mensagens por segundo no formato JSON.
**Message Broker** — Confluent Cloud Kafka com um tópico `clickstream-events` em 6 partições, autenticação SASL_SSL, compressão Snappy e retenção de 7 dias.
**Camada de Processamento** — Serviço produtor utilizando `kafka-python` com configurações otimizadas (batches de 16 KB, compressão gzip, buffer de 32 MB, taxa limitada a 10K msg/seg) e serviço consumidor com processamento em lotes (1000 eventos por lote, auto-commit de offsets e transformação de campos).
**Armazenamento Analítico** — ClickHouse Cloud com tabela de engine MergeTree ordenada por `(timestamp, user_id)`, compressão LZ4 e TTL de 30 dias.
## Fluxo de Dados
1. O produtor consulta a API externa a cada segundo, coleta mensagens em uma fila e as publica em lotes no tópico do Kafka.
2. Múltiplas instâncias do consumidor se inscrevem nas partições atribuídas, desserializam JSON, transformam timestamps e campos, e gravam lotes no ClickHouse.
3. Os usuários consultam a tabela resultante `ecommerce.clickstream_events` via SQL.
## Configuração Principal
### Produtor
- Tamanho do lote: 16 KB
- Compressão: gzip
- Memória do buffer: 32 MB
- Retransmissões: 3
### Consumidor
- Tamanho do lote: 1000 eventos
- Timeout do lote: 5 segundos
- Auto-commit de offsets habilitado
### Tabela 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)
```
## Desempenho Relatado
- Vazão sustentada: 650+ eventos/segundo
- Total de eventos processados: 510.414+
- Latência end-to-end sub-segundo
## Início Rápido
1. Instale as dependências com `uv sync`
2. Copie `.env.example` para `.env` e preencha as credenciais do Confluent Cloud e do ClickHouse
3. Execute o produtor: `uv run python services/producer/src/simple_producer.py`
4. Execute o consumidor: `uv run python services/consumer/src/simple_consumer.py`
5. Verifique os dados: `uv run python check_clickhouse.py`
## Solução de Problemas
O repositório inclui scripts de diagnóstico (`debug_kafka.py`, `test_clickhouse_insert.py`) e aborda modos comuns de falha, incluindo timeouts de conexão da API, lag do consumidor Kafka, erros de desserialização de mensagens e incompatibilidades de schema/insert no ClickHouse.
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.