Skip to content

Event-Driven Architecture (EDA)

O que é EDA

Definição: Event-Driven Architecture (EDA)

Padrão de arquitetura de software voltado ao design de aplicações que se comunicam entre si por meio de eventos, de forma assíncrona — permitindo que sistemas continuem funcionando mesmo que algumas partes estejam temporariamente indisponíveis. Em vez de depender de um fluxo de controle rígido e bloqueante (uma chamada síncrona esperando resposta), as aplicações reagem aos eventos que ocorrem no sistema, o que promove maior flexibilidade. Ganhou popularidade a partir dos anos 2000, junto com o crescimento de sistemas distribuídos e aplicações em nuvem.

Definição: Componentes da EDA

  • Evento — uma mudança de estado ou uma ação que ocorre num sistema, gerada por entradas internas ou externas (uma alteração num banco de dados, uma operação realizada por um usuário).
  • Produtores (producers) — componentes que geram eventos (ex.: um serviço de pagamento que emite um evento quando uma transação é finalizada).
  • Consumidores (consumers) — componentes que reagem a eventos (ex.: um serviço que atualiza o status do pedido ao receber um evento de pagamento aprovado).
  • Bus de Eventos (Event Bus) — o sistema de mensagens que transporta os eventos dos produtores aos consumidores; pode ser um sistema de streaming, como o Apache Kafka, ou uma fila mais simples, como o RabbitMQ.

O maior benefício da EDA é promover baixo acoplamento entre os componentes do sistema: produtores não precisam saber quais consumidores estão ouvindo, nem depender deles para concluir suas próprias operações — e vice-versa. Essa independência torna o sistema mais resiliente, escalável e fácil de evoluir.

Categoria Ponto Descrição
Vantagem Desacoplamento Componentes independentes, facilitando manutenção e escalabilidade
Vantagem Escalabilidade Consumidores e produtores podem ser adicionados/removidos facilmente conforme demanda
Vantagem Resiliência Permite operação parcial mesmo com falhas, com reprocessamento de eventos
Vantagem Flexibilidade Novos recursos podem ser adicionados sem afetar o sistema existente
Desvantagem Complexidade Maior esforço para gerenciar eventos e orquestrar fluxos
Desvantagem Difícil depuração Problemas são mais difíceis de rastrear devido ao fluxo distribuído
Desvantagem Consistência eventual Dados podem ficar temporariamente inconsistentes em alguns cenários
Desvantagem Dependência de infraestrutura Necessita ferramentas robustas de mensageria, gerando custos extras

Padrões de EDA

Dois padrões principais organizam como os eventos fluem entre produtores e consumidores.

Definição: Event Sourcing (transmissão de eventos)

O estado de um sistema é armazenado como uma sequência de eventos, registrados numa espécie de log, em vez de apenas atualizar e guardar o estado atual dos dados. Cada alteração é registrada como um evento separado, formando um histórico completo das mudanças ocorridas ao longo do tempo — valioso para entender o que motivou o estado atual do sistema, além de útil para auditoria, depuração e até para prever tendências futuras com base em ações já realizadas. Os consumidores não se inscrevem num fluxo fixo e constante — podem ler o log a partir de qualquer ponto e a qualquer momento, ingressando nele sob demanda.

flowchart LR
    Client --> S1[Service 1] --> DB1[(Database 1)]
    S1 <--> ES[(Event Store)]
    ES <--> S2[Service 2] --> DB2[(Database 2)]

Definição: Pub/Sub (Publish/Subscribe ou Publicar/Assinar)

Os produtores publicam mensagens num canal, enquanto os consumidores se inscrevem para recebê-las — infraestrutura de mensageria baseada na assinatura de fluxos de eventos. Sempre que um evento ocorre (é publicado), ele é enviado a todos os consumidores inscritos que precisam ser notificados para realizar os processamentos correspondentes.

flowchart LR
    P1[Producer 1] --> EC["Event<br/>channels"]
    P2[Producer 2] --> EC
    P3[Producer 3] --> EC
    EC --> C1[Consumer 1]
    EC --> C2[Consumer 2]
    EC --> C3[Consumer 3]

Pensar em eventos, não em comandos

Na comunicação tradicional (request-response), uma aplicação precisa conhecer o endpoint de outra e mandar um comando ("processe este pagamento"), o que cria acoplamento. Na EDA, a aplicação anuncia um fato — "o pagamento foi aprovado" — e quem se interessa reage. Muda a forma de pensar o negócio: em vez de "o comando Y deve ser executado", "o evento X ocorreu". Em uma academia, "aluno entrou" pode gerar "notificar a série de exercícios" e "avisar o professor"; em uma companhia aérea, "voo atrasado" pode gerar remarcação e aviso aos passageiros. São eventos de negócio e as oportunidades que eles abrem — o EDA permite explorá-los em tempo real.

Definição: comunicação síncrona x assíncrona; consistência forte x eventual

Síncrona: o cliente espera a resposta antes de continuar. Assíncrona: envia e segue, e a resposta (se houver) chega depois. Consistência forte: após uma atualização, todos os leitores veem o dado novo (obrigatória em saldo bancário, por exemplo). Consistência eventual: os leitores podem ver o dado antigo por um tempo, até convergirem. EDA é distribuída, assíncrona e eventualmente consistente: não use EDA onde a consistência forte é indispensável.

