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):
- 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).
- Requisitos: funcionais, não funcionais (volume de pedidos, latência, disponibilidade) e restrições (por exemplo, implantação on-premises por exigência legal).
- 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).
- Design: coreografia de eventos, com Saga e fluxo de compensação; AsyncAPI/CloudEvents; diagramas ("big picture" e de sequência).
- Teste: cenário do caminho feliz em BDD, automatizado (Qualidade).
- 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
intJava não preenchido vira0e 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 (preferirdecimaladouble). - Campos opcionais: use
unioncomnulle 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, comoPedido→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) eKTable(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.idtambé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 odocker-composesobe 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).