KafkaTemplate.send认证失败无限重试:如何设置3次重试后停止
问题描述
调用kafkaTemplate.send(message)时,因凭证不匹配导致认证失败后,生产者会无限重试连接,日志翻译如下:
2022-10-19 16:03:32.410 INFO [,,,] 828 --- [ad | producer-2] o.apache.kafka.common.network.Selector : [Producer客户端ID=producer-2] 认证失败(无效的用户名或密码)
2022-10-19 16:03:32.412 ERROR [,,,] 828 --- [ad | producer-2] org.apache.kafka.clients.NetworkClient : [Producer客户端ID=producer-2] 连接节点-2失败,原因:认证失败:无效的用户名或密码
2022-10-19 16:03:34.417 INFO [,,,] 828 --- [ad | producer-2] o.apache.kafka.common.network.Selector : [Producer客户端ID=producer-2] 认证失败(无效的用户名或密码)
2022-10-19 16:03:34.419 ERROR [,,,] 828 --- [ad | producer-2] org.apache.kafka.clients.NetworkClient : [Producer客户端ID=producer-2] 连接节点-2失败,原因:认证失败:无效的用户名或密码
2022-10-19 16:03:36.338 INFO [,,,] 828 --- [ad | producer-2] o.apache.kafka.common.network.Selector : [Producer客户端ID=producer-2] 认证失败(无效的用户名或密码)
2022-10-19 16:03:36.340 ERROR [,,,] 828 --- [ad | producer-2] org.apache.kafka.clients.NetworkClient : [Producer客户端ID=producer-2] 连接节点-1失败,原因:认证失败:无效的用户名或密码
解决方案
1. 配置生产者重试参数
直接通过Kafka生产者配置,限制重试次数并排除认证失败类异常的重试:
spring: kafka: producer: retries: 3 retry-backoff-ms: 1000 properties: retry: retry-on-errors: "retriable_authorization_exception,retriable_topic_exception" exclude-on-errors: "authentication_exception"
如果使用Java配置类:
@Configuration public class KafkaProducerConfig { @Bean public ProducerFactory<String, Object> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka集群地址"); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); // 重试次数 configProps.put(ProducerConfig.RETRIES_CONFIG, 3); // 重试间隔(毫秒) configProps.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 1000); // 仅对可恢复异常重试,排除认证失败 configProps.put("retry.retry-on-errors", "retriable_authorization_exception,retriable_topic_exception"); configProps.put("retry.exclude-on-errors", "authentication_exception"); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public KafkaTemplate<String, Object> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } }
2. 自定义监听实现精细重试控制
实现ProducerListener接口,针对认证失败异常单独计数,达到3次后停止重试:
@Component public class CustomKafkaProducerListener implements ProducerListener<String, Object> { private final Map<String, AtomicInteger> retryCounter = new ConcurrentHashMap<>(); @Override public void onError(ProducerRecord<String, Object> record, Exception exception, boolean isRetry) { if (exception instanceof AuthenticationException) { String producerId = record.headers().lastHeader("client_id") != null ? new String(record.headers().lastHeader("client_id").value()) : "default-producer"; AtomicInteger count = retryCounter.computeIfAbsent(producerId, k -> new AtomicInteger(0)); int currentRetry = count.incrementAndGet(); if (currentRetry >= 3) { log.error("认证失败重试已达3次,停止重试,生产者ID:{},异常:{}", producerId, exception.getMessage()); retryCounter.remove(producerId); throw new RuntimeException("Kafka认证失败重试次数耗尽", exception); } else { log.warn("认证失败,第{}次重试,生产者ID:{}", currentRetry, producerId); } } } }
注入到KafkaTemplate:
@Bean public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> producerFactory, CustomKafkaProducerListener listener) { KafkaTemplate<String, Object> template = new KafkaTemplate<>(producerFactory); template.setProducerListener(listener); return template; }
3. 全局异常捕获
添加Spring全局异常处理器,统一处理认证失败及重试耗尽的情况:
@RestControllerAdvice public class GlobalExceptionHandler { @ExceptionHandler(AuthenticationException.class) public ResponseEntity<String> handleKafkaAuthError(AuthenticationException e) { log.error("Kafka认证失败:{}", e.getMessage()); return ResponseEntity.status(HttpStatus.UNAUTHORIZED).body("Kafka凭证无效,认证失败"); } @ExceptionHandler(RuntimeException.class) public ResponseEntity<String> handleRetryExhausted(RuntimeException e) { if (e.getMessage().contains("认证失败重试次数耗尽")) { log.error("Kafka重试耗尽:{}", e.getMessage()); return ResponseEntity.status(HttpStatus.SERVICE_UNAVAILABLE).body("Kafka发送失败,重试次数已用完"); } throw e; } }
关键说明
- 通过
retry.exclude-on-errors直接排除认证异常,从根源阻止生产者对这类不可恢复错误的重试 - 自定义监听器可以针对不同生产者实例单独控制重试次数,适配多场景需求
- 全局异常处理器可以统一返回友好提示,避免异常扩散影响业务流程
内容的提问来源于stack exchange,提问作者charan charan