Benefícios: baixo acoplamento, escalabilidade (cada aplicação escala sozinha), extensibilidade (novos consumidores sem alterar os existentes), disponibilidade (o broker retém eventos enquanto um consumidor está fora) e tempo real. Desafios: idempotência, ordem, duplicidade, depuração de fluxos distribuídos, governança de esquemas e consistência eventual.

Eventos e mensagens: anatomia

Conceito Explicação
Produtor, consumidor, broker Quem publica (também publisher, fonte), quem reage (subscriber) e o intermediário (message broker, event bus, roteador, hub) que roteia, traduz, persiste e entrega. Uma aplicação pode ser as duas coisas
Mensagem Termo geral para o que as aplicações trocam; tem a mesma estrutura seja evento, comando ou consulta
Evento discreto Fato independente que relata uma mudança de estado: "pedido criado", "preço atualizado"
Evento de série (event stream) Fluxo contínuo e ordenado de eventos medidos no tempo: leituras de sensores, métricas, cliques
Comando Pedido para executar uma ação ou alterar um estado; a resposta é opcional
Query Pedido para recuperar informação; exige resposta (em REST, GET); pode ser assíncrona via request-reply
Cabeçalho (header) Metadados: identificador, tipo, origem, data/hora, chave de correlação; o broker os usa para roteamento e rastreamento
Corpo (body) O dado transmitido (pedido, leitura, aluno)
Formato e esquema Formato = estrutura de nomes e tipos; esquema = contrato que valida campos, tipos e obrigatoriedade

Formato de texto x binário: texto (JSON, XML, CSV) é legível e universal, mas maior (carrega os nomes dos campos); binário (Avro, Protobuf) é compacto e rápido, mas exige serialização e esquema. Recomendação: comece com JSON (ou XML) se atender aos requisitos de desempenho; use binário quando volume e latência pedirem; para APIs externas a parceiros, prefira texto. O que importa é o formato ser padrão de mercado, independente de linguagem e com bom suporte em bibliotecas (detalhes em Schema Registry, Avro e Protobuf).

Protocolos: AMQP (RabbitMQ; troca de mensagens com exchanges, filas e bindings), MQTT (leve, para IoT), HTTP/HTTPS (webhooks, APIs), WebSocket (bidirecional em tempo real), e o protocolo próprio do Kafka.

Destino (canal): uma fila entrega cada mensagem a no máximo um consumidor (competição); um tópico entrega a todos os consumidores inscritos. Um consumidor durável recebe as mensagens mesmo que estivesse offline quando chegaram (graças à persistência); um não durável só recebe se estiver conectado.

Garantia de entrega: o broker persiste o evento e exige confirmação (ack) em dois momentos: do broker ao produtor (recebimento) e do consumidor ao broker (processamento). Semânticas no máximo uma vez, ao menos uma vez e exatamente uma vez estão em garantias de entrega.

Padrões de comunicação de mudança de estado

Padrão Como funciona Quando usar
Notificação de evento (event notification) O evento leva o mínimo (um ID e o tipo); quem precisa de mais detalhes consulta o produtor O mais comum; baixo acoplamento de dados, mas gera chamadas de volta
Transferência de estado no evento (event-carried state transfer) O evento carrega todo o estado relevante da entidade Consumidor autônomo, sem chamar o produtor; pré-requisito do event sourcing; eventos maiores
Claim check O evento carrega só uma referência a um dado grande guardado num serviço externo (banco, armazenamento de objetos) Anexos, imagens, mensagens que excedem o limite do broker
Event sourcing Guarda-se a sequência de eventos, e o estado é reconstruído a partir deles Auditoria, histórico, replay (CQRS e Event Sourcing)
CQRS Separa escrita e leitura em modelos distintos, sincronizados por eventos Leituras e escritas com necessidades diferentes
Saga Transação distribuída como sequência de transações locais com ações de compensação Processos de negócio que passam por vários serviços
Outbox / CDC Gravação atômica do evento junto do dado; publicação confiável Evitar o "dual write" (Mensageria confiável)
DLQ, FIFO, webhook Fila de mensagens com falha; ordem de entrega; notificação por HTTP Veja as seções sobre resiliência e integração

Saga: orquestrada x coreografada.

Orquestrada Coreografada
Controle Um orquestrador central comanda os passos (ex.: AWS Step Functions) Cada participante executa sua transação local e publica um evento que dispara o próximo
Vantagens Visão clara da sequência; fácil de entender e monitorar Desacoplada, sem ponto único de falha, alinhada à EDA
Desvantagens O orquestrador é um ponto de falha e de acoplamento Fluxo difícil de enxergar quando há muitos participantes
Indicada para Processos complexos, com muitos passos Processos simples, poucos participantes

Em ambas é preciso desenhar o fluxo de compensação desde o início (estornar o pagamento, devolver o estoque, cancelar o pedido). Perguntas de projeto: o ganho de paralelismo no caminho feliz compensa o custo da compensação? Quais eventos de erro são genéricos (um por serviço, com o motivo no corpo) e quais específicos (um por tipo de falha, que permite assinaturas diretas)? O resultado, "the big picture", é o desenho da coreografia, com diagramas de sequência e a documentação dos eventos.

