SpringBoot KafkaListener异常处理:ConnectException直送DLT,4xx/5xx重试
解决方案:异常处理逻辑实现与DLT序列化问题修复
一、先解决DLT序列化异常
问题根源
你看到的SerializationException是因为死信队列(DLT)的生产者默认使用了StringSerializer,但要发送的是Protobuf类型的IngestionHttpRequest$HttpRequest对象,类型强制转换失败导致。
修复步骤
配置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()); }自定义死信发布恢复器
使用上面配置的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
相关产品推荐
相关产品推荐

