Mostrando postagens com marcador Kafka. Mostrar todas as postagens
Mostrando postagens com marcador Kafka. Mostrar todas as postagens

sexta-feira, 26 de dezembro de 2025

KAFKA - Principais conceitos

        Nesse artigo vou trazer um resumo sobre o Kafka. Sempre gosto de fazer resumos que podem me ajudar no dia a dia, afinal, é impossível lembrar de tudo o tempo todo.


- Fila, Tópico e Grupos


Apresenta o conceito de fila ou tópico (Pub/Sub):

Kafka -> fila -> Só tem um grupo, o primeiro nó do grupo que ler trata a mensagem.

        '--> tópico -> a mesma mensagem pra vários grupos, onde em cada grupo o primeiro nó que pegar lê e trata a mensagem (Pub/Sub).


Vamos ver o coneito de grupo no exemplo abaixo:


@KafkaListener(

    topics = "pedidos",

    groupId = "pagamento-service"

)

public void consumir(Pedido pedido) {

}


        Se o microsserviço A é do grupo "pagamento-service" e apenas esse grupo acessa o tópico "pedidos", é uma fila. Pode ter várias instancias de A que não impacta nesse comportamento de fila.


        Se crio o microsserviço B e ele tem a seguinte definição de grupo: groupId = "estoque-service"


@KafkaListener(

    topics = "pedidos",

    groupId = "estoque-service"

)

public void consumir(Pedido pedido) {

}


        Então agora temos o comportamento de Pub/Sub (tópico), porque o grupo "estoque-service" também acessa "pedido".


*Note que há um overhead de conceitos aí na questão de tópico, uma coisa é o tópic do Kafka e outra é o conceito de tópico no sentido de Pub/Sub, onde alguém publica e outros são (sub)inscritos para receber a mensagem.


Usando mensagens diferentes no mesmo tópico


        Embora pareça estranho, usar mensagens diferentes no mesmo tópico não é errado. É assim por exemplo que se garante a ordem de execução, usando a mesma key (vamos ver esse coneito mais a frente) e tópico o que joga na mesma partição (vamos ver esse coneito mais a frente), mas o corpo da mensagem é diferente, e cada consumer que tem seu group-id consome o que lhe interessa e descarta o resto. Ignorar evento NÃO é erro.


Exemplo: Podemos ter o Tópico: pedidos e usar a key = pedidoId


Eventos (Cada evento é uma mensagem diferente, com campos diferente):

- PEDIDO_CRIADO

- PEDIDO_PAGO

- PEDIDO_ENVIADO

- PEDIDO_CANCELADO


        Isso casa perfeitamente com Event Sourcing, DDD, CQRS, Auditoria, Reprocessamento...


        Por que usar o MESMO tópico nesse caso? Porque você quer garantir:


✔️ Ordem dos eventos do pedido

✔️ Processamento sequencial

✔️ Consistência de estado

✔️ Um único consumer por pedido

✔️ Reprocessamento confiável



- Partições:


        Partições -> paralelismo/escalabilidade. São divisões dentro de um tópico. Cada partição permite apenas um consumidor por grupo. Se tiver mais consumidores que partições de um mesmo grupo eles ficam ociosos.


        Partição NÃO é algo que você “escala dinamicamente” como pod de Kubernetes. Partição é uma decisão estrutural, onde alterar depois tem efeitos colaterais


        Kafka foi desenhado para escalar consumidores, mas com partições pensadas antes. Ao alterar partições, as partições novas começam vazias, as chaves (keys) podem ir para partições diferentes, a ordem global não é preservada, pode quebrar a lógica baseada em key.


👉 Por isso: erramos para mais, não para menos. Partições devem acompanhar o pico de consumo esperado, não o consumo médio.


- Keys


        As Keys servem para decidir em qual partição a mensagem vai cair. Mesma key significa que a mensagem sempre vai cair na mesma partição. Isso é importante porque o Kafka só garante ordem dentro da partição.


Quando usar key:


✔️ Existe uma entidade de negócio

✔️ Eventos precisam ser processados em ordem

✔️ Existe estado

✔️ Existe atualização incremental


Quando NÃO usar key:


❌ Eventos são independentes

❌ Não existe estado

❌ Ordem não importa

❌ Quer máximo throughput


Keys influenciam na escalabilidade (trade-off real)


Quanto mais granular a key:

  • Mais espalhamento
  • Mais paralelismo
  • Menos ordem global


Quanto mais concentrada a key:


  • Menos paralelismo
  • Mais ordem
  • Possível gargalo


Exemplo de chave ruim: key = "PEDIDOS" -> Tudo cai em uma partição só e o Kafka vira single-thread.


Exemplo de chave boa: pedidoId, userId, contaId, CPF, CNPJ...


        Quando se trabalha com Keys e é necessário aumentar partições, a mesma key pode ir para partições diferentes no momento dessa alteração e a ordem não será garantida. Isso acontece porque o calculo de pra qual partição a mensagem vai é: "hash(key) % N" onde N é o número de partições. Por isso sistemas que dependem de key precisam pensar bem antes de aumentar partições. É preciso fazer uma análise dos riscos.


Exemplo de comando para criar um tópico em Kafka:

./kafka-topics.sh --create --topic <nome_do_topico> --bootstrap-server <endereco_do_broker> --partitions <numero_de_particoes> --replication-factor <fator_de_replicacao>