Modelando e documentando eventos

  • EventStorming: workshop colaborativo (criado por Alberto Brandolini) em que desenvolvimento, especialistas do domínio e arquitetura mapeiam o negócio no quadro. Primeiro os eventos de domínio (post-its laranja, no passado: "Pedido criado"), em ordem temporal; depois comandos/gatilhos (o que provoca cada evento), agregados, políticas ("sempre que X, faça Y") e, por fim, contextos delimitados e o mapa de contexto. Acelera o aprendizado do negócio, alinha a linguagem ubíqua e é o ponto de partida do DDD (DDD). Dicas: facilitador ativo, todos na sala, sessões curtas e várias.
  • AsyncAPI: especificação (JSON/YAML) para descrever interfaces assíncronas — canais, mensagens, esquemas e servidores —, o "OpenAPI dos eventos". Gera documentação e código, valida e alimenta um portal do desenvolvedor. Organize por aplicação (o que ela publica e assina) ou por domínio de negócio (melhor para expor a outros domínios).
  • CloudEvents: especificação (CNCF) que padroniza os metadados de um evento (id, source, specversion, type, time, datacontenttype, subject...), para que eventos de sistemas e provedores diferentes sejam interoperáveis. Convive com a AsyncAPI (uma descreve a interface; a outra, o envelope).

O broker de eventos

O broker é o coração da EDA: recebe, valida, roteia, persiste e entrega. Ao escolher um, avalie: tipo, protocolos e SDKs, entrega, retenção, desempenho, operação e governança (monitoramento, segurança, multi-tenant), gerenciamento de esquemas, padrão de implantação (autogerenciado, nuvem gerenciada ou serverless) e custo.

Tipos de broker

Tipo Característica Exemplos Indicado para
Orientado a fila A mensagem é removida após o ack; consumidores competem; push; roteamento flexível (exchanges), prioridade, DLQ, FIFO RabbitMQ, ActiveMQ, Amazon SQS Distribuição de tarefas, publish/subscribe com filas por assinante, integração entre aplicações
Orientado a log Eventos ficam retidos no log (tópico dividido em partições, replicadas); consumidores leem por offset (pull) e podem reler; ordem por partição (chave) Apache Kafka, Redpanda, Amazon Kinesis, Apache Pulsar Streaming, grande volume, replay, event sourcing, vários consumidores independentes
Orientado a assinatura Distribui por regras de filtro, normalmente push (a webhooks, funções) e remove após a confirmação; retenção curta Amazon EventBridge, Google Pub/Sub, Azure Event Grid Integrações em nuvem, eventos de serviços, aplicações serverless

Funcionalidades comuns

Funcionalidade O que faz
Push x *pull* Push: o broker "empurra" o evento (endpoint HTTP, função, serviço) — bom para webhooks e consumidores nativos de nuvem. Pull: o consumidor "puxa" no seu ritmo, com conexão persistente — melhor para alto volume e streaming
SDK Bibliotecas por linguagem que escondem o protocolo (e permitem ajustar timeouts, tentativas, confirmação)
Lote (batch) Publicar/consumir vários eventos por viagem (batch.size, linger.ms): mais vazão, com pequeno atraso; prefetch limita o que o consumidor guarda em buffer
Roteamento inteligente Entrega só aos consumidores certos por filtro (nome do canal, cabeçalho ou conteúdo): evita descartar eventos e dispensa um tópico por tipo
Confirmação (ack) Automática (ao receber) ou manual (após processar) — define a garantia contra perda
Retenção Por quanto tempo o evento fica disponível (de segundos a indefinido)
Reprodução (replay) Reler eventos passados a partir de um offset ou timestamp (recuperar falhas, criar novas visões)
Visibilidade Tracejamento e métricas (latência, atraso do consumidor, tamanho de fila) para operar a solução
Agendamento Entrega atrasada ou programada
Gerenciamento de esquema Registro e validação de contratos (Schema Registry)

EDA e outros estilos

  • Microsserviços: a comunicação síncrona entre serviços cria acoplamento forte e falhas em cascata (se um serviço cai, quem o chama também falha). Com EDA, os serviços se comunicam pelo broker, que retém os eventos enquanto o consumidor está fora. Na prática convivem os dois estilos: síncrono (request-response) para consultas e interações que exigem resposta imediata; assíncrono para o fluxo de negócio e a propagação de mudanças (Microsserviços).
  • Serverless: funções e serviços gerenciados acionados por eventos, sem servidores para provisionar; broker, banco e API também podem ser sem servidor. Elasticidade rápida e pagamento por uso: combina naturalmente com EDA (Nuvem).
  • Streaming de dados: em vez de lotes noturnos de ETL ("D+1"), os dados fluem como eventos de série e são processados em tempo real (ingestão, armazenamento, processamento, envio). Plataformas como Kafka, Kinesis e Flink, e o Kafka Streams implementam essa arquitetura.
  • Plataforma de integração: um gerenciador de APIs (gateway, segurança, portal) cobre o mundo request-response, e o broker cobre o assíncrono; uma plataforma de integração (iPaaS) reúne conectores, orquestração, API e eventos, e atende integrações com sistemas legados, parceiros e SaaS.

Caso de estudo: plataforma de e-commerce

