Streamline Pipeline com Apache Kafka

- Published on

No mundo do e-commerce moderno, cada transação, cada clique e cada interação de um cliente carrega informações valiosas que, se processadas em tempo real, podem gerar insights imediatos sobre comportamento de compra, riscos de fraude e oportunidades de upsell. A capacidade de reagir a dados em milissegundos não é mais um diferencial competitivo — é uma necessidade.
Foi pensando nesse desafio que desenvolvi o projeto kafka-transactions-streamline, uma pipeline de dados de transações de e-commerce de moda usando Apache Kafka, Kafka Connect e PostgreSQL, com monitoramento e alertas em tempo real. Neste artigo, vou detalhar como essa arquitetura funciona, os conceitos por trás de cada componente e por que o processamento de dados em tempo real está se tornando o padrão da indústria.
Por que Dados em Tempo Real?
Antes de mergulharmos na arquitetura técnica, é importante entender o contexto de mercado que torna soluções como esta tão relevantes.
O mercado de streaming analytics está em franca expansão. De acordo com o relatório Streaming Analytics Market Report 2026 da The Business Research Company, o mercado foi avaliado em 35,17 bilhões de dólares em 2025 e deve crescer para 46,78 bilhões de dólares em 2026, a uma taxa de crescimento anual composta (CAGR) de 33%. A expectativa é que esse crescimento exponencial continue, com o mercado projetado para alcançar 146,59 bilhões de dólares até 2030, mantendo um CAGR de 33%. Esse crescimento é impulsionado pela expansão da adoção de big data analytics, o uso crescente de sistemas de monitoramento em tempo real, o aumento do volume de transações digitais e a demanda por insights instantâneos.
O mercado de ferramentas de data pipeline, por sua vez, também apresenta um cenário de forte expansão. Avaliado em 12,53 bilhões de dólares em 2025, ele deve crescer para 15,14 bilhões de dólares em 2026 e alcançar 52,53 bilhões de dólares até 2032, com um CAGR de 22,71%. Paralelamente, o mercado de ferramentas de streaming de dados em tempo real, foco central deste projeto, foi avaliado em 8,2 bilhões de dólares em 2025 e projeta-se que chegue a 25 bilhões de dólares até 2035, a um CAGR de 11,8%.
Os números são impressionantes, mas o que eles realmente significam? As empresas estão migrando de arquiteturas baseadas em batch (processamento em lote) para arquiteturas orientadas a eventos (streaming) porque:
- Velocidade de decisão: 82% das organizações estão usando ou planejam implementar capacidades de processamento de dados em tempo real.
- Adoção massiva: o streaming de dados tornou-se uma prioridade estratégica para 86% dos líderes de TI.
- Adoção do Kafka: o Apache Kafka se consolidou como o padrão da indústria para streaming de dados, detendo uma participação de 38,7% no mercado de ferramentas de fila, mensageria e processamento em background, à frente de concorrentes como RabbitMQ (28,6%) e IBM MQ (6,3%). A plataforma está presente em mais de 80% das empresas da Fortune 100.
A Arquitetura do Projeto
O coração do projeto é uma pipeline de dados que simula transações de um e-commerce de moda, processa esses eventos em tempo real e os disponibiliza para análise e monitoramento.
Visão Geral da Arquitetura
A arquitetura segue um fluxo linear e bem definido, orquestrado inteiramente via Docker Compose:

