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.