Um roteiro de projeto de EDA (aplicável em entrevistas de system design):

  1. Entender o problema e o escopo: uma loja que não suporta o pico de pedidos perde receita e credibilidade. Mapear o processo de negócio (pedido → reserva de estoque → pagamento → preparação → envio) e o fluxo de erro (central de operações).
  2. Requisitos: funcionais, não funcionais (volume de pedidos, latência, disponibilidade) e restrições (por exemplo, implantação on-premises por exigência legal).
  3. Decisões e ADR: escolher o broker e o estilo (microsserviços + EDA), registrando o racional em um ADR (Architectural Decision Record): contexto, alternativas, decisão e consequências (Boas práticas).
  4. Design: coreografia de eventos, com Saga e fluxo de compensação; AsyncAPI/CloudEvents; diagramas ("big picture" e de sequência).
  5. Teste: cenário do caminho feliz em BDD, automatizado (Qualidade).
  6. Implementação e desafios adicionais (idempotência, DLQ, observabilidade, replay).

Apache Kafka

Definição: Apache Kafka

Plataforma de streaming distribuída com o objetivo de mover, armazenar e direcionar dados entre sistemas em tempo real, garantindo alta performance e resiliência. Originalmente desenvolvido pelo LinkedIn, hoje é um projeto de código aberto mantido pela Apache Software Foundation. De forma distribuída, processa uma vasta quantidade de dados e os entrega em tempo real, trabalhando tanto com padrões de fila quanto de pub/sub, além de atuar como um banco de dados ao persistir as mensagens geradas em disco — com performance equivalente ao processamento diretamente em memória, diferenciando-se de sistemas tradicionais de filas como o RabbitMQ.

Componentes do Apache Kafka

Definição: Mensagem, Tópico e Offset

Uma mensagem é o evento na EDA — uma unidade de dados composta por uma chave, um valor e um timestamp; a chave pode direcionar a mensagem a uma partição específica de um tópico, enquanto o valor contém o conteúdo real. Um tópico é um pipeline de dados que orquestra a persistência e a entrega das mensagens aos consumidores, dividido em partições (numeradas a partir de 0, definidas na criação do tópico) para permitir alta performance e escalabilidade. A cada mensagem armazenada numa partição é atribuído um offset — a posição da mensagem naquela partição — e cada consumidor pode estar lendo mensagens num offset diferente do outro, processando de forma independente.

flowchart LR
    Producers --> T["Tópico<br/>(partições 0, 1, 2, ...)"]
    T --> C1["Consumer 1<br/>(offset=4)"]
    T --> C2["Consumer 2<br/>(offset=6)"]

Definição: Broker e Cluster

Um broker é um servidor responsável por armazenar mensagens e atender às solicitações de leitura e escrita dos clientes (produtores e consumidores) — cada broker é identificado por um ID único e armazena dados de uma ou mais partições. Um cluster é um conjunto de brokers que trabalham juntos, distribuindo as partições de um tópico entre si — o que reforça a natureza distribuída do Kafka e aumenta a resiliência: se um broker fica indisponível, nem todas as mensagens do tópico são perdidas, só as das partições que ele hospedava.

Definição: Apache Zookeeper

Responsável pela descoberta dos brokers e pela orquestração do gerenciamento do cluster — coordenação, gerenciamento de configuração e sincronização entre brokers. Sem o Zookeeper, o Kafka não conseguiria operar de forma distribuída e confiável em ambientes de produção.

Grupos de consumidores (Consumer Groups)

Definição: Consumer Group

Por padrão, quando vários consumidores estão inscritos no mesmo tópico, cada mensagem é entregue a apenas um deles — o que inviabiliza casos em que dois ou mais consumidores diferentes precisam receber uma cópia da mesma mensagem (ex.: um pagamento realizado interessando tanto a um serviço de processamento quanto a um de auditoria). Um Consumer Group resolve isso: o Kafka garante que cada grupo receba uma cópia da mensagem — dentro do mesmo grupo, ela continua sendo distribuída entre os consumidores membros, nunca duplicada para dois consumidores do mesmo grupo.

flowchart LR
    P[Producer] --> T[Tópico<br/>3 partições]
    T --> CG0["Consumer Group 0<br/>(audit-service)"]
    T --> CG1["Consumer Group 1<br/>(order-service x2)"]

Definição: Regra de ouro — nunca mais consumidores que partições, no mesmo grupo

Dentro de um mesmo grupo, uma partição não pode ser lida por mais de um consumidor ao mesmo tempo — se houver mais consumidores num grupo do que partições no tópico, o consumidor excedente fica ocioso, sem nenhuma mensagem para processar. Por padrão, se nenhum grupo for definido, o Kafka cria automaticamente um grupo próprio para cada consumidor, garantindo que todos recebam ao menos uma cópia da mensagem. O Kafka, por padrão, também não garante ordem de entrega no nível do tópico — só assegura ordem dentro de cada partição individualmente.

Resiliência: replication factor

Definição: Replication factor (fator de replicação)