Tecnologias Utilizadas
| Componente | Tecnologia |
|---|---|
| Mensageria | Apache Kafka 4.3 (KRaft mode, sem Zookeeper) |
| Conectores | Kafka Connect + JDBC Sink Connector |
| Banco de Dados | PostgreSQL 16 |
| Producer | Python 3.12 + kafka-python-ng |
| Monitor | Python (alertas por email) |
| UI | Kafka UI (provectuslabs) |
| Orquestração | Docker Compose |
Componentes em Detalhe
1. Producer (Python)
O producer é o ponto de entrada dos dados na pipeline. Desenvolvido em Python 3.12 com a biblioteca kafka-python-ng, ele simula transações de um e-commerce de moda com características realistas.
Estrutura do Código
O arquivo producer.py começa com as importações necessárias e a configuração do servidor Kafka:
import json
import uuid
import random
import time
import os
from datetime import datetime, timezone
from kafka import KafkaProducer
KAFKA_SERVER = os.environ.get("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092")
TOPIC = "transactions"
producer = KafkaProducer(
bootstrap_servers=KAFKA_SERVER,
value_serializer=lambda v: json.dumps(v, ensure_ascii=False).encode("utf-8"),
)
Schema da Transação
O producer define um schema estruturado para as transações:
SCHEMA = {
"type": "struct",
"fields": [
{"field": "transaction_id", "type": "string", "optional": False},
{"field": "timestamp", "type": "string", "optional": False},
{"field": "customer_name", "type": "string", "optional": False},
{"field": "customer_email", "type": "string", "optional": False},
{"field": "products", "type": "string", "optional": False},
{"field": "total_amount", "type": "double", "optional": False},
{"field": "payment_method", "type": "string", "optional": False},
{"field": "status", "type": "string", "optional": False},
{"field": "error_type", "type": "string", "optional": True},
{"field": "error_message", "type": "string", "optional": True},
],
"optional": False,
"name": "transaction",
}
Dados Gerados
O producer simula um e-commerce de moda com:
- 10 clientes fictícios:
CUSTOMERS = [
{"name": "Ana Silva", "email": "ana.silva@email.com"},
{"name": "Carlos Oliveira", "email": "carlos.oliveira@email.com"},
{"name": "Mariana Santos", "email": "mariana.santos@email.com"},
{"name": "Pedro Costa", "email": "pedro.costa@email.com"},
{"name": "Juliana Lima", "email": "juliana.lima@email.com"},
{"name": "Rafael Souza", "email": "rafael.souza@email.com"},
{"name": "Fernanda Rocha", "email": "fernanda.rocha@email.com"},
{"name": "Lucas Pereira", "email": "lucas.pereira@email.com"},
{"name": "Beatriz Martins", "email": "beatriz.martins@email.com"},
{"name": "Thiago Barbosa", "email": "thiago.barbosa@email.com"},
]
- 12 categorias de produtos: vestidos, camisas, calças, blazers, tênis, acessórios (bolsas, relógios, óculos), jaquetas, saias, bermudas e macacões
- 4 formas de pagamento:
credit_card,pix,boleto,debit_card - 6 tipos de erro:
payment_declined,insufficient_stock,fraud_suspected,invalid_cvv,timeout,card_expired
Geração de Transações
A função generate_transaction é o coração do producer:
def generate_transaction(is_error=False):
customer = random.choice(CUSTOMERS)
num_items = random.randint(1, 4)
items = random.sample(PRODUCTS, num_items)
products_with_qty = []
for p in items:
qty = random.randint(1, 2)
products_with_qty.append({**p, "quantity": qty})
total = sum(p["price"] * p["quantity"] for p in products_with_qty)
payment = random.choice(PAYMENT_METHODS)
payload = {
"transaction_id": str(uuid.uuid4()),
"timestamp": datetime.now(timezone.utc).isoformat(),
"customer_name": customer["name"],
"customer_email": customer["email"],
"products": json.dumps(products_with_qty, ensure_ascii=False),
"total_amount": round(total, 2),
"payment_method": payment,
}
if is_error:
error = random.choice(ERROR_TYPES)
payload["status"] = "error"
payload["error_type"] = error[0]
payload["error_message"] = error[1]
else:
payload["status"] = "success"
payload["error_type"] = None
payload["error_message"] = None
return {"schema": SCHEMA, "payload": payload}
Loop Principal
O loop principal publica transações continuamente, com ~15% de taxa de erro e intervalo aleatório de 0,5 a 3 segundos entre transações:
def main():
print(f"Producer iniciado. Enviando transações para {KAFKA_SERVER} ...")
while True:
try:
is_error = random.random() < 0.15
tx = generate_transaction(is_error=is_error)
producer.send(TOPIC, tx)
tx_id = tx["payload"]["transaction_id"][:8]
status = tx["payload"]["status"]
amount = tx["payload"]["total_amount"]
customer = tx["payload"]["customer_name"]
print(f"[{status.upper()}] {tx_id} | {customer} | R$ {amount:.2f}")
interval = random.uniform(0.5, 3.0)
time.sleep(interval)
except KeyboardInterrupt:
break
except Exception as e:
print(f"Erro ao enviar: {e}")
time.sleep(5)
producer.flush()
producer.close()