Exemplo de comando para deletar um tópico em Kafka rodando no Docker:

docker exec -it nome-container /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic batalha


Exemplo de comando para fazer o Kafka ignorar as mensagens anteriores e ir para a última (Kafka rodando no Docker):

docker exec -it  nome-container /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group id-do-grupo --topic topico-aqui --reset-offsets --to-latest --execute


Exemplo de comando para criar mensagens em um tópico Kafka (Kafka rodando no Docker):

docker exec -it nome-container /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic nome-topico


Exemplo de comando para aumentar o tamanho das mensagens suportadas em um tópico Kafka (Kafka rodando no Docker):

docker exec -it nome-container /opt/kafka/bin/kafka-configs.sh --bootstrap-server localhost:9092 --entity-type topics --entity-name nome-do-topico --alter --add-config max.message.bytes=2097152




sábado, 24 de agosto de 2024

Kafka - Rodando no WSL 2 (Ubuntu) e usando (Producer e Consumer) a partir do Windows

        Lá estava eu estudando Kafka e me deparei com a seguinte situação: O Kafka instalado no Ubuntu através do WSL 2, rodando legal e eu criando produtores e consumidores a partir do Java usando o Intellij no Windows. Daí o que ocorre é: o Kafka sobe normal, mas o projeto Java não consegue encontrar o serviço, mesmo colocando no projeto o IP do WSL.

      Então descobri que para fazer o projeto Java se comunicar com o Kafka no WSL 2 era preciso configurar a propriedade advertised.listeners no arquivo server.properties no diretório config do Kafka: 


advertised.listeners=PLAINTEXT://ip_do_wsl:9092


        Dessa forma a aplicação mesmo em outra rede consegue encontrar o Kafka e criar e ler os tópicos, produzindo e consumindo suas menagens. A propriedade advertised.listeners no Apache Kafka desempenha um papel crucial na comunicação entre os clientes (produtores e consumidores) e os brokers do Kafka, especialmente em cenários de redes distribuídas ou ambientes com várias interfaces de rede.

        Ela especifica os endereços IP ou nomes de host e portas que os brokers Kafka devem anunciar para os clientes (produtores, consumidores, etc.). Esses são os endereços que os clientes usarão para se conectar ao broker após receberem a metadata inicial do Kafka.

        Uma questão é: o IP do WSL pode mudar toda vez que for iniciado. Ficar mudando isso tanto no arquivo de configuração do Kafka quanto na aplicação (ou aplicações como é o caso microsseviços) é chato. Então pra mudar na aplicação (no meu caso Java) muda muito de caso pra caso, mas a ideia inicial seria usar variável de ambiente ou arquivo de propriedades tirando do código a dependência. 

     Para o Kafka recorri ao GPT que gerou o script bash a seguir que troca o IP no arquivo config/server.properties:


#!/bin/bash

# Passo 1: Obter o IP do WSL 2

WSL_IP=$(hostname -I | awk '{print $1}')


# Passo 2: Caminho para o arquivo server.properties do Kafka

KAFKA_CONFIG_PATH="/caminho/para/seu/kafka/config/server.properties"


# Passo 3: Atualizar o advertised.listeners com o IP atual do WSL 2

# Verifique se o arquivo existe

if [ -f "$KAFKA_CONFIG_PATH" ]; then

    # Usar sed para substituir o advertised.listeners existente ou adicionar se não existir

    if grep -q "^advertised.listeners=" "$KAFKA_CONFIG_PATH"; then

        # Substitui a linha existente

        sed -i "s/^advertised.listeners=.*/advertised.listeners=PLAINTEXT:\/\/$WSL_IP:9092/" "$KAFKA_CONFIG_PATH"

    else

        # Adiciona a linha se não existir

        echo "advertised.listeners=PLAINTEXT://$WSL_IP:9092" >> "$KAFKA_CONFIG_PATH"

    fi

    echo "advertised.listeners atualizado com o IP: $WSL_IP"

else

    echo "Erro: Arquivo server.properties não encontrado em $KAFKA_CONFIG_PATH"

    exit 1

fi


# Passo 4: Iniciar ou reiniciar o Kafka para aplicar as mudanças

# Ajuste o comando abaixo para reiniciar seu Kafka

# Exemplo para iniciar Kafka:


/opt/kafka/bin/kafka-server-stop.sh


while ps ax | grep -i 'kafka\.Kafka' | grep -v grep > /dev/null; do

    echo "Aguardando o Kafka encerrar..."

    sleep 1

done


/opt/kafka/bin/kafka-server-start.sh -daemon $KAFKA_CONFIG_PATH


echo "Script concluído. Certifique-se de reiniciar o Kafka para aplicar as mudanças."


       Basta salvar o script como atualiza_kafka.sh na pasta do Kafka, dar permissão de execução com chmod +x atualiza_kafka.sh e executar antes de subir o Kafka.

        Ainda não executei o Kafka pelo Docker subindo também os produtores e consumidores em uma mesma rede com Docker Compose então não sei se muda muita coisa, mas como seria a mesma rede, acredito que essa configuração de advertised.listeners não seja necessária, de qualquer forma esse artigo pode ser útil pra mim mesmo no futuro e outros que estão aprendendo a ferramenta.