Define quantas réplicas cada partição de um tópico terá, obrigatoriamente distribuídas em brokers diferentes — com replication factor = 1 (padrão), cada partição existe só num broker, e a indisponibilidade dele significa perda completa das mensagens daquela partição. Com um fator maior, o Kafka define qual réplica é a "master" (de onde os consumidores leem) e mantém as demais como cópias — se o broker da master ficar indisponível, uma das réplicas é automaticamente promovida a master, sem perda de dados. O valor deve ser definido considerando a criticidade do sistema e o número de brokers disponíveis. Consumer Groups também são resilientes a indisponibilidade: quando o problema é corrigido, o consumo retoma a partir do último offset confirmado antes da falha.

Principais comandos da CLI

Comando Descrição
kafka-topics.sh --bootstrap-server <host> --list Lista todos os tópicos disponíveis
kafka-topics.sh --bootstrap-server <host> --create --topic <nome> --partitions N --replication-factor N Cria um novo tópico
kafka-topics.sh --bootstrap-server <host> --delete --topic <nome> Deleta um tópico existente
kafka-console-producer.sh --topic <nome> --bootstrap-server <host> Inicia um produtor de mensagens interativo
kafka-console-consumer.sh --topic <nome> --from-beginning --bootstrap-server <host> Inicia um consumidor, lendo desde o início
kafka-configs.sh --describe --entity-type topics --entity-name <nome> --bootstrap-server <host> Descreve as configurações de um tópico
kafka-consumer-groups.sh --list --bootstrap-server <host> Lista todos os grupos de consumidores
kafka-consumer-groups.sh --describe --group <nome> --bootstrap-server <host> Descreve informações de um grupo específico

Apache Kafka com Spring Boot

Integrar uma aplicação Spring Boot ao Kafka exige configurar um produtor (ProducerFactory - KafkaTemplate) e/ou um consumidor (ConsumerFactory + @KafkaListener).

@Configuration
public class KafkaProducerConfig {

    @Value(value = "${spring.kafka.bootstrap-servers}")
    private String bootstrapAddress;

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}
@Component
@RequiredArgsConstructor
public class KafkaMovieProducerDataProviderAdapter implements MovieGateway {
    private final KafkaTemplate<String, String> kafkaTemplate;
    private final ObjectMapper objectMapper;
    private static final String moviesTopic = "movies-topic";

    @Override
    public Integer create(Movie toCreate) {
        MovieMessage movieMessage = new MovieMessage(toCreate.getName(),
            toCreate.getGenre().toString(), toCreate.getAvailableTotal());
        try {
            String messageAsJson = objectMapper.writeValueAsString(movieMessage);
            kafkaTemplate.send(moviesTopic, messageAsJson);
        } catch (JsonProcessingException e) {
            LOGGER.error("Falha ao converter mensagem: {}", toCreate, e);
        }
        return new Random().nextInt(); // simula um ID gerado
    }
}

kafkaTemplate.send(topico, mensagem) publica a mensagem no tópico indicado — aqui, convertida para String em formato JSON antes do envio.

@EnableKafka
@Configuration
public class KafkaConsumerConfig {

    @Value(value = "${spring.kafka.bootstrap-servers}")
    private String bootstrapAddress;
    @Value(value = "${customer.marketing.consumer.group.id}")
    private String groupId;

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, MovieDeserializer.class);
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}
@Service
@RequiredArgsConstructor
public class MovieCreatedKafkaConsumer {
    private final SendCommunicationToCustomerWhenMovieCreatedInputBoundary inputBoundary;
    private final IMovieMapper mapper;

    @KafkaListener(topics = "${movies.topic.name}",
                   groupId = "${customer.marketing.consumer.group.id}")
    public void receive(CreateMovieDto message) {
        LOGGER.info("Received Message in group catalog-consumer-group: " + message);
        inputBoundary.execute(mapper.movieCreateDtoToMovie(message));
    }
}

Definição: @KafkaListener é um entrypoint, não um REST controller

Uma reflexão importante sobre arquitetura: entrypoints não precisam ser sempre REST controllers — na Clean Architecture, um consumidor Kafka (via @KafkaListener) é tão entrypoint quanto um endpoint HTTP. Entrypoints devem ser vistos como portas de entrada para a aplicação, podendo ser REST controllers, rotinas automatizadas, consumidores de fila, prompts de IA, ou qualquer outro mecanismo de acionamento — e podem ser trocados livremente sem que o domínio de negócio por trás precise mudar.

Integrando com sistemas externos de forma assíncrona

Um cenário muito comum em EDA: processar um pagamento (ou qualquer outra operação que dependa de um sistema externo, como um gateway de pagamento — PayPal, Cielo, Mercado Pago, PagSeguro) sem travar o fluxo principal da aplicação esperando a resposta.

Definição: Requisição síncrona x assíncrona a um sistema externo

Uma integração síncrona submete o pedido e espera a resposta na mesma requisição — mais simples de implementar, mas o tempo de resposta da aplicação requisitante fica refém da demora (às vezes de segundos) do sistema externo. Integrações assíncronas — via mensageria (postar a requisição num tópico) ou webhooks (o sistema externo chama de volta a aplicação quando o processamento terminar) — evitam esse acoplamento de tempo de resposta, ao custo de uma orquestração mais complexa (rastrear o estado da operação até a confirmação chegar).

Definição: Desacoplamento aparente — o gargalo escondido