2. Apache Kafka (KRaft Mode)
O Kafka é a espinha dorsal da pipeline. A versão utilizada é o Kafka 4.3 em modo KRaft, que elimina a dependência do Zookeeper, simplificando a operação e melhorando a escalabilidade.
A configuração no docker-compose.yml define um broker Kafka com:
- Cluster ID:
MkU3OEVBNTcwNTJENDM2Qk - Node ID: 1
- Process Roles: broker e controller (KRaft)
- Listeners: PLAINTEXT (9092), PLAINTEXT_INTERNAL (29092), CONTROLLER (9093)
- Healthcheck: verificação via
kafka-broker-api-versions
O Kafka atua como um broker de mensagens distribuído, garantindo durabilidade, escalabilidade, desacoplamento entre produtores e consumidores, e a capacidade de replay de mensagens.
3. PostgreSQL e Schema
O banco de dados PostgreSQL atua como o data warehouse da aplicação. A configuração no Docker Compose expõe o PostgreSQL na porta 5444 para evitar conflitos:
postgres:
image: postgres:16-alpine
environment:
POSTGRES_DB: ecommerce
POSTGRES_USER: reulison
POSTGRES_PASSWORD: reulison123
ports:
- "5444:5432"
volumes:
- postgres-data:/var/lib/postgresql/data
- ./sql/init.sql:/docker-entrypoint-initdb.d/init.sql
Tabela de Transações
O arquivo init.sql cria a tabela principal com índices otimizados:
CREATE TABLE IF NOT EXISTS transactions (
id SERIAL PRIMARY KEY,
transaction_id UUID UNIQUE NOT NULL,
timestamp TEXT NOT NULL,
customer_name VARCHAR(255),
customer_email VARCHAR(255),
products JSONB,
total_amount DECIMAL(10,2),
payment_method VARCHAR(50),
status VARCHAR(20) NOT NULL,
error_type VARCHAR(100),
error_message TEXT,
created_at TIMESTAMPTZ DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS idx_transactions_status ON transactions(status);
CREATE INDEX IF NOT EXISTS idx_transactions_timestamp ON transactions(timestamp);
CREATE INDEX IF NOT EXISTS idx_transactions_error ON transactions(error_type)
WHERE status = 'error';
Destaques do schema:
transaction_idcomo UUID únicoproductscomo JSONB, permitindo consultas flexíveis sobre os itens da compra- Índices em
status,timestampeerror_type(com condição parcial) para consultas eficientes
Views Analíticas
O arquivo views.sql cria views para facilitar a análise no Looker Studio:
v_transacoes_ao_longo_tempo: Visão completa das transações com timestamp convertido para timestamptz.v_transacoes_por_categoria: Agregação por categoria de produto usandojsonb_to_recordset:
CREATE VIEW v_transacoes_por_categoria AS
SELECT
p.category AS categoria,
COUNT(*) AS total_transacoes,
ROUND(SUM((p.price * p.quantity))::numeric, 2) AS valor_total,
ROUND(AVG(p.price)::numeric, 2) AS preco_medio
FROM transactions t
CROSS JOIN LATERAL jsonb_to_recordset(t.products::jsonb)
AS p(sku text, name text, category text, price numeric, quantity int)
GROUP BY p.category
ORDER BY total_transacoes DESC;
v_erros_por_tipo: Análise de erros com cálculo de percentual:
CREATE VIEW v_erros_por_tipo AS
SELECT
error_type,
COUNT(*) AS total_ocorrencias,
ROUND(SUM(total_amount)::numeric, 2) AS valor_impactado,
ROUND(COUNT(*) * 100.0 / SUM(COUNT(*)) OVER (), 2) AS percentual
FROM transactions
WHERE status = 'error'
GROUP BY error_type
ORDER BY total_ocorrencias DESC;
v_transacoes_por_pagamento: Distribuição por meio de pagamento com ticket médio.v_transacoes_por_valor: Faixas de valor para análise de ticket.

4. Kafka Connect com JDBC Sink Connector
O Kafka Connect é a ponte entre o Kafka e o PostgreSQL. O conector JDBC Sink é configurado automaticamente via um container connector-setup:
curl -X POST http://kafka-connect:8083/connectors \
-H "Content-Type: application/json" \
-d '{
"name": "jdbc-sink-tx",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"connection.url": "jdbc:postgresql://postgres:5432/ecommerce",
"connection.user": "reulison",
"connection.password": "reulison123",
"topics": "transactions",
"table.name.format": "transactions",
"insert.mode": "insert",
"pk.mode": "none",
"auto.create": "false",
"auto.evolve": "true",
"value.converter.schemas.enable": "true",
"errors.tolerance": "all"
}
}'
O conector é configurado para:
- Ler mensagens do tópico
transactions - Converter o payload JSON em registros SQL
- Inserir os dados na tabela
transactions - Gerenciar automaticamente a consistência e o offset

