Streamline Pipeline com Apache Kafka

Por Réulison Silva
Réulison Silva
Published on
Streamline Pipeline com Apache Kafka

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:

Visão Geral da Arquitetura
Visão Geral da Arquitetura

Tecnologias Utilizadas

ComponenteTecnologia
MensageriaApache Kafka 4.3 (KRaft mode, sem Zookeeper)
ConectoresKafka Connect + JDBC Sink Connector
Banco de DadosPostgreSQL 16
ProducerPython 3.12 + kafka-python-ng
MonitorPython (alertas por email)
UIKafka UI (provectuslabs)
OrquestraçãoDocker 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:

Python
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:

Python
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:
Python
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:

Python
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:

Python
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()
Exemplo de saída do producer no terminal mostrando transações sendo enviadas
Exemplo de saída do producer no terminal mostrando transações sendo enviadas

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:

YAML
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:

SQL
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_id como UUID único
  • products como JSONB, permitindo consultas flexíveis sobre os itens da compra
  • Índices em status, timestamp e error_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 usando jsonb_to_recordset:
SQL
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:
SQL
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.
Diagrama do schema da tabela transactions e suas views
Diagrama do schema da tabela transactions e suas views

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:

Prompt
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
Dashboard do Kafka UI mostrando o conector JDBC Sink ativo
Dashboard do Kafka UI mostrando o conector JDBC Sink ativo

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
Python
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:

Python
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:

Python
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:

Python
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)
E-mails chegando na Caixa de entrada
E-mails chegando na Caixa de entrada
Exemplo de email de alerta recebido
Exemplo de email de alerta recebido

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.

Dashboard do Kafka UI mostrando tópicos, métricas e conectores
Dashboard do Kafka UI mostrando tópicos, métricas e conectores

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çoPortaDescrição
Kafka9092Broker Kafka
PostgreSQL5444Banco de dados
Kafka Connect8083REST API do Connect
Kafka UI8080Interface 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
Dashboard final no Looker Studio com as visualizações
Dashboard final no Looker Studio com as visualizações

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.