Publicar uma mensagem num tópico Kafka para processar algo de forma assíncrona não garante, sozinho, que o sistema esteja livre de gargalos: se o consumidor dessa mensagem, ao processá-la, faz uma chamada síncrona para um sistema externo lento, a thread daquele consumidor fica bloqueada esperando a resposta — o gargalo não desapareceu, só migrou de lugar (do cliente original para o próprio consumidor). Uma solução mais robusta usa webhooks (ou uma interface assíncrona do próprio sistema externo) para que o sistema externo notifique a conclusão, publicando então uma mensagem nova no tópico de origem — só nesse ponto o fluxo fica de fato ponta-a-ponta assíncrono.

Definição: Idempotência em consumidores de eventos

Como filas e tópicos podem, em cenários de falha, entregar a mesma mensagem mais de uma vez (reprocessamento após um consumidor cair antes de confirmar o processamento, por exemplo), o consumidor deve ser idempotente — processar a mesma mensagem duas vezes não pode gerar um efeito colateral duplicado (cobrar um pagamento duas vezes, criar dois pedidos idênticos). Isso normalmente é garantido verificando, antes de processar, se aquela operação (identificada por um ID de mensagem ou de pedido) já foi concluída anteriormente.

Kafka avançado: garantias de entrega e operação

Complementa Apache Kafka e grupos de consumidores.

Partições, chaves e ordem

  • A ordem só é garantida dentro de uma partição. Mensagens com a mesma chave (ex.: pedidoId) vão sempre para a mesma partição — use a chave de negócio quando a ordem por entidade importar.
  • O número de partições define o paralelismo máximo do grupo: com 6 partições, no máximo 6 consumidores do mesmo grupo trabalham ao mesmo tempo (os demais ficam ociosos). Aumentar partições depois muda o mapeamento chave → partição, então planeje com folga.
  • Rebalanceamento: quando um consumidor entra ou sai do grupo, as partições são redistribuídas e o consumo pausa por instantes. Reduza o impacto com static membership, processamento rápido (max.poll.interval.ms) e atribuição cooperativa.

Confirmação de offset e semânticas de entrega

Semântica Como ocorre Risco
At-most-once (no máximo uma) Confirma o offset antes de processar Pode perder mensagens se cair no meio
At-least-once (pelo menos uma) Confirma depois de processar (padrão recomendado) Pode duplicar → consumidor idempotente
Exactly-once (exatamente uma) Produtor idempotente + transações do Kafka (enable.idempotence=true, transactional.id) Maior custo; vale dentro do ecossistema Kafka (consumir → processar → produzir)
# Produtor confiável
acks=all                      # espera todas as réplicas em sincronia
enable.idempotence=true       # evita duplicatas por retentativa do produtor
retries=2147483647
# Consumidor
enable.auto.commit=false      # confirma manualmente após processar
isolation.level=read_committed

Definição: acks e ISR

acks define quantas réplicas confirmam a gravação: 0 (nenhuma), 1 (só o líder) ou all (todas as réplicas em sincronia, o ISR — In-Sync Replicas). Com acks=all e min.insync.replicas=2, a gravação sobrevive à queda de um broker.

Retenção, compactação e schemas

  • Retenção: por padrão o Kafka guarda as mensagens por tempo (retention.ms) ou tamanho (retention.bytes), independentemente de terem sido consumidas — por isso novos consumidores podem reler o histórico (replay).
  • Compactação (log compaction): mantém só a última mensagem de cada chave — ideal para guardar o "estado atual" (ex.: cadastro de clientes).
  • Schema Registry: serviço que guarda os esquemas (Avro, Protobuf ou JSON Schema) dos eventos e valida a compatibilidade entre versões (backward, forward, full), evitando que um produtor quebre os consumidores.
  • Kafka Streams / ksqlDB: processamento contínuo de fluxos (filtros, agregações, joins, janelas de tempo) direto sobre os tópicos.

Contratos de mensagens: Schema Registry, Avro e Protobuf

O Kafka aceita qualquer sequência de bytes: nada impede que um produtor publique uma mensagem sem campos, com tipos trocados ou um formato novo. Enquanto o sistema é pequeno, "combinar na conversa" funciona; com o tempo, cada mudança em um produtor pode quebrar silenciosamente vários consumidores. Falta um contrato — o papel que, em uma API REST, é da própria API (e, em um banco, das constraints e dos tipos).

Definição: Schema Registry

Serviço central que guarda e versiona os esquemas (schemas) das mensagens e valida produtores e consumidores contra eles. O produtor serializa a mensagem com um schema registrado e inclui na mensagem apenas o ID do schema (poucos bytes); o consumidor usa esse ID para buscar o schema e desserializar. Também recusa evoluções incompatíveis de schema.

Problemas que o contrato evita

  • Campos obrigatórios ausentes: uma mensagem sem nome do item de cardápio é publicada normalmente, e o consumidor recebe null.
  • Valores padrão silenciosos: um int Java não preenchido vira 0 e passa por um id válido de "restaurante zero".
  • Mudança de tipo ou renomeação de campo sem avisar quem consome.

Validar só no controller do produtor não basta: outros produtores e outros caminhos de escrita escapam. A validação precisa estar no tópico.

Conceitos do Schema Registry