5. Consumer de Alertas
Um consumer dedicado monitora continuamente o fluxo de transações em busca de padrões anômalos.
Configuração e Conexão
import json
import os
import smtplib
import ssl
import threading
import time
from datetime import datetime, timezone
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
from kafka import KafkaConsumer
KAFKA_SERVER = os.environ.get("KAFKA_BOOTSTRAP_SERVERS", "localhost:9092")
TOPIC = "transactions"
SMTP_SERVER = os.environ.get("SMTP_SERVER", "smtp.gmail.com")
SMTP_PORT = int(os.environ.get("SMTP_PORT", "587"))
SMTP_USERNAME = os.environ.get("SMTP_USERNAME", "")
SMTP_PASSWORD = os.environ.get("SMTP_PASSWORD", "")
ALERT_EMAIL = os.environ.get("ALERT_EMAIL", "eu@reulison.com.br")
Envio de Emails
A função send_email utiliza SMTP com suporte a TLS/SSL:
def send_email(subject, body):
if not SMTP_USERNAME or not SMTP_PASSWORD:
print(f"[ALERTA] Email não configurado (falta SMTP_USERNAME/SMTP_PASSWORD)")
print(f"[ALERTA] Para: {ALERT_EMAIL} | Assunto: {subject}")
return
msg = MIMEMultipart()
msg["From"] = SMTP_USERNAME
msg["To"] = ALERT_EMAIL
msg["Subject"] = subject
msg.attach(MIMEText(body, "plain", "utf-8"))
try:
context = ssl.create_default_context()
if SMTP_PORT == 465:
with smtplib.SMTP_SSL(SMTP_SERVER, SMTP_PORT, context=context) as server:
server.login(SMTP_USERNAME, SMTP_PASSWORD)
server.send_message(msg)
else:
with smtplib.SMTP(SMTP_SERVER, SMTP_PORT) as server:
server.starttls(context=context)
server.login(SMTP_USERNAME, SMTP_PASSWORD)
server.send_message(msg)
print(f"[EMAIL] Alerta enviado para {ALERT_EMAIL}: {subject}")
except Exception as e:
print(f"[ERRO EMAIL] Falha ao enviar: {e}")
Monitoramento de Inatividade
O heartbeat_monitor verifica a cada 60 segundos se há transações nas últimas 4 horas:
SILENCE_WINDOW = 14400 # 4 horas
def heartbeat_monitor():
global last_transaction_at
while True:
time.sleep(60)
with lock:
elapsed = time.time() - last_transaction_at
if elapsed > SILENCE_WINDOW:
last_seen = datetime.fromtimestamp(
last_transaction_at, tz=timezone.utc
).isoformat()
subject = "[URGENTE] Parada de Transações - E-commerce Moda"
body = (
f"ALERTA DE INATIVIDADE\n\n"
f"Nenhuma transação processada nas últimas 4 horas.\n"
f"Última transação: {last_seen}\n"
f"Tempo sem transações: {elapsed / 3600:.1f} horas\n\n"
f"Verifique o sistema de vendas imediatamente."
)
send_email(subject, body)
Consumo e Detecção de Erros
O loop principal consome mensagens e dispara alertas para transações com erro:
def main():
global last_transaction_at
consumer = KafkaConsumer(
TOPIC,
bootstrap_servers=KAFKA_SERVER,
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
auto_offset_reset="earliest",
group_id="alert-consumer-group",
)
t = threading.Thread(target=heartbeat_monitor, daemon=True)
t.start()
print(f"Consumer de alertas iniciado. Monitorando '{TOPIC}'...")
for msg in consumer:
raw = msg.value
tx = raw.get("payload", raw) if isinstance(raw, dict) else raw
with lock:
last_transaction_at = time.time()
if tx.get("status") == "error":
subject = f"[ERRO] Transação Falhou - {tx.get('error_type', 'desconhecido')}"
body = (
f"Transação com erro detectada!\n\n"
f"ID: {tx.get('transaction_id', 'N/A')}\n"
f"Cliente: {tx.get('customer_name', 'N/A')}\n"
f"Email: {tx.get('customer_email', 'N/A')}\n"
f"Valor: R$ {tx.get('total_amount', 0):.2f}\n"
f"Pagamento: {tx.get('payment_method', 'N/A')}\n"
f"Erro: {tx.get('error_type', 'N/A')}\n"
f"Detalhe: {tx.get('error_message', 'N/A')}\n"
f"Timestamp: {tx.get('timestamp', 'N/A')}\n"
)
send_email(subject, body)

![Captura de tela de um e-mail de alerta com assunto '[ERRO] Transação Falhou - fraud_suspected'. O remetente é eu@reulison.com.br e o destinatário está marcado como 'me'. O corpo do e-mail contém o título 'Transação com erro detectada!' e os detalhes: ID: 82dff3a3-e28b-4c1a-89f4-9337d0cd3772, Cliente: Pedro Costa, Email: pedro.costa@email.com, Valor: R$ 319.80, Pagamento: debit_card, Erro: fraud_suspected, Detalhe: Transação suspeita de fraude bloqueada pelo antifraude, Timestamp: 2026-07-24T22:44:31.987088+00:00. Na parte inferior, botões 'Reply' e 'Forward'. Exemplo de email de alerta recebido](/_next/image?url=%2Fphotos%2Fexemplo-email-alerta.png&w=1120&q=75)
Kafka UI
A interface gráfica do Kafka (provectuslabs/kafka-ui) permite monitorar em tempo real:
- Tópicos e partições
- Produção e consumo de mensagens
- Latência e throughput
- Status dos conectores
Acessível em http://localhost:8080.

Looker Studio
O Looker Studio (antigo Google Data Studio) consome os dados do PostgreSQL para criar dashboards interativos. As views SQL preparadas facilitam a criação de visualizações como:
- Faturamento total e ticket médio
- Distribuição por meio de pagamento
- Evolução temporal das vendas
- Análise de erros e tipos de falha
Orquestração com Docker Compose
Toda a infraestrutura é orquestrada via Docker Compose, o que permite subir todos os serviços com um único comando.
Serviços e Portas
| Serviço | Porta | Descrição |
|---|---|---|
| Kafka | 9092 | Broker Kafka |
| PostgreSQL | 5444 | Banco de dados |
| Kafka Connect | 8083 | REST API do Connect |
| Kafka UI | 8080 | Interface gráfica |
| Producer | - | Gerador de transações |
| Consumer Alert | - | Monitor e alertas |
volumes:
kafka-data:
postgres-data:
connect-data:
networks:
kafka-net:
driver: bridge
Detecção de Fraude em Tempo Real
Um dos cenários mais críticos que uma pipeline como esta pode endereçar é a detecção de fraudes. O projeto inclui simulação de transações com erro fraud_suspected, representando parte dos ~15% de erros no fluxo.
Em sistemas financeiros reais, a detecção de fraudes em tempo real é fundamental. Estudos mostram que arquiteturas baseadas em Kafka permitem latência de subsegundo na detecção de fraudes com fluxos de alta volumetria. O valor disso é imenso: cada minuto de atraso na detecção de uma fraude pode representar perdas significativas para o negócio e danos à reputação da marca.
Como Executar o Projeto
Pré-requisitos
- Docker e Docker Compose instalados
Configuração
cp .env.example .env
# Edite .env com suas credenciais SMTP para receber alertas
Execução
docker-compose up -d --build
Acessos
- Kafka UI: http://localhost:8080
- PostgreSQL: localhost:5444
- Kafka Connect REST: http://localhost:8083
Conectando ao Looker Studio
O banco pode ser exposto via ngrok para conexão com o Looker Studio:
- Host:
tcp://0.tcp.sa.ngrok.io - Porta: (definida pelo ngrok)
- Banco:
ecommerce - Usuário:
reulison

