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

SpringBoot KafkaListener异常处理:ConnectException直送DLT,4xx/5xx重试

解决方案:异常处理逻辑实现与DLT序列化问题修复

一、先解决DLT序列化异常

问题根源

你看到的SerializationException是因为死信队列(DLT)的生产者默认使用了StringSerializer,但要发送的是Protobuf类型的IngestionHttpRequest$HttpRequest对象,类型强制转换失败导致。

修复步骤

  1. 配置Protobuf序列化器与生产者工厂
    明确指定生产者使用Confluent的Protobuf序列化器,绑定Schema Registry地址:

    @Bean
    public ProducerFactory<String, IngestionHttpRequest.HttpRequest> protobufProducerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka broker地址");
        config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaProtobufSerializer.class);
        config.put(KafkaProtobufSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "你的Schema Registry地址");
        return new DefaultKafkaProducerFactory<>(config);
    }
    
    @Bean
    public KafkaTemplate<String, IngestionHttpRequest.HttpRequest> protobufKafkaTemplate() {
        return new KafkaTemplate<>(protobufProducerFactory());
    }
    
  2. 自定义死信发布恢复器
    使用上面配置的Protobuf KafkaTemplate创建恢复器,确保DLT发送的消息用Protobuf序列化:

    @Bean
    public DeadLetterPublishingRecoverer dltRecoverer(KafkaTemplate<String, IngestionHttpRequest.HttpRequest> protobufKafkaTemplate) {
        // 指定DLT主题命名规则:原主题名 + "-dlt"
        return new DeadLetterPublishingRecoverer(protobufKafkaTemplate,
                (record, ex) -> new TopicPartition(record.topic() + "-dlt", record.partition()));
    }
    

二、实现异常区分处理逻辑

需要让ConnectException直接触发DLT,4xx/5xx状态码异常保留重试,可通过SeekToCurrentErrorHandler实现简洁的异常控制:

1. 配置错误处理器

@Bean
public SeekToCurrentErrorHandler kafkaErrorHandler(DeadLetterPublishingRecoverer dltRecoverer) {
    // 对可重试异常(4xx/5xx)最多重试3次,间隔1秒
    FixedBackOff backOff = new FixedBackOff(1000, 3);
    SeekToCurrentErrorHandler errorHandler = new SeekToCurrentErrorHandler(dltRecoverer, backOff);
    
    // 标记ConnectException为不可重试,直接发DLT
    errorHandler.addNotRetryableExceptions(ConnectException.class);
    
    // 补充处理嵌套的ConnectException(比如RestTemplate抛出的ResourceAccessException的cause)
    errorHandler.setRetryListeners(new RetryListener() {
        @Override
        public <T, E extends Throwable> boolean open(RetryContext context, RetryCallback<T, E> callback) {
            Throwable ex = context.getLastThrowable();
            if (ex instanceof ResourceAccessException && ex.getCause() instanceof ConnectException) {
                context.setAttribute("skipRetry", true);
                return false;
            }
            return true;
        }

        @Override
        public <T, E extends Throwable> void close(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {}

        @Override
        public <T, E extends Throwable> void onError(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {}
    });
    return errorHandler;
}

2. 在KafkaListener中绑定错误处理器

@KafkaListener(
        topics = "ingestion-topic",
        groupId = "ingestion-group",
        errorHandler = "kafkaErrorHandler"
)
public void processMessage(IngestionHttpRequest.HttpRequest message) {
    ResponseEntity<String> response = restTemplate.postForEntity(
            "你的HTTP端点地址",
            message,
            String.class
    );
    
    // 根据状态码抛出对应异常,触发重试逻辑
    if (response.getStatusCode().is4xxClientError()) {
        throw new HttpClientErrorException(response.getStatusCode());
    }
    if (response.getStatusCode().is5xxServerError()) {
        throw new HttpServerErrorException(response.getStatusCode());
    }
}

三、验证要点

  • 确保DLT主题的配置与原主题一致,Schema Registry已注册IngestionHttpRequest.HttpRequest的Protobuf schema
  • 检查生产者配置中的schema.registry.url是否正确,Confluent Protobuf序列化器依赖Schema Registry完成序列化
  • 测试嵌套异常场景:比如RestTemplate抛出ResourceAccessException时,确认其内部的ConnectException会被识别并直接发DLT

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 12:20:35