Conceito Significado
Subject Nome lógico que agrupa e versiona os schemas de um tópico — por convenção <tópico>-value e <tópico>-key (a chave e o valor têm schemas independentes)
Schema ID Identificador global e único de um schema registrado (é o que vai na mensagem)
Versão Número sequencial de um schema dentro do subject; não confundir com o Schema ID
Compatibilidade Política por subject que define quais evoluções são permitidas
API REST Cadastro, consulta, listagem de subjects/versões, verificação de compatibilidade e exclusão (soft delete ou permanente) por HTTP

Formatos aceitos: Avro, Protobuf e JSON Schema. Com o JSON puro, os serializadores comuns (JsonSerializer) não validam nada; os serializadores do Schema Registry (KafkaJsonSchemaSerializer, KafkaAvroSerializer, KafkaProtobufSerializer) registram e validam o schema na hora de enviar e de ler. Cadastrar um schema pode ser feito pela API, por um console web (como o Redpanda Console) ou automaticamente pelo produtor (auto.register.schemas) — em produção, prefira o cadastro controlado (no pipeline de CI), e não pelo primeiro produtor que subir.

Avro

Apache Avro descreve o schema em JSON (arquivos .avsc) e serializa em binário compacto: sem os nomes dos campos nas mensagens, o que reduz muito o tamanho e o custo de rede em comparação com o JSON, e exige o schema para ler. Pontos principais:

  • Tipos: primitivos (string, int, long, boolean, bytes...), complexos (record, array, map, enum, union) e tipos lógicos (decimal, date, timestamp-millis, uuid) para representar valores como dinheiro sem perder precisão (preferir decimal a double).
  • Campos opcionais: use union com null e um valor padrão (default) — é o que permite adicionar campos de forma compatível.
  • Geração de código: plugins (como o avro-maven-plugin) geram as classes (stubs) a partir dos .avsc/.avdl, de modo semelhante ao que o SOAP/WSDL fazia com XML; o programador trabalha com objetos tipados.
  • Avro IDL (.avdl): linguagem mais legível que o JSON para escrever schemas (e um schema que depende de outro, como Pedido → ItemDoPedido → ItemCardapio).
{ "type": "record", "name": "ItemCardapio", "namespace": "com.exemplo.cardapio",
  "fields": [
    { "name": "id",    "type": "string" },
    { "name": "nome",  "type": "string" },
    { "name": "preco", "type": ["null", {"type": "bytes", "logicalType": "decimal", "precision": 11, "scale": 2}],
      "default": null }
  ] }

Protobuf

Protocol Buffers (Google) também é binário, compacto e multilinguagem, com arquivos .proto. A diferença conceitual: os campos são identificados por números de tag (string nome = 2;), e não por nome — é isso que permite renomear com segurança, mas nunca reutilize ou renumere uma tag existente. Oferece tipos bem-conhecidos (Timestamp, Duration, wrappers) no lugar dos tipos lógicos do Avro, e gera código por meio do compilador protoc. Combina bem com gRPC (Backend).

JSON Schema Avro Protobuf
Formato na rede Texto JSON Binário Binário
Tamanho Maior Pequeno Pequeno
Legível por humanos Sim Não Não
Identificação de campos Nome Posição + schema Número da tag
Evolução de schema Regras do Schema Registry Forte, com default e union Forte, com tags
Uso típico Compatibilidade com APIs REST Padrão do ecossistema Kafka, Big Data Microsserviços com gRPC, mobile

Compatibilidade e evolução de schemas

Schemas mudam. A política de compatibilidade do subject decide quais mudanças o Schema Registry aceita:

Modo Garante Mudanças seguras (exemplos)
BACKWARD (padrão) Consumidores com o schema novo conseguem ler dados antigos (consumidores atualizam primeiro) Remover um campo; adicionar campo com valor padrão
FORWARD Consumidores com o schema antigo conseguem ler dados novos (produtores atualizam primeiro) Adicionar campo; remover campo que tinha padrão
FULL Ambos Só adicionar/remover campos opcionais com padrão
*_TRANSITIVE O mesmo em relação a todas as versões anteriores, e não só à última
NONE Nenhuma verificação Qualquer coisa — use com muito cuidado em produção

Adicionar um campo obrigatório sem valor padrão (como um novo preco) quebra a compatibilidade: mensagens antigas não têm o campo e o consumidor não consegue desserializá-las; o Schema Registry recusa o registro. Caminhos: torná-lo opcional com padrão, ou mudar o modo (NONE) de forma deliberada e coordenada, ou criar um novo tópico/versão do evento. Regra geral: mude de forma aditiva, nunca renomeie nem mude o tipo de um campo existente.

Boas práticas: um schema por tipo de evento; documentar campos; testar a compatibilidade na integração contínua (a API tem um endpoint de verificação); usar USE_LATEST_VERSION com cautela; testes de produtores/consumidores com assincronia tratada por ferramentas como Awaitility (esperar a mensagem chegar em vez de sleep).

Kafka Connect

Definição: Kafka Connect

Framework do Kafka para integrar sistemas externos com tópicos sem escrever código: conectores de origem (source) trazem dados de bancos, arquivos, APIs e filas para tópicos; conectores de destino (sink) levam dados dos tópicos para bancos, buscadores, data lakes. Roda como cluster de workers e é configurado por JSON/REST; converters controlam a serialização (por exemplo, Avro com Schema Registry).

Exemplo do livro: um conector de origem para MongoDB observa a collection de pedidos (por CDC — change data capture, usando o fluxo de mudanças do banco) e publica cada alteração em um tópico de auditoria, sem alterar uma linha do microsserviço de pedidos: ele só grava no banco, e o Connect cuida do resto. Com o Schema Registry integrado, o conector registra o schema automaticamente. É uma alternativa de baixo acoplamento ao outbox pattern (veja Mensageria confiável); o Debezium é o conjunto de conectores de CDC mais usado para bancos relacionais.

Kafka Streams e configurações

Kafka Streams

Kafka Streams é uma biblioteca Java (não um servidor à parte) para processar fluxos contínuos de eventos direto dos tópicos: a aplicação lê um ou mais tópicos, transforma e escreve em outro. Primeiro se declara a topologia (as operações) e só depois streams.start() inicia o processamento. Ideias principais:

  • KStream (sequência de eventos) e KTable (estado atual por chave, como uma tabela atualizada por eventos).
  • Operações: filter, map, groupByKey, count, aggregate, join, e saída para console ou para outro tópico (to("topico")).
  • Janelas de tempo (windowedBy): agregam por intervalos (por exemplo, compras por comprador a cada 5 segundos), úteis para métricas em tempo real.
  • Serdes (serializer/deserializer): classes que dizem como converter chave e valor; o application.id também vira o nome do consumer group.
  • O processamento é feito em micro-lotes de tempo: várias mensagens chegadas no mesmo intervalo são processadas juntas. O estado é guardado localmente e tolerante a falhas (via tópicos internos). Alternativa em SQL: ksqlDB.

Configurações importantes

Onde Configuração Efeito
Broker num.partitions Partições padrão de novos tópicos (o padrão é 1: sem paralelismo)
log.retention.hours (padrão 168 = 7 dias), log.dirs Retenção e diretório dos dados
delete.topic.enable, auto.create.topics.enable Permitir apagar tópicos e criá-los sob demanda (em produção, desligue a criação automática para evitar tópicos por engano)
Consumidor group.id Grupo de consumo
auto.offset.reset (earliest/latest) Desde quando ler se não há offset salvo
max.poll.records (padrão 500) Mensagens por busca
enable.auto.commit Se o offset é confirmado automaticamente; desligue para confirmar só depois de processar
heartbeat.interval.ms, session.timeout.ms Detecção de consumidor morto (dispara o rebalance)
Produtor acks, retries, enable.idempotence, linger.ms, batch.size, compression.type Garantia de entrega, tentativas, envio em lote e compressão (veja garantias de entrega)

Serializar é converter o objeto para um formato de troca (JSON, Avro, Protobuf) e desserializar, o caminho inverso; o Spring Boot faz isso de forma transparente, mas a escolha do formato é decisão de arquitetura (veja a seção anterior). O Kafka também tem uma API de administração (AdminClient) para criar, listar e apagar tópicos por código, e clientes em várias linguagens (Python, Go, Node).

Testes e execução

  • Testes de unidade do produtor e do consumidor: simule o KafkaTemplate/o consumidor com mocks e teste a lógica de negócio sem Kafka; em integração, use o EmbeddedKafka ou Testcontainers (Qualidade).
  • Contêineres: cada aplicação ganha um Dockerfile, e o docker-compose sobe Kafka, a rede interna e os serviços, com o endereço dos brokers vindo de variável de ambiente (Containers e Docker).

Mensageria confiável: padrões de resiliência

Vale para Kafka, RabbitMQ, SQS e outros brokers.

Padrão Para quê
Confirmação (ack/nack) O consumidor só confirma depois de processar com sucesso; sem confirmação, a mensagem é reentregue
Retentativas (retry) com backoff Repetir com intervalos crescentes (e jitter) para falhas transitórias, sem sobrecarregar o serviço em apuros
Dead Letter Queue (DLQ) Fila/tópico para onde vão as mensagens que esgotaram as tentativas (ex.: pedidos.DLT), preservando-as para análise e reprocessamento manual
Idempotência e deduplicação Guardar o eventId já processado (tabela com chave única) e ignorar repetições; ver Idempotência em consumidores
Mensagem "veneno" (poison pill) Mensagem que sempre falha e travaria a fila; vai para a DLQ após N tentativas
Ordenação Particione por chave; evite paralelismo dentro da mesma chave

Outbox pattern: gravar no banco e publicar sem perder

Definição: Transactional Outbox

Resolve o problema de gravar no banco e publicar um evento de forma atômica (sem transação distribuída). Na mesma transação do banco, o serviço grava o dado de negócio e uma linha em uma tabela outbox; um processo separado (poller ou CDC, como o Debezium) lê essa tabela e publica no broker, marcando a linha como enviada.

sequenceDiagram
    participant S as Serviço
    participant DB as Banco (pedido + outbox)
    participant R as Relay (Debezium/poller)
    participant K as Broker
    S->>DB: BEGIN; grava pedido; grava evento na outbox; COMMIT
    R->>DB: lê eventos não publicados
    R->>K: publica PedidoCriado
    R->>DB: marca como publicado

Como o relay pode publicar duas vezes (at-least-once), os consumidores precisam ser idempotentes. É a base de sagas e de integração confiável entre microsserviços (Microsserviços; CQRS e Event Sourcing).