Pular para o conteúdo principal
Sumário do livro

Parte V — Java e Spring Boot

Consumer com Spring Kafka

@KafkaListener na prática — groupId, concurrency, Acknowledgment, desserialização e tratamento de exceções com retry e DLQ.

Nesta página

O Capítulo 13 mostrou o lado que publica. Este capítulo mostra o lado que consome — onde a maioria dos conceitos das Partes II a IV (Consumer Group, offset/commit, retry/DLQ, idempotência) se materializa em código real.

@KafkaListener

@KafkaListener

@KafkaListener é a anotação do Spring Kafka que registra um método como consumer de um ou mais tópicos. Por trás dela, o Spring Kafka gerencia o ciclo de vida completo do KafkaConsumer nativo — poll loop, deserialização, e (dependendo da configuração) commit de offset.

@Component
public class PagamentoEventListener {

    @KafkaListener(topics = "pagamentos.aprovados", groupId = "notificacao-service")
    public void processar(PagamentoAprovadoEvent evento) {
        notificacaoService.enviarPush(evento.getContaId(), evento.getValor());
    }
}

groupId define o Consumer Group (Capítulo 6) — se omitido no método, cai no valor padrão de spring.kafka.consumer.group-id da aplicação.

Concurrency

ProducerTopic: pagamentosPartition 0Partition 1Partition 2Consumer Group: notificacao-serviceConsumer 0Consumer 1Consumer 2
Cada partition do tópico é atribuída a exatamente um consumer do Consumer Group — concurrency define quantas threads a própria instância da aplicação registra como consumers adicionais.
@KafkaListener(
    topics = "pagamentos.aprovados",
    groupId = "notificacao-service",
    concurrency = "3"
)
public void processar(PagamentoAprovadoEvent evento) { /* ... */ }

Concurrency não ultrapassa o número de partitions

concurrency = "3" cria 3 threads consumidoras dentro da mesma instância da aplicação, cada uma se comportando como um membro adicional do Consumer Group. Configurar concurrency maior que o número de partitions disponíveis para essa instância cria threads ociosas — o mesmo limite do Capítulo 6 se aplica aqui, só que dentro de um único processo em vez de entre múltiplos pods.

Acknowledgment: commit manual na prática

Acknowledgment

Acknowledgment é o objeto que o Spring Kafka injeta no método do listener quando o AckMode está configurado como manual, permitindo à aplicação decidir explicitamente quando o offset deve ser commitado.

@KafkaListener(topics = "pagamentos.aprovados", groupId = "notificacao-service")
public void processar(PagamentoAprovadoEvent evento, Acknowledgment ack) {
    notificacaoService.enviarPush(evento.getContaId(), evento.getValor());
    ack.acknowledge();
}

Isso exige, na configuração do ContainerFactory, enable.auto.commit=false e AckMode.MANUAL (ou MANUAL_IMMEDIATE) — o padrão discutido no Capítulo 7, aqui aplicado em código: ack.acknowledge() só é chamado depois que a lógica de negócio (aqui, o envio da notificação) termina com sucesso.

Desserialização

O value-deserializer configurado (JSON, Avro, Protobuf — espelhando a escolha do producer no Capítulo 13) converte o payload bruto de volta para o tipo Java esperado pelo método do listener. Uma incompatibilidade de schema entre o que o producer serializou e o que o consumer espera desserializar é uma fonte comum de erro permanente (Capítulo 9) — a mensagem nunca vai desserializar corretamente, não importa quantas vezes seja tentada.

Tratamento de exceções: retry e DLQ na prática

@Bean
public DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
    var recoverer = new DeadLetterPublishingRecoverer(template);
    var backoff = new ExponentialBackOff(1000L, 2.0);
    backoff.setMaxInterval(30_000L);
    return new DefaultErrorHandler(recoverer, backoff);
}

Esse DefaultErrorHandler, registrado no ContainerFactory, implementa exatamente o padrão do Capítulo 9: tentativas com backoff exponencial e, ao esgotá-las, publicação automática no tópico de DLQ (por convenção, <topico>.DLT) via DeadLetterPublishingRecoverer — sem que o método do listener precise tratar isso manualmente.

"Uma exceção não tratada no listener derruba a aplicação"

Uma exceção lançada dentro do método @KafkaListener não derruba a aplicação — ela é capturada pelo container do Spring Kafka e encaminhada ao ErrorHandler configurado (retry, depois DLQ). O erro real de não configurar um ErrorHandler adequado não é a aplicação cair, é o comportamento padrão (retry indefinido, ou log genérico) não ser o que o time espera para aquele fluxo de negócio.

Idempotência no consumer

Todo o cuidado de commit manual e retry não substitui a idempotência (Capítulo 11) — eles resolvem quando o offset avança e o que fazer quando o processamento falha, não o que acontece se a mesma mensagem for entregue duas vezes.

@KafkaListener(topics = "pix.recebido", groupId = "saldo-service")
@Transactional
public void processar(PixRecebidoEvent evento, Acknowledgment ack) {
    try {
        eventosProcessadosRepository.insert(evento.getEventId());
    } catch (DataIntegrityViolationException e) {
        ack.acknowledge();
        return;
    }
    contaRepository.creditar(evento.getContaId(), evento.getValor());
    ack.acknowledge();
}

Como aparece em entrevistas

"Como você implementaria um consumer Kafka em Spring Boot com garantias de at-least-once e proteção contra duplicidade?" é uma pergunta que amarra praticamente todo o livro. A resposta completa cita: @KafkaListener com groupId explícito, AckMode.MANUAL com Acknowledgment.acknowledge() após o processamento, DefaultErrorHandler com backoff e DLQ, e checagem de eventId na mesma transação do efeito de negócio.

Dica de entrevista

Se o entrevistador pedir para "desenhar" um consumer completo, monte a resposta nessa ordem: listener e group, ack manual, error handler com retry/DLQ, idempotência. Essa sequência espelha exatamente a jornada dos Capítulos 6 a 11 — e mostra que os conceitos se conectam, não são tópicos isolados.

Relação com sistemas reais

Consumer completo do saldo-service

O saldo-service roda com concurrency = "6" (igual ao número de partitions do tópico pix.recebido), AckMode.MANUAL, um DefaultErrorHandler com backoff de 1s/5s/30s e DLQ, e checagem de eventId na mesma transação do crédito. O resultado: nenhuma mensagem perdida (at-least-once), nenhum crédito duplicado (idempotência), e mensagens problemáticas isoladas na DLQ sem travar o processamento das demais contas.

Resumo

@KafkaListener registra o consumer; groupId define o Consumer Group; concurrency cria threads adicionais dentro da instância, limitadas pelo número de partitions disponíveis. Acknowledgment viabiliza commit manual após o processamento. DefaultErrorHandler com DeadLetterPublishingRecoverer implementa retry com backoff e DLQ declarativamente. Nenhum desses mecanismos substitui a checagem de idempotência via eventId na mesma transação do efeito de negócio — eles são complementares, não alternativos.

Pode vir a seguir

Prováveis follow-ups: "o que acontece se ack.acknowledge() for chamado antes da exceção ser lançada?" e "como você testaria esse consumer, incluindo o cenário de mensagem duplicada?".