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

Spring Cloud Stream Kafka批量消费不重试问题排查求助

问题排查:Spring Cloud Stream Kafka批量消费重试不生效,直接进入DLQ

问题描述

使用Spring Boot + Spring Cloud Stream 4.0.3 + Kafka Binder,以批量模式消费Kafka Topic时,抛出异常后整个批次直接进入DLQ,配置的重试逻辑完全不生效。

相关配置代码

重试与DLQ配置类

@Configuration
public class KafkaRetryConfig {

    @Bean
    public ListenerContainerCustomizer<AbstractMessageListenerContainer<byte[], byte[]>> customizer(KafkaOperations<Object, Object> bytesTemplate) {
        return (container, destinationName, group) -> {
            container.setCommonErrorHandler(new DefaultErrorHandler(new DeadLetterPublishingRecoverer(bytesTemplate), new FixedBackOff(5000L, 5L)));
        };
    }
}

消费者代码

@Bean
public Consumer<Message<List<records>>> recordsConsumer() {
    return message -> {
        List<records> records= message.getPayload();
        int index = IntStream.range(0, records.size())
                .filter(streamIndex -> records.get(streamIndex).getId().equals("abc123"))
                .findFirst()
                .orElse(-1);
        Acknowledgment acknowledgment = message.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class);
        assert acknowledgment != null;
        try {
            if (index > -1) {
                throw new RuntimeException("runtime exception");
            }
            //message processing logic
            acknowledgment.acknowledge();
        } catch (Exception e) {
            throw new BatchListenerFailedException(records.get(index).toString(),index);
        }
    };
}

应用配置(application.yml)

spring:
  cloud:
    stream:
      default-binder: kafka
      default:
        contentType: application/*+avro
        consumer:
          useNativeDecoding: true
          autoStartup: false
        producer:
          useNativeEncoding: true
      kafka:
        binder:
          autoCreateTopics: false
          brokers: broker
          configuration:
            enable:
              auto.commit: false
              idempotence: true
            max.in.flight.requests.per.connection: 1
            request.timeout.ms: 5000
            security.protocol: SASL_SSL
            sasl:
              kerberos:
                service:
                  name: service-name
              jaas:
                config: com.sun.security.auth.module.Krb5LoginModule required
                  doNotPrompt=true
                  useKeyTab=true
                  useTicketCache=false
                  storeKey=true
                  keyTab="xyz.keytab"
                  principal="principal@PAYCHEX.COM";
            ssl:
              endpoint.identification.algorithm:
              truststore:
                type: JKS
                location: /config/global/payx-cacerts/cacerts
                password: changeit
          consumer-properties:
            client.id: hrs-productsubscription-consumer-test-9
            key.deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer
            value.deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer
            specific.avro.reader: true
            schema.registry.url: schema-registry-url
            max.poll.records: 200
          requiredAcks: -1
          producer-properties:
            key.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
            value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
        bindings:
          recordsConsumer-in-0:
            consumer:
              startOffset: earliest
              resetOffsets: false
              autoCommitOffset: false
              enableDlq: true
              dlqName: dlq-topic-name
              dlqPartitions: 1
              dlqProducerProperties:
                configuration:
                  key.serializer: org.apache.kafka.common.serialization.StringSerializer
                  value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
                  schema.registry.url: schema-registry-url
              configuration:
                group.id: group-id
                schema.registry.url: schema-registry-url
                autoStartup: true
                key.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
                value.deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
                spring:
                  deserializer:
                    key:
                      delegate.class: org.apache.kafka.common.serialization.StringDeserializer
                    value:
                      delegate.class: io.confluent.kafka.streams.serdes.avro.SpecificAvroDeserializer
      bindings:
        recordsConsumer-in-0:
          consumer:
            batch-mode: true
            max-attempts: 2
          destination: topic-name
          group: group-name
          partitioned: true
          concurrency: 8

问题原因及解决方案

1. Binder自带DLQ与自定义ErrorHandler冲突

你同时启用了Spring Cloud Stream Kafka Binder自带的enableDlq: true配置,以及自定义的DefaultErrorHandler。当启用自带DLQ时,Binder会自动生成一套错误处理逻辑,直接覆盖你自定义的重试配置,导致异常发生时直接转发到DLQ,跳过重试。

  • 解决:删除enableDlq: true、dlqName、dlqPartitions等Binder自带的DLQ配置,完全依赖自定义的DefaultErrorHandler实现重试+DLQ逻辑。

2. 批量消费场景下的ErrorHandler配置不足

默认的DefaultErrorHandler对批量消费的支持需要明确配置,确保它能正确识别BatchListenerFailedException并触发重试。

  • 优化配置:在自定义ErrorHandler中添加批量错误处理的适配逻辑,示例如下:
@Configuration
public class KafkaRetryConfig {

    @Bean
    public ListenerContainerCustomizer<AbstractMessageListenerContainer<byte[], byte[]>> customizer(KafkaOperations<Object, Object> bytesTemplate) {
        return (container, destinationName, group) -> {
            DefaultErrorHandler errorHandler = new DefaultErrorHandler(
                new DeadLetterPublishingRecoverer(bytesTemplate),
                new FixedBackOff(5000L, 5L) // 间隔5秒,重试5次
            );
            // 确保BatchListenerFailedException被纳入重试逻辑
            errorHandler.setRetryListeners((record, ex, deliveryAttempt) -> 
                System.out.printf("重试第%d次,异常:%s%n", deliveryAttempt, ex.getMessage())
            );
            container.setCommonErrorHandler(errorHandler);
        };
    }
}

3. 配置层级重复与冲突

你的application.yml中存在重复配置:kafka.bindings.recordsConsumer-in-0.consumer和bindings.recordsConsumer-in-0.consumer都配置了消费者属性,比如max-attempts: 2是Binder层面的重试配置,会和自定义ErrorHandler的重试次数冲突;同时autoStartup在两处都有设置,可能导致配置被覆盖。

  • 解决:清理重复配置,移除bindings.recordsConsumer-in-0.consumer下的max-attempts,确保所有错误处理逻辑统一由自定义DefaultErrorHandler控制。

4. 手动提交与批量消费的配合

你使用了手动提交acknowledgment.acknowledge(),在批量消费场景下,一旦抛出BatchListenerFailedException,容器会终止当前批次处理并触发重试,此时不需要手动提交,重试时会重新拉取整个批次。确保autoCommitOffset: false配置正确,由ErrorHandler控制提交时机。

修改后的核心配置示例

简化后的application.yml(关键部分)

spring:
  cloud:
    stream:
      kafka:
        bindings:
          recordsConsumer-in-0:
            consumer:
              startOffset: earliest
              resetOffsets: false
              autoCommitOffset: false
              # 移除所有自带DLQ相关配置
              dlqProducerProperties:
                configuration:
                  key.serializer: org.apache.kafka.common.serialization.StringSerializer
                  value.serializer: io.confluent.kafka.serializers.KafkaAvroSerializer
                  schema.registry.url: schema-registry-url
              # ... 其他配置保留
      bindings:
        recordsConsumer-in-0:
          consumer:
            batch-mode: true
            # 移除max-attempts配置
          # ... 其他配置保留

内容的提问来源于stack exchange,提问作者Mira

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 01:17:03