Spring Kafka:如何屏蔽RetryTopicConfigurer的错误日志
问题
我有一个包含两个分区的Kafka Topic,正在使用@RetryableTopic注解。应用控制台会输出INFO级别的日志:
INFO o.s.k.r.RetryTopicConfigurer - Received message in dlt listener: {带第二个分区的Topic名称}
但这是错误的,因为消息来自原Topic的第二个分区,而非死信队列(DLT)Topic。我部署了两个应用实例,分别处理该Topic的分区0和1,属于同一个消费组。请问如何屏蔽或避免这类错误日志?
相关代码与配置
@RetryableTopic注解配置
@RetryableTopic( attempts = "1", backoff = @Backoff(delay = 100, multiplier = 3.0), autoCreateTopics = "false", topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE, numPartitions = "2")
Kafka监听方法
@KafkaListener(id= "ccn2_listener",topics = "test", groupId = "test", autoStartup = "${listen.auto.start:true}", topicPartitions = { @TopicPartition(topic = "ccn2-bam-raw-data", partitions = {"1")}) public void listen(ConsumerRecord<String, String> consumerRecord, Acknowledgment acknowledgment) throws IOException, InterruptedException { log.info(consumerRecord.key()); log.info(consumerRecord.value()); // 业务处理代码 if(满足特定条件) { throw new RodaTableMappingException("Kafka记录映射出错,将发送至DLT Topic"); } // 手动提交确认 acknowledgment.acknowledge(); }
应用配置文件(application.properties)
elastic.apm.enabled=true elastic.apm.server-url=url elastic.apm.service-name=name elastic.apm.secret-token=token elastic.apm.environment=prod elastic.apm.application-packages=package elastic.apm.log-level=INFO apminsight.console.logger=true
控制台错误日志示例
2023-01-24 02:16:08,824 [topic_listener-dlt-0-C-1] INFO o.s.k.r.RetryTopicConfigurer - Received message in dlt listener: topic-1@38262.
消费者配置类
@Bean public ConsumerFactory<String, String> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfigurations()); } @Bean public Map<String, Object> consumerConfigurations() { Map<String, Object> configurations = new HashMap<>(); configurations.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBroker0); configurations.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); configurations.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); configurations.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); configurations.put("ssl.truststore.location", truststoreLocation); configurations.put("ssl.truststore.password", truststorePassword); configurations.put("security.protocol", "SSL"); configurations.put("ssl.keystore.location", keystoreLocation); configurations.put("ssl.keystore.password", keyPassword); configurations.put("ssl.key.password", keyPassword); configurations.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configurations.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configurations.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); return configurations; } @Bean ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; }
生产者配置类
@Bean public Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put("security.protocol", "SSL"); props.put("ssl.truststore.location", truststoreLocation); props.put("ssl.truststore.password", truststorePassword); props.put("ssl.keystore.location", keystoreLocation); props.put("ssl.keystore.password", keyPassword); props.put("ssl.key.password", keyPassword); return props; } @Bean public ProducerFactory<String, String> producerFactory() { return new DefaultKafkaProducerFactory<>(producerConfigs()); } @Bean public KafkaTemplate<String, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); }
解决方案
方案1:调整RetryTopicConfigurer日志级别
直接针对org.springframework.kafka.retrytopic.RetryTopicConfigurer类降低日志级别,屏蔽误报的INFO日志,同时保留警告/错误级别的关键信息。在应用配置文件中添加:
logging.level.org.springframework.kafka.retrytopic.RetryTopicConfigurer=WARN
方案2:修复@RetryableTopic后缀策略
你使用的TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE会给重试/DLT Topic添加数字后缀(如topic-1),和原Topic的分区命名逻辑冲突,导致日志误判。建议更换为更明确的后缀策略:
@RetryableTopic( attempts = "1", backoff = @Backoff(delay = 100, multiplier = 3.0), autoCreateTopics = "false", topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_RETRY_ATTEMPT, numPartitions = "2", dltTopicSuffix = "-dlt" // 显式指定DLT后缀,彻底区分分区与DLT命名 )
修改后重试Topic会命名为topic-retry-1,DLT Topic命名为topic-dlt,从根源上解决日志误报问题。
方案3:修正@KafkaListener配置
你的@KafkaListener同时指定了topics和topicPartitions参数,可能导致监听逻辑混乱。应只保留topicPartitions明确指定消费的Topic和分区:
@KafkaListener( id= "ccn2_listener", groupId = "test", autoStartup = "${listen.auto.start:true}", topicPartitions = { @TopicPartition(topic = "ccn2-bam-raw-data", partitions = {"1"})} )
确保两个实例的partitions参数分别设为{"0"}和{"1"},避免跨分区消费引发的日志判断异常。
内容的提问来源于stack exchange,提问作者Kamalama
相关产品推荐
相关产品推荐

