Spring-Kafka:如何为ConcurrentMessageListenerContainer配置RetryTopicConfiguration?
手动配置Kafka消费者关联重试主题与死信队列
步骤1:创建RetryTopicConfiguration Bean
定义重试规则的配置Bean,可适配泛型消息类型,按需调整重试策略:
import org.springframework.kafka.config.RetryTopicConfiguration; import org.springframework.kafka.config.RetryTopicConfigurationBuilder; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.context.annotation.Bean; @Bean public <V> RetryTopicConfiguration genericRetryTopicConfig(KafkaProperties props, KafkaTemplate<String, V> kafkaTemplate) { return RetryTopicConfigurationBuilder.newInstance() .fixedBackOff(2000) // 固定退避间隔2秒 .maxAttempts(3) // 最大重试次数(包含首次消费) .useSingleTopicForFixedDelays() // 固定延迟场景复用单个重试主题 .deadLetterTopicSuffix("-dlt") // 死信队列后缀 .retryTopicSuffix("-retry") // 重试主题后缀 .create(kafkaTemplate); }
步骤2:修改MyCustomListener,注入并应用重试配置
将重试配置注入Listener,替换原容器创建逻辑,让重试规则生效:
import org.springframework.kafka.config.RetryTopicConfiguration; import javax.validation.constraints.NotNull; import java.util.HashMap; import org.apache.kafka.common.serialization.StringDeserializer; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; public class MyCustomListener<V> implements MessageListener<String, V> { private final String topicName; private final Class<V> recordType; private final String recordTypeName; private final KafkaProperties props; private final Logger logger; private final ConcurrentMessageListenerContainer<String, V> listenerContainer; private final RetryTopicConfiguration retryTopicConfiguration; // 构造方法新增重试配置注入 public MyCustomListener( @NotNull String topicName, @NotNull Class<V> recordType, @NotNull KafkaProperties props, Logger logger, RetryTopicConfiguration retryTopicConfiguration ) { this.topicName = topicName; this.recordType = recordType; this.recordTypeName = recordType.getSimpleName(); this.props = props; this.logger = logger; this.retryTopicConfiguration = retryTopicConfiguration; this.listenerContainer = setup(); } protected ConcurrentMessageListenerContainer<String, V> setup() { logger.info("Setting up Kafka listener container for topic {}/{}", topicName, recordTypeName); var consumerFactory = new DefaultKafkaConsumerFactory<>( new HashMap<>() {{ put(BOOTSTRAP_SERVERS_CONFIG, props.getBootstrapServers()); put(GROUP_ID_CONFIG, props.getConsumer().getGroupId()); }}, new StringDeserializer(), new MyDataDeserializer<>() ); var containerProps = new ContainerProperties(topicName); containerProps.setMessageListener(this); // 用重试配置初始化并配置容器,替代直接实例化 return retryTopicConfiguration.configureContainer( new ConcurrentMessageListenerContainer<>(consumerFactory, containerProps), consumerFactory ); } @Override public final void onMessage(ConsumerRecord<String, V> record) { logger.info( "(#{}) Received {} message in Kafka: ({})", Thread.currentThread().getId(), recordTypeName, record.key() ); try { // 业务处理逻辑抽离,便于异常捕获触发重试 processRecord(record); } catch (Exception e) { logger.error("Processing failed for record key: {}", record.key(), e); // 抛出异常触发重试机制,达到重试次数后消息会进入死信队列 throw new RuntimeException("Record processing failed", e); } } // 业务处理方法,子类可重写 protected void processRecord(ConsumerRecord<String, V> record) { // 你的业务逻辑实现 } }
关键注意事项
- 重试触发条件:必须在业务处理中抛出异常,重试机制才会将消息转发到重试主题;达到最大重试次数后,消息自动进入死信队列。
- 主题命名:默认重试主题为
原主题名-retry,死信队列为原主题名-dlt,可通过RetryTopicConfigurationBuilder的方法自定义后缀。 - 生命周期管理:建议将MyCustomListener注册为Spring Bean,交由Spring管理容器的启动/停止,避免手动调用
container.start()引发的问题。 - 多规则适配:若不同类型消息需不同重试策略,可创建多个RetryTopicConfiguration Bean,注入时通过
@Qualifier指定对应Bean。
内容的提问来源于stack exchange,提问作者Higher-Kinded Type
相关产品推荐
相关产品推荐

