You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 06:25:28