Olá! Quer saber mais?
Por que Esta Arquitetura é Relevante Hoje
O Fim do Batch Processing
Por décadas, o processamento batch foi o padrão: dados eram coletados, armazenados e processados em janelas de tempo (diárias, horárias). Mas em um mundo onde decisões precisam ser tomadas em milissegundos, o batch já não é suficiente.
O stream processing trata os dados não como registros estáticos, mas como fluxos contínuos de eventos. Isso permite:
- Reação imediata: detectar e responder a eventos assim que ocorrem
- Visibilidade total: entender o estado do negócio a qualquer momento
- Tomada de decisão proativa: agir antes que problemas se agravem
O Papel do Kafka como Padrão Industrial
O Apache Kafka se consolidou como o padrão de fato para streaming de dados:
- 38,7% de market share na categoria de filas e mensageria — à frente de RabbitMQ (28,6%) e IBM MQ (6,3%)
- Adotado por mais de 150.000 organizações mundialmente
- Presente em mais de 80% das empresas da Fortune 100
Conclusão
O projeto kafka-transactions-streamline é uma demonstração prática de como construir uma arquitetura de dados moderna, escalável e orientada a eventos. Ele incorpora as melhores práticas da indústria:
- Mensageria robusta com Apache Kafka em modo KRaft
- Integração simplificada com Kafka Connect e JDBC Sink Connector
- Armazenamento confiável com PostgreSQL e índices otimizados
- Monitoramento em tempo real com alertas automáticos por email
- Visualização intuitiva com Looker Studio via views SQL preparadas
Em um mercado onde o streaming analytics cresce a taxas superiores a 30% ao ano e onde 86% dos líderes de TI consideram data streaming um investimento estratégico, dominar essas tecnologias não é apenas relevante — é essencial.
Fonte: Streaming Analytics Market Report 2026
O futuro dos dados é em tempo real. E pipelines como esta são a fundação sobre a qual esse futuro será construído.
Projeto no Github: Kafka Transactions Streamline
Fique ligado
Seja um Expert em Growth
Receba insights práticos sobre marketing, dados, performance e tecnologia direto no seu email.