From b14f55e6ac76827cd439d61081e00af6cee7df41 Mon Sep 17 00:00:00 2001 From: hoangtuzami Date: Sat, 16 May 2026 11:18:20 +0700 Subject: [PATCH] chore(config): add Kafka consumer config with DLT recoverer Co-Authored-By: Claude Opus 4.7 (1M context) --- .../configurations/KafkaConsumerConfig.java | 70 +++++++++++++++++++ 1 file changed, 70 insertions(+) create mode 100644 src/main/java/com/isums/scheduleservice/configurations/KafkaConsumerConfig.java diff --git a/src/main/java/com/isums/scheduleservice/configurations/KafkaConsumerConfig.java b/src/main/java/com/isums/scheduleservice/configurations/KafkaConsumerConfig.java new file mode 100644 index 0000000..7c73574 --- /dev/null +++ b/src/main/java/com/isums/scheduleservice/configurations/KafkaConsumerConfig.java @@ -0,0 +1,70 @@ +package com.isums.scheduleservice.configurations; + +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.serialization.StringSerializer; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; +import org.springframework.kafka.listener.DefaultErrorHandler; +import org.springframework.util.backoff.ExponentialBackOff; + +import java.util.Map; + +@Configuration +public class KafkaConsumerConfig { + + @Value("${spring.kafka.bootstrap-servers}") + private String bootstrapServers; + + @Bean + public KafkaTemplate objectKafkaTemplate() { + Map props = Map.of( + ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers, + ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class, + ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, + org.springframework.kafka.support.serializer.JsonSerializer.class, + org.springframework.kafka.support.serializer.JsonSerializer.ADD_TYPE_INFO_HEADERS, false + ); + return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props)); + } + + @Bean + public KafkaTemplate dltKafkaTemplate() { + Map props = Map.of( + ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers, + ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class, + ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class + ); + return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props)); + } + + @Bean + public DefaultErrorHandler kafkaErrorHandler(KafkaTemplate dltKafkaTemplate) { + DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer( + dltKafkaTemplate, + (record, ex) -> new TopicPartition(record.topic() + ".DLT", record.partition()) + ); + + ExponentialBackOff backOff = new ExponentialBackOff(1_000L, 2.0); + backOff.setMaxInterval(60_000L); + backOff.setMaxAttempts(Long.MAX_VALUE); + + DefaultErrorHandler handler = new DefaultErrorHandler(recoverer, backOff); + + handler.addNotRetryableExceptions( + com.fasterxml.jackson.core.JsonProcessingException.class, + com.fasterxml.jackson.databind.exc.InvalidDefinitionException.class, + com.fasterxml.jackson.databind.exc.UnrecognizedPropertyException.class, + IllegalArgumentException.class, + org.springframework.messaging.converter.MessageConversionException.class, + org.springframework.dao.DataIntegrityViolationException.class, + org.hibernate.exception.ConstraintViolationException.class + ); + + return handler